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:
Phil Mossman
2026-09-08 09:16:10 +00:00
committed by kshitij
parent 034eb7649e
commit ac10770894
5 changed files with 57 additions and 7 deletions

View File

@@ -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())

View File

@@ -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)

View File

@@ -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()

View File

@@ -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):

View File

@@ -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)