diff --git a/cron/jobs.py b/cron/jobs.py index 195fddba62..903021c38c 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -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.""" diff --git a/cron/scheduler_provider.py b/cron/scheduler_provider.py index 6f13778382..e008e3223d 100644 --- a/cron/scheduler_provider.py +++ b/cron/scheduler_provider.py @@ -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 diff --git a/hermes_cli/web_server_cron.py b/hermes_cli/web_server_cron.py index fb0c5ddff5..e894a23494 100644 --- a/hermes_cli/web_server_cron.py +++ b/hermes_cli/web_server_cron.py @@ -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)) diff --git a/tests/cron/test_claim_job_for_fire.py b/tests/cron/test_claim_job_for_fire.py index da6382cb5f..aa7f351135 100644 --- a/tests/cron/test_claim_job_for_fire.py +++ b/tests/cron/test_claim_job_for_fire.py @@ -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) diff --git a/tests/cron/test_scheduler_provider.py b/tests/cron/test_scheduler_provider.py index 81e208c064..f0f22dcb34 100644 --- a/tests/cron/test_scheduler_provider.py +++ b/tests/cron/test_scheduler_provider.py @@ -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): diff --git a/tests/hermes_cli/test_web_server_cron_profiles.py b/tests/hermes_cli/test_web_server_cron_profiles.py index 1be4ee2818..9e345c5faf 100644 --- a/tests/hermes_cli/test_web_server_cron_profiles.py +++ b/tests/hermes_cli/test_web_server_cron_profiles.py @@ -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"