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:
+2
-1
@@ -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."""
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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))
|
||||
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user