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
12 changes: 12 additions & 0 deletions cron/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -886,6 +886,9 @@ def _classify_dispatch_lateness(lateness_seconds: float, grace_seconds: int) ->
_persisted_error_recoveries: int = 0
# Bounded in-memory history kept by every probe-visible fire-path counter.
_TELEMETRY_RECENT_HISTORY = 20
# A fire_claim younger than this is treated as a live run in another process
# (heartbeat cadence is 60 s; matches the claim TTL used by claim_job_fire).
_STALE_ERROR_FIRE_CLAIM_TTL_SECONDS = 300.0
_persisted_error_recoveries_recent: list = []


Expand All @@ -906,6 +909,15 @@ def _job_is_stale_error_recurring(
return False
if _job_running_in_this_process(str(job.get("id") or "")):
return False
# Cross-process liveness (2026-09-02 brain incident): a run owned by
# ANOTHER scheduler process sharing this jobs store (e.g. several Desktop
# profile tabs + the gateway) heartbeats ``fire_claim`` every
# _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
last_run = job.get("last_run_at")
last_run_dt = _parse_aware(last_run) if last_run else None
if last_run_dt is None:
Expand Down
30 changes: 30 additions & 0 deletions tests/cron/test_recurring_persisted_error_recovery.py
Original file line number Diff line number Diff line change
Expand Up @@ -161,3 +161,33 @@ def test_recent_error_within_cadence_not_force_rearmed(self, cron_env, monkeypat
"a within-cadence error is a normal retry — must not be "
"force-re-armed onto an immediate fire"
)

def test_live_fire_claim_from_other_process_not_rearmed(self, cron_env, monkeypatch):
"""A stale-error job whose ``fire_claim`` is being heartbeated by a run
in ANOTHER process is running, not wedged. Re-arming it would make this
ticker claim-fight the live run (2026-09-02 multi-Desktop-tab incident:
171 re-arms, ~175 junk executions, live run killed)."""
S, E, J, env = _setup(cron_env, monkeypatch)
job_id = env["job_id"]

_persist_stale_error(J, job_id, error_age_minutes=110)
now = datetime.now(timezone.utc)
J.update_job(
job_id,
{"fire_claim": {"at": now.isoformat(), "by": "other-host:deadbeef"}},
)

job = J.get_job(job_id)
assert not J._job_is_stale_error_recurring(job, job["schedule"], now)

with mock.patch("cron.jobs.load_jobs", return_value=[job]):
S.tick(verbose=False, sync=True)
assert E.latest_execution(job_id) is None, (
"a job with a live fire_claim held by another process must not be "
"re-armed by stale-error recovery"
)

# Once the foreign claim has expired (owner died), recovery resumes.
expired = (now - timedelta(seconds=J._STALE_ERROR_FIRE_CLAIM_TTL_SECONDS + 1)).isoformat()
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)