From 3e3cfec49d2f57734a2a40aeadf7b6ae8d2b9b2e Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 16:19:44 +0530 Subject: [PATCH] fix(cron): fire a cron job's unreachable-model retry instead of re-anchoring it MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit unreachable_retry.plan_retry parks a cron job at now + ladder delay, an instant that is off the expression's lattice. The stale-cron guard on the due scan classified it as a direct schedule edit (STALE_CRON_EXPR_EDIT) and re-anchored it to the natural occurrence without firing, so the 5/15/30-minute ladder never ran for cron jobs — the same failure mode the quota-hold recovery fire had to exempt itself from. Record the ladder instant in the retry state and let _reanchor_stale_cron accept either planner's parked instant as an authorized one-shot off-lattice fire; interval jobs and legacy state without the instant keep the previous behaviour. --- cron/jobs.py | 14 ++++++++------ cron/unreachable_retry.py | 10 ++++++++-- tests/cron/test_unreachable_retry.py | 20 +++++++++++++++++--- 3 files changed, 33 insertions(+), 11 deletions(-) diff --git a/cron/jobs.py b/cron/jobs.py index 8d482645e5..31c7f2dae3 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -2988,16 +2988,18 @@ def _reanchor_stale_cron(d: _DueJob) -> bool: """Stale-schedule guard for a due cron instant; True when re-anchored without firing. A direct edit of schedule.expr leaves next_run_at on the old lattice, so re-anchor first (from - the current expr, so this converges). Two cases intentionally authorize one off-lattice fire - instead: an offset-representation migration that would otherwise swallow a never-fired - occurrence, and a quota recovery before a sparse cron's next natural occurrence. - Both fall through to fire ONCE (at-most-once holds: completion rewrites the instant).""" + the current expr, so this converges). Some instants are off the lattice on purpose and + authorize one fire instead: an offset-representation migration that would otherwise swallow + a never-fired occurrence, and the failure-path planners' parked instants (a quota recovery + before a sparse cron's next natural occurrence, an unreachable-model retry rung). All fall + through to fire ONCE (at-most-once holds: completion rewrites the instant).""" stale_class = _classify_stale_cron_next_run(d.schedule, d.raw_next_run_dt, d.next_run_dt) if stale_class == STALE_CRON_EXPR_EDIT: from cron.quota_hold import is_recovery_fire - if is_recovery_fire(d.job, d.next_run): + from cron.unreachable_retry import is_retry_fire + if is_recovery_fire(d.job, d.next_run) or is_retry_fire(d.job, d.next_run): logger.info( - "cron.quota_hold.recovery_fire job='%s' id=%s expr=%r at=%s", + "cron.off_lattice_fire job='%s' id=%s expr=%r at=%s", d.label, d.job.get("id"), d.schedule.get("expr"), d.next_run) return False new_next = d.recompute_next() diff --git a/cron/unreachable_retry.py b/cron/unreachable_retry.py index bd4e02f999..738e2cf204 100644 --- a/cron/unreachable_retry.py +++ b/cron/unreachable_retry.py @@ -31,7 +31,8 @@ logger = logging.getLogger("cron.scheduler") RETRY_DELAYS_SECONDS: tuple[int, ...] = (300, 900, 1800) # Persisted on the job while a retry cycle is active: {"attempt": <1-based count of -# retries already scheduled>}. Cleared by any run that reached the model. +# retries already scheduled>, "at": }. Cleared by any run +# that reached the model. STATE_KEY = "unreachable_retry" @@ -81,6 +82,11 @@ def clear_state(job: Dict[str, Any]) -> None: job.pop(STATE_KEY, None) +def is_retry_fire(job: Dict[str, Any], next_run: str) -> bool: + """True for the exact ladder instant parked by ``plan_retry`` (off the cron lattice).""" + return (job.get(STATE_KEY) or {}).get("at") == next_run + + def plan_retry(job: Dict[str, Any]) -> bool: """Called under the jobs lock AFTER ``_advance_after_run`` computed the schedule's natural ``next_run_at`` for a failed, flagged run. Pulls ``next_run_at`` earlier to @@ -113,7 +119,7 @@ def plan_retry(job: Dict[str, Any]) -> bool: clear_state(job) return False retry_at = retry_dt.isoformat() - job[STATE_KEY] = {"attempt": attempt + 1} + job[STATE_KEY] = {"attempt": attempt + 1, "at": retry_at} job["next_run_at"] = retry_at if job.get("state") != "paused": job["state"] = "scheduled" diff --git a/tests/cron/test_unreachable_retry.py b/tests/cron/test_unreachable_retry.py index d0a71f9504..bc0465cd95 100644 --- a/tests/cron/test_unreachable_retry.py +++ b/tests/cron/test_unreachable_retry.py @@ -11,7 +11,7 @@ from datetime import datetime, timedelta, timezone import pytest from cron import unreachable_retry as ur -from cron.jobs import create_job, get_job, mark_job_run +from cron.jobs import create_job, get_due_jobs, get_job, mark_job_run @pytest.fixture @@ -26,9 +26,13 @@ def _iso(dt: datetime) -> str: return dt.isoformat() -def test_unreachable_failure_pulls_next_run_earlier_then_ladder_exhausts(tmp_cron_home): +def test_unreachable_failure_pulls_next_run_earlier_then_ladder_exhausts( + tmp_cron_home, monkeypatch, +): """Failed-unreachable runs re-fire on the 5/15/30-minute ladder instead of waiting a - full period, and the ladder stops after its last rung (falls back to the schedule).""" + full period, and the ladder stops after its last rung (falls back to the schedule). A + cron job's ladder instant is off its lattice yet must be due, not re-anchored as a stale + expression edit.""" # Interval, not a cron expression: the natural next fire is always a full day out. A # fixed clock time ("0 3 * * *") makes the 30-minute rung land past the natural fire # in the half hour before it, and plan_retry rightly yields to the schedule (CI red). @@ -51,6 +55,16 @@ def test_unreachable_failure_pulls_next_run_earlier_then_ladder_exhausts(tmp_cro assert j.get(ur.STATE_KEY) is None assert datetime.fromisoformat(j["next_run_at"]) - now > timedelta(hours=1) + pinned = datetime(2026, 9, 18, 12, 1, tzinfo=timezone.utc) + monkeypatch.setattr("cron.jobs._hermes_now", lambda: pinned) + monkeypatch.setattr(ur, "_hermes_now", lambda: pinned) + weekly = create_job("weekly digest", "0 12 * * 5") + assert mark_job_run(weekly["id"], False, "ConnectError: dns", model_unreachable=True) + retry_at = datetime.fromisoformat(get_job(weekly["id"])["next_run_at"]) + assert retry_at == pinned + timedelta(seconds=ur.RETRY_DELAYS_SECONDS[0]) + monkeypatch.setattr("cron.jobs._hermes_now", lambda: retry_at + timedelta(seconds=1)) + assert weekly["id"] in {due["id"] for due in get_due_jobs()} + def test_reaching_the_model_resets_ladder_and_oneshots_never_retry(tmp_cron_home): """Any run that reached the model clears retry state; one-shots (pre-claimed