fix(cron): don't stamp the next occurrence on an off-tick manual run
claim_job_for_fire() derives the occurrence identity from next_run_at before the same function advances it. On a scheduler tick next_run_at is the occurrence being run, which is correct; on an off-tick manual run it is the NEXT occurrence, so the execution is stamped with the identity of a slot that has not happened yet. _job_is_due() then finds a completed execution carrying that identity and skips the real slot, returning before the last_dispatch write — no error, no log line, no dispatch record. The manual flag already guards this and both _job_is_due() and claim_job_for_fire() honour it; the agent-facing run-now path never declared itself. Add a keyword-only manual= parameter and pass it from _claim_for_manual_run(). Deliberately not force=True: force also calls _activate_job_record(), which would resume a paused or disabled job, and the run-now tool depends on continuing to refuse those. The local flag is renamed to manual_fire so the new parameter is not shadowed inside the apply closure, which would raise UnboundLocalError. Three existing tests in tests/tools/ pinned the old call signature via assert_called_once_with; they now pin manual=True, so dropping the flag again fails loudly rather than silently reintroducing the skip. Restores the intent stated in #104790 — the column records the scheduled instant an execution was claimed for, and an off-tick manual run was claimed for none. Fixes #105690 Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
10
cron/jobs.py
10
cron/jobs.py
@@ -2491,7 +2491,8 @@ def _machine_id() -> str:
|
||||
|
||||
|
||||
def claim_job_for_fire(
|
||||
job_id: str, *, claim_ttl_seconds: int = FIRE_CLAIM_TTL_SECONDS, force: bool = False, return_job: bool = False,
|
||||
job_id: str, *, claim_ttl_seconds: int = FIRE_CLAIM_TTL_SECONDS, force: bool = False,
|
||||
manual: bool = False, return_job: bool = False,
|
||||
) -> Union[bool, Dict[str, Any]]:
|
||||
"""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
|
||||
@@ -2513,8 +2514,11 @@ def claim_job_for_fire(
|
||||
return False # someone holds a fresh claim
|
||||
from cron.occurrences import completed_occurrence, scheduled_instant
|
||||
|
||||
manual = force or job.get("manual_run_at") == job.get("next_run_at")
|
||||
instant = None if manual else scheduled_instant(job.get("next_run_at"))
|
||||
# ``manual`` (an off-tick run-now) must NOT stamp an occurrence identity: outside a
|
||||
# scheduler tick ``next_run_at`` is the NEXT occurrence, not the one being run, so
|
||||
# stamping it would make completed_occurrence() skip that slot when it arrives.
|
||||
manual_fire = force or manual or job.get("manual_run_at") == job.get("next_run_at")
|
||||
instant = None if manual_fire else scheduled_instant(job.get("next_run_at"))
|
||||
if instant and completed_occurrence(job, instant):
|
||||
if job.get("schedule", {}).get("kind") in {"cron", "interval"}:
|
||||
nxt = compute_next_run(job["schedule"], now.isoformat())
|
||||
|
||||
@@ -263,3 +263,49 @@ def test_same_thread_fire_fence_reentrancy_preserves_ownership(temp_home):
|
||||
assert result == {"outer": True, "inner": True}
|
||||
thread.join(timeout=2)
|
||||
assert thread.is_alive() is False
|
||||
|
||||
|
||||
def test_manual_claim_does_not_stamp_a_future_occurrence(temp_home):
|
||||
"""An off-tick run-now must not consume the NEXT scheduled slot.
|
||||
|
||||
Outside a scheduler tick ``next_run_at`` is the occurrence that has NOT happened
|
||||
yet, so stamping it as a completed occurrence makes ``_job_is_due`` skip that slot
|
||||
when it arrives — silently, with no error and no dispatch record. ``manual=True``
|
||||
is the caller's declaration that this is an off-tick fire.
|
||||
"""
|
||||
from cron.jobs import create_job, claim_job_for_fire, get_job
|
||||
|
||||
job = create_job(prompt="x", schedule="every 5m", name="m")
|
||||
pending = get_job(job["id"])["next_run_at"]
|
||||
|
||||
claimed = claim_job_for_fire(job["id"], manual=True, return_job=True)
|
||||
assert isinstance(claimed, dict)
|
||||
assert claimed["_scheduled_instant"] is None, (
|
||||
f"manual fire stamped the future occurrence {pending}")
|
||||
|
||||
|
||||
def test_manual_claim_still_refuses_a_paused_job(temp_home):
|
||||
"""``manual=True`` suppresses only the occurrence stamp — unlike ``force=True`` it
|
||||
must not resume a paused job, which the run-now tool relies on to refuse it."""
|
||||
from cron.jobs import create_job, claim_job_for_fire, get_job, pause_job
|
||||
|
||||
job = create_job(prompt="x", schedule="every 5m", name="mp")
|
||||
pause_job(job["id"])
|
||||
|
||||
assert claim_job_for_fire(job["id"], manual=True) is False
|
||||
assert get_job(job["id"]).get("paused_at") is not None
|
||||
|
||||
|
||||
def test_scheduler_tick_claim_still_stamps_its_occurrence(temp_home):
|
||||
"""The dedupe path stays intact for ordinary (non-manual) claims: a tick fire still
|
||||
records the occurrence it ran, so a re-delivery cannot double-fire it."""
|
||||
from cron.jobs import create_job, claim_job_for_fire, get_job
|
||||
|
||||
job = create_job(prompt="x", schedule="every 5m", name="s")
|
||||
pending = get_job(job["id"])["next_run_at"]
|
||||
|
||||
claimed = claim_job_for_fire(job["id"], return_job=True)
|
||||
assert isinstance(claimed, dict)
|
||||
assert claimed["_scheduled_instant"] is not None
|
||||
from cron.occurrences import scheduled_instant
|
||||
assert claimed["_scheduled_instant"] == scheduled_instant(pending)
|
||||
|
||||
@@ -79,7 +79,7 @@ class TestBackgroundDispatch:
|
||||
assert res["claimed"] is True
|
||||
assert res["dispatched"] is True
|
||||
assert res["delegation_id"]
|
||||
m_claim.assert_called_once_with("job-bg-01", return_job=True)
|
||||
m_claim.assert_called_once_with("job-bg-01", manual=True, return_job=True)
|
||||
# The job actually starts on the daemon executor.
|
||||
assert run_started.wait(timeout=5.0), "job never started in background"
|
||||
finally:
|
||||
@@ -326,5 +326,5 @@ class TestCronjobRunToolIntegration:
|
||||
assert out["success"] is True
|
||||
assert out["job"]["executed"] is True
|
||||
assert out["job"]["execution_success"] is True
|
||||
m_claim.assert_called_once_with("job-bg-13", return_job=True)
|
||||
m_claim.assert_called_once_with("job-bg-13", manual=True, return_job=True)
|
||||
m_run.assert_called_once()
|
||||
|
||||
@@ -39,7 +39,7 @@ class TestCronjobRunExecutesImmediately:
|
||||
assert out["success"] is True
|
||||
assert out["job"]["executed"] is True
|
||||
assert out["job"]["execution_success"] is True
|
||||
m_claim.assert_called_once_with("job-run-1", return_job=True)
|
||||
m_claim.assert_called_once_with("job-run-1", manual=True, return_job=True)
|
||||
m_run.assert_called_once_with(claimed, adapters=None, loop=None, extra_prompt=None)
|
||||
|
||||
def test_run_reconciles_external_provider_after_claimed_execution(self):
|
||||
|
||||
@@ -184,7 +184,7 @@ def _claim_for_manual_run(job_id: str, log_label: str):
|
||||
``(None, error_dict)`` in the ``_execute_job_now`` shape. A lost claim is labelled precisely —
|
||||
claim_job_for_fire also returns False for paused/disabled/missing jobs, not just in-flight ones."""
|
||||
try:
|
||||
claimed_job = claim_job_for_fire(job_id, return_job=True)
|
||||
claimed_job = claim_job_for_fire(job_id, manual=True, return_job=True)
|
||||
if isinstance(claimed_job, dict):
|
||||
return claimed_job, None
|
||||
refreshed = get_job(job_id)
|
||||
|
||||
Reference in New Issue
Block a user