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
36 changes: 30 additions & 6 deletions cron/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -3621,13 +3621,37 @@ def _sweep_completed_oneshots(


def heartbeat_fire_claim(job_id: str, *, expected_owner: str) -> bool:
"""Refresh this owner's ``fire_claim`` lease.

A heartbeat is not an external side effect: it only compare-and-swaps
``fire_claim["at"]`` when ``fire_claim["by"]`` still equals
``expected_owner``. ``_jobs_lock()`` already serializes that CAS with both
an in-process lock and a cross-process flock, so it is safe on its own.

Taking the fire fence here is actively harmful. The run thread holds the
per-job fence across its whole side effect (see ``_fire_job_lock``), and the
fence's reentrancy is thread-local, so it does not extend to the heartbeat
thread. Any execution outliving ``_RUN_CLAIM_HEARTBEAT_SECONDS`` therefore
blocked here for ``_JOBS_LOCK_TIMEOUT_SECONDS`` and returned False, which
the scheduler reads as "ownership lost" and uses to interrupt a perfectly
healthy run -- reported to the operator as "Interrupted by shutdown before
terminal completion."

So: prefer the fence to keep the common path serialized with other owner
mutations, but when it is unavailable fall back to the CAS instead of
reporting a loss of ownership we have not actually observed. A genuine
takeover still returns False, because ``by`` no longer matches.
"""
with _fire_job_lock(job_id) as acquired:
if not acquired:
return False
return _heartbeat_fire_claim_locked(
job_id,
expected_owner=expected_owner,
)
if acquired:
return _heartbeat_fire_claim_locked(
job_id,
expected_owner=expected_owner,
)

# Fence unavailable (held by this job's own run thread, or contended
# cross-process). The owner check below is authoritative either way.
return _heartbeat_fire_claim_locked(job_id, expected_owner=expected_owner)


def _heartbeat_fire_claim_locked(job_id: str, *, expected_owner: str) -> bool:
Expand Down
44 changes: 44 additions & 0 deletions tests/cron/test_claim_job_for_fire.py
Original file line number Diff line number Diff line change
Expand Up @@ -263,3 +263,47 @@ def reentrant_claimant():
assert result == {"outer": True, "inner": True}
thread.join(timeout=2)
assert thread.is_alive() is False


def test_heartbeat_survives_fence_held_by_its_own_run(temp_home, monkeypatch):
"""A heartbeat must not report ownership lost merely because the fence is busy.

Regression for the false "Interrupted by shutdown before terminal
completion." that killed every cron job outliving the heartbeat interval:
the run thread holds the per-job fence across its side effect, the
heartbeat thread (a different thread, so the fence's thread-local
reentrancy does not apply) could not acquire it, and that timeout was
reported to the scheduler as a lost claim.
"""
import cron.jobs as jobs

job = jobs.create_job(prompt="x", schedule="every 5m", name="long-run")
assert jobs.claim_job_for_fire(job["id"]) is True
owner = jobs.get_job(job["id"])["fire_claim"]["by"]

# Keep the failing path fast; the real value is 30s.
monkeypatch.setattr(jobs, "_JOBS_LOCK_TIMEOUT_SECONDS", 0.2)

fence_held = threading.Event()
release = threading.Event()

def hold_fence():
with jobs.fire_claim_fence(job["id"], expected_owner=owner) as owns:
assert owns is True
fence_held.set()
release.wait(timeout=5)

holder = threading.Thread(target=hold_fence, daemon=True)
holder.start()
try:
assert fence_held.wait(timeout=5)
# The run thread owns the fence; the heartbeat still confirms ownership.
assert jobs.heartbeat_fire_claim(job["id"], expected_owner=owner) is True
# A genuine takeover is still detected while the fence is busy.
assert (
jobs.heartbeat_fire_claim(job["id"], expected_owner="replacement-owner")
is False
)
finally:
release.set()
holder.join(timeout=5)