refactor(cron): one FIRE_CLAIM_TTL_SECONDS for claiming, one-shot re-arm, and stale-error recovery

The TTL was three copies of the literal 300 (claim_job_for_fire default,
rearm_oneshot, and the new stale-error guard). Hoisted so the three lanes
cannot drift; incident narrative in the guard comment cut to the WHY.
This commit is contained in:
kshitijk4poor
2026-09-07 00:04:22 +05:30
committed by kshitij
parent 37a0c1b95f
commit 193f05dec5
2 changed files with 9 additions and 14 deletions
+8 -13
View File
@@ -886,9 +886,9 @@ def _classify_dispatch_lateness(lateness_seconds: float, grace_seconds: int) ->
_persisted_error_recoveries: int = 0 _persisted_error_recoveries: int = 0
# Bounded in-memory history kept by every probe-visible fire-path counter. # Bounded in-memory history kept by every probe-visible fire-path counter.
_TELEMETRY_RECENT_HISTORY = 20 _TELEMETRY_RECENT_HISTORY = 20
# A fire_claim younger than this is treated as a live run in another process # A fire_claim younger than this is a live run (heartbeat cadence is 60 s). One value
# (heartbeat cadence is 60 s; matches the claim TTL used by claim_job_fire). # for claiming, one-shot re-arm, and stale-error recovery so they cannot disagree.
_STALE_ERROR_FIRE_CLAIM_TTL_SECONDS = 300.0 FIRE_CLAIM_TTL_SECONDS = 300
_persisted_error_recoveries_recent: list = [] _persisted_error_recoveries_recent: list = []
@@ -909,14 +909,9 @@ def _job_is_stale_error_recurring(
return False return False
if _job_running_in_this_process(str(job.get("id") or "")): if _job_running_in_this_process(str(job.get("id") or "")):
return False return False
# Cross-process liveness (2026-09-02 brain incident): a run owned by # A fresh fire_claim means the job is running in ANOTHER process sharing this
# ANOTHER scheduler process sharing this jobs store (e.g. several Desktop # store, not wedged; re-arming it here would only claim-fight the live run.
# profile tabs + the gateway) heartbeats ``fire_claim`` every if _claim_is_live(job.get("fire_claim"), now, FIRE_CLAIM_TTL_SECONDS):
# _RUN_CLAIM_HEARTBEAT_SECONDS. While that claim is fresh the job is not
# wedged — it is running elsewhere — and re-arming it here makes every
# other ticker claim-fight the live run (171 re-arms, ~175 junk
# "Fire claim lost" execution rows, then the live run was killed).
if _claim_is_live(job.get("fire_claim"), now, _STALE_ERROR_FIRE_CLAIM_TTL_SECONDS):
return False return False
last_run = job.get("last_run_at") last_run = job.get("last_run_at")
last_run_dt = _parse_aware(last_run) if last_run else None last_run_dt = _parse_aware(last_run) if last_run else None
@@ -2057,7 +2052,7 @@ def rearm_oneshot(job_id: str, run_at: Any) -> Optional[Dict[str, Any]]:
now = _hermes_now() now = _hermes_now()
if _claim_is_live(job.get("run_claim"), now, _oneshot_run_claim_ttl_seconds()): if _claim_is_live(job.get("run_claim"), now, _oneshot_run_claim_ttl_seconds()):
raise ValueError("Cannot re-arm one-shot over a live run claim.") raise ValueError("Cannot re-arm one-shot over a live run claim.")
if _claim_is_live(job.get("fire_claim"), now, 300): if _claim_is_live(job.get("fire_claim"), now, FIRE_CLAIM_TTL_SECONDS):
raise ValueError("Cannot re-arm one-shot over a live fire claim.") raise ValueError("Cannot re-arm one-shot over a live fire claim.")
if job.get("schedule", {}).get("kind") != "once": if job.get("schedule", {}).get("kind") != "once":
raise ValueError(_REARM_RECURRING_ERROR) raise ValueError(_REARM_RECURRING_ERROR)
@@ -2488,7 +2483,7 @@ def _machine_id() -> str:
def claim_job_for_fire( def claim_job_for_fire(
job_id: str, *, claim_ttl_seconds: int = 300, force: bool = False, return_job: bool = False, job_id: str, *, claim_ttl_seconds: int = FIRE_CLAIM_TTL_SECONDS, force: bool = False, return_job: bool = False,
) -> Union[bool, Dict[str, Any]]: ) -> Union[bool, Dict[str, Any]]:
"""Atomically claim a job for one external 'fire' (multi-machine at-most-once); True iff THIS """Atomically claim a job for one external 'fire' (multi-machine at-most-once); True iff THIS
caller won (``CronScheduler.fire_due``: exactly one of N replicas runs a job). Under the caller won (``CronScheduler.fire_due``: exactly one of N replicas runs a job). Under the
@@ -188,6 +188,6 @@ class TestPersistedStaleErrorRecovery:
) )
# Once the foreign claim has expired (owner died), recovery resumes. # Once the foreign claim has expired (owner died), recovery resumes.
expired = (now - timedelta(seconds=J._STALE_ERROR_FIRE_CLAIM_TTL_SECONDS + 1)).isoformat() expired = (now - timedelta(seconds=J.FIRE_CLAIM_TTL_SECONDS + 1)).isoformat()
J.update_job(job_id, {"fire_claim": {"at": expired, "by": "other-host:deadbeef"}}) J.update_job(job_id, {"fire_claim": {"at": expired, "by": "other-host:deadbeef"}})
assert J._job_is_stale_error_recurring(J.get_job(job_id), job["schedule"], now) assert J._job_is_stale_error_recurring(J.get_job(job_id), job["schedule"], now)