diff --git a/cron/jobs.py b/cron/jobs.py index 5c8808505d..27b7166342 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -2062,13 +2062,35 @@ def trigger_job(job_id: str, extra_prompt: Optional[str] = None) -> Optional[Dic }) +def _claim_owner_is_dead(claim: Dict[str, Any]) -> bool: + """True when the claim's ``by`` names a process on THIS host that provably no longer exists. + ``_machine_id()`` stamps ``host:pid[:token]``; a foreign host, an explicit HERMES_MACHINE_ID, + or any liveness-probe failure returns False (fail safe: only a proven death shortens the TTL).""" + parts = str(claim.get("by") or "").split(":") + if len(parts) < 2 or not parts[1].isdigit(): + return False + try: + import socket + if parts[0] != socket.gethostname(): + return False + from gateway.status import _pid_exists + return not _pid_exists(int(parts[1])) + except Exception: + return False + + def _claim_is_live(claim: Any, now: datetime, ttl_seconds: float) -> bool: - """True for a well-formed claim aged within ``[0, ttl)``: future-dated (clock/TZ skew) or - malformed claims count as stale so they can never wedge a job.""" + """True for a well-formed claim aged within ``[0, ttl)`` whose owner is not provably dead: + future-dated (clock/TZ skew) or malformed claims count as stale so they can never wedge a + job, and a same-host owner pid that has exited releases the claim immediately instead of + after the TTL (a killed ``hermes cron run`` otherwise blocks the next manual run for the + full window with "already being fired").""" if not isinstance(claim, dict) or not claim.get("at"): return False claimed_at = _parse_aware(claim["at"]) - return claimed_at is not None and 0 <= (now - claimed_at).total_seconds() < ttl_seconds + if claimed_at is None or not (0 <= (now - claimed_at).total_seconds() < ttl_seconds): + return False + return not _claim_owner_is_dead(claim) _REARM_RECURRING_ERROR = ( diff --git a/tests/cron/test_claim_job_for_fire.py b/tests/cron/test_claim_job_for_fire.py index aa7f351135..2242dd8f35 100644 --- a/tests/cron/test_claim_job_for_fire.py +++ b/tests/cron/test_claim_job_for_fire.py @@ -294,3 +294,34 @@ def test_manual_claim_still_refuses_a_paused_job(temp_home): assert claim_job_for_fire(job["id"], manual=True) is False assert get_job(job["id"]).get("paused_at") is not None + + +def test_fresh_claim_from_a_dead_same_host_owner_is_reclaimable(temp_home): + """A claim younger than the TTL whose owner pid (same host) has exited is stale at once: a + ``hermes cron run`` killed mid-flight must not block the next manual run for the whole TTL + with "already being fired". A live owner's fresh claim still blocks.""" + import os + import socket + import subprocess + import sys + + from cron.jobs import claim_job_for_fire, create_job, load_jobs, save_jobs + + jid = create_job(prompt="x", schedule="every 5m", name="s")["id"] + assert claim_job_for_fire(jid) is True + + # Live same-host owner (this process) → still blocked. + jobs = load_jobs() + job = next(j for j in jobs if j["id"] == jid) + job["fire_claim"]["by"] = f"{socket.gethostname()}:{os.getpid()}:tok" + save_jobs(jobs) + assert claim_job_for_fire(jid) is False + + # Owner that has provably exited → reclaimable despite the fresh timestamp. + child = subprocess.Popen([sys.executable, "-c", "pass"]) + child.wait() + jobs = load_jobs() + job = next(j for j in jobs if j["id"] == jid) + job["fire_claim"]["by"] = f"{socket.gethostname()}:{child.pid}:tok" + save_jobs(jobs) + assert claim_job_for_fire(jid) is True