fix(cron): the dashboard "Trigger" run-now no longer stamps the next occurrence either

Second entry of the same bug class: POST /api/cron/jobs/{id}/trigger →
_fire_cron_job_for_profile → CronScheduler.fire_due → claim_fire built its claim
without `manual`, so an off-tick run from the web UI stamped the future slot exactly
like the tools path #105704 fixes. fire_due/claim_fire gain `manual` (forwarded only
when set, mirroring `force`, so third-party providers keep working) and the dashboard
trigger passes it when the provider's signature accepts it. Webhook and misfire
catch-up fires run the slot that is due and keep the stamp.

Also drops the base-green tick-stamp test (the same contract is pinned by
tests/cron/test_scheduled_occurrence.py) and documents `manual` vs `force`.
This commit is contained in:
kshitijk4poor
2026-09-09 11:25:27 +05:30
committed by kshitij
parent ac10770894
commit a73b750391
6 changed files with 30 additions and 22 deletions
+2 -1
View File
@@ -2498,7 +2498,8 @@ def claim_job_for_fire(
caller won (``CronScheduler.fire_due``: exactly one of N replicas runs a job). Under the
fence + file lock: reject missing/terminal/paused jobs unless ``force`` (explicit manual
fire, which also resumes the job atomically; external callbacks must leave it false so a
stale callback cannot resurrect a paused job). Lose if a claim younger than
stale callback cannot resurrect a paused job). ``manual`` = off-tick run-now without the
resume: no occurrence stamp, so the still-pending ``next_run_at`` slot is not skipped. Lose if a claim younger than
``claim_ttl_seconds`` exists (the TTL lets another fire reclaim after a crash; mark_job_run
clears the claim). Otherwise stamp ``fire_claim`` and, for recurring jobs, advance
``next_run_at`` so a stale re-delivery cannot re-fire."""
+15 -4
View File
@@ -136,16 +136,20 @@ class CronScheduler(ABC):
def fire_due(
self, job_id: str, *, adapters: Any = None, loop: Any = None, force: bool = False,
manual: bool = False,
) -> bool:
"""Run one job NOW (inbound fire webhook entry). Store CAS claim (multi-machine
at-most-once) then shared ``run_one_job``. True if THIS caller claimed and processed the
attempt (even if the job failed); False if the claim was lost or the job is gone."""
claimed_job = self.claim_fire(job_id, force=force)
attempt (even if the job failed); False if the claim was lost or the job is gone.
``manual`` marks an off-tick run-now (dashboard trigger): the claim must not stamp
``next_run_at`` as the occurrence, or that slot is skipped when it arrives. Webhook and
misfire fires run the slot that is due and keep the stamp."""
claimed_job = self.claim_fire(job_id, force=force, manual=manual)
if claimed_job is None:
return False
return self.fire_claimed(claimed_job, adapters=adapters, loop=loop)
def claim_fire(self, job_id: str, *, force: bool = False) -> dict | None:
def claim_fire(self, job_id: str, *, force: bool = False, manual: bool = False) -> dict | None:
"""Durably claim one fire + create its audit attempt. Transports call this synchronously
before acknowledging, then pass the exact snapshot to ``fire_claimed`` off-thread."""
from cron.executions import create_execution, finish_execution, set_execution_occurrence
@@ -155,6 +159,8 @@ class CronScheduler(ABC):
claim_kwargs = {"return_job": True}
if force:
claim_kwargs["force"] = True
if manual:
claim_kwargs["manual"] = True
try:
claimed_job = claim_job_for_fire(job_id, **claim_kwargs)
if isinstance(claimed_job, dict):
@@ -189,6 +195,11 @@ class CronScheduler(ABC):
def provider_supports_force_fire(provider: Any) -> bool:
"""Return whether a provider can safely receive ``fire_due(force=...)`` (signature-detected)."""
return provider_fire_due_accepts(provider, "force")
def provider_fire_due_accepts(provider: Any, name: str) -> bool:
"""Whether ``provider.fire_due`` takes keyword ``name`` (third-party providers may predate it)."""
try:
parameters = inspect.signature(provider.fire_due).parameters.values()
except (TypeError, ValueError):
@@ -196,7 +207,7 @@ def provider_supports_force_fire(provider: Any) -> bool:
return any(
p.kind is inspect.Parameter.VAR_KEYWORD
or (
p.name == "force"
p.name == name
and p.kind in (inspect.Parameter.POSITIONAL_OR_KEYWORD, inspect.Parameter.KEYWORD_ONLY)
)
for p in parameters
+5 -1
View File
@@ -278,7 +278,7 @@ def _fire_cron_job_for_profile(profile: str, job_id: str, *, force: bool = False
and external callers on the web_deps late-binding seam; do not add new uses.
"""
_profile_name, home = _cron_profile_home(profile)
from cron.scheduler_provider import provider_supports_force_fire, resolve_cron_scheduler
from cron.scheduler_provider import provider_fire_due_accepts, provider_supports_force_fire, resolve_cron_scheduler
with _cron_store_scope(home):
provider = resolve_cron_scheduler()
if force:
@@ -289,6 +289,10 @@ def _fire_cron_job_for_profile(profile: str, job_id: str, *, force: bool = False
f"Cron provider '{getattr(provider, 'name', 'custom')}' "
"does not support atomic forced firing of paused jobs"))
return bool(provider.fire_due(job_id, adapters=None, loop=None, force=True))
# Off-tick run-now: never stamp next_run_at as the occurrence (#104790); third-party
# providers without the kwarg keep the legacy call.
if provider_fire_due_accepts(provider, "manual"):
return bool(provider.fire_due(job_id, adapters=None, loop=None, manual=True))
return bool(provider.fire_due(job_id, adapters=None, loop=None))
-15
View File
@@ -294,18 +294,3 @@ def test_manual_claim_still_refuses_a_paused_job(temp_home):
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)
+4
View File
@@ -439,6 +439,10 @@ def test_fire_due_forwards_manual_force_to_store_claim(monkeypatch):
assert InProcessCronScheduler().fire_due("j1", force=True) is True
assert claims == [("j1", {"force": True, "return_job": True})]
# An off-tick run-now forwards ``manual`` so the claim does not stamp the next occurrence;
# the default (webhook / misfire) fire keeps the occurrence stamp.
assert InProcessCronScheduler().fire_due("j1", manual=True) is True
assert claims[-1] == ("j1", {"manual": True, "return_job": True})
def test_fire_due_lost_claim_does_not_run(monkeypatch):
@@ -538,12 +538,13 @@ async def test_trigger_cron_job_fires_only_selected_job_and_returns_refreshed_st
fired = []
class RecordingProvider:
def fire_due(self, job_id, *, adapters=None, loop=None, force=False):
def fire_due(self, job_id, *, adapters=None, loop=None, force=False, manual=False):
fired.append(
{
"job_id": job_id,
"jobs_file": cron_jobs._current_cron_store().jobs_file,
"force": force,
"manual": manual,
}
)
cron_jobs.mark_job_run(job_id, success=True)
@@ -571,6 +572,8 @@ async def test_trigger_cron_job_fires_only_selected_job_and_returns_refreshed_st
"job_id": selected["id"],
"jobs_file": isolated_profiles["worker_alpha"] / "cron" / "jobs.json",
"force": False,
# Off-tick run-now: the claim must not stamp next_run_at as the occurrence (#104790).
"manual": True,
}
]
assert triggered["last_status"] == "ok"