diff --git a/cron/jobs.py b/cron/jobs.py index 160bd6e481..195fddba62 100644 --- a/cron/jobs.py +++ b/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()) diff --git a/tests/cron/test_claim_job_for_fire.py b/tests/cron/test_claim_job_for_fire.py index 200199efd7..da6382cb5f 100644 --- a/tests/cron/test_claim_job_for_fire.py +++ b/tests/cron/test_claim_job_for_fire.py @@ -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) diff --git a/tests/tools/test_cronjob_run_background.py b/tests/tools/test_cronjob_run_background.py index 58190889c5..122913e8b4 100644 --- a/tests/tools/test_cronjob_run_background.py +++ b/tests/tools/test_cronjob_run_background.py @@ -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() diff --git a/tests/tools/test_cronjob_run_immediate.py b/tests/tools/test_cronjob_run_immediate.py index aa0eb8b97f..868f4b01d2 100644 --- a/tests/tools/test_cronjob_run_immediate.py +++ b/tests/tools/test_cronjob_run_immediate.py @@ -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): diff --git a/tools/cronjob_tools.py b/tools/cronjob_tools.py index fb0de38258..71208d70bd 100644 --- a/tools/cronjob_tools.py +++ b/tools/cronjob_tools.py @@ -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)