fix(cron): measure hold, retry, resume and alert windows in real time across the fall-back hour

Same class as the due-scan fix: the configured zone is one cached ZoneInfo, so two aware
datetimes in it compare and subtract by wall clock, and aware + timedelta is wall-clock
arithmetic that drops fold. Inside the repeated 01:xx hour:

- unreachable_retry.plan_retry: now(01:30 EST) + 300s landed on 01:35 EDT, 55 min in the
  past, so the "re-run in 5 min" fired on the next tick.
- quota_hold: the provider-window end was an hour off (hold too long, or already expired)
  and hold_active compared wall clocks.
- resume_job: the "slot elapsed while paused" check compared wall clocks and could drop
  the kept slot (#113603 path).
- _repeat_alert_withheld: the repeat-alert window was measured in wall time.

Add jobs._seconds_after (UTC add, back to the caller's zone) next to the instant helpers
and use the instant comparisons there. Found by the C13 virtual-clock cron soak.
This commit is contained in:
teknium1
2026-09-23 07:22:22 -07:00
committed by Teknium
parent e3618cdec2
commit e85291e27b
4 changed files with 27 additions and 15 deletions

View File

@@ -875,6 +875,12 @@ def _instant_before(left: datetime, right: datetime) -> bool:
return left.astimezone(timezone.utc) < right.astimezone(timezone.utc)
def _seconds_after(dt: datetime, seconds: float) -> datetime:
"""*dt* plus real *seconds*, in *dt*'s zone. Aware ``+ timedelta`` is wall-clock arithmetic
that drops ``fold``, so inside a fall-back hour it lands an hour off."""
return (dt.astimezone(timezone.utc) + timedelta(seconds=seconds)).astimezone(dt.tzinfo)
def _parse_aware(value: Any) -> Optional[datetime]:
"""``_ensure_aware(datetime.fromisoformat(value))``, or None when *value* is not a parseable ISO
string."""
@@ -2088,7 +2094,7 @@ def resume_job(job_id: str) -> Optional[Dict[str, Any]]:
if (
job["schedule"].get("kind") in {"cron", "interval"}
and stored_dt is not None
and stored_dt <= _hermes_now()
and _instant_at_or_before(stored_dt, _hermes_now())
):
next_run_at = stored_next
logger.info(

View File

@@ -16,7 +16,7 @@ from __future__ import annotations
import logging
import re
from datetime import datetime, timedelta
from datetime import datetime
from typing import Any, Dict, Optional
from hermes_time import now as _hermes_now
@@ -56,29 +56,35 @@ def hold_seconds_from_failure(exc: BaseException) -> Optional[float]:
def hold_active(job: Dict[str, Any], now: Optional[datetime] = None) -> bool:
"""True while the job is parked inside a provider window (an expired marker is inert)."""
from cron.jobs import _parse_aware # late: jobs imports this module's helpers
from cron.jobs import _instant_after, _parse_aware # late: jobs imports this module's helpers
until = _parse_aware(job.get(STATE_KEY)) if job.get(STATE_KEY) else None
return until is not None and until > (now or _hermes_now())
return until is not None and _instant_after(until, now or _hermes_now())
def clear_state(job: Dict[str, Any]) -> None:
job.pop(STATE_KEY, None)
def _window_end(hold_seconds: float) -> datetime:
from cron.jobs import _seconds_after
return _seconds_after(_hermes_now(), float(hold_seconds) + HOLD_SLACK_SECONDS)
def plan_hold(job: Dict[str, Any], hold_seconds: float) -> bool:
"""Called under the jobs lock AFTER ``_advance_after_run`` computed the schedule's natural
``next_run_at`` for a failed run. Parks a recurring job at its first occurrence after the
provider window when that is later than the natural one. Returns True when parked."""
from cron.jobs import _parse_aware, compute_next_run
from cron.jobs import _instant_before, _parse_aware, compute_next_run
schedule = job.get("schedule") or {}
if schedule.get("kind") not in {"cron", "interval"} or job.get("state") == "paused":
clear_state(job)
return False
window_end = _hermes_now() + timedelta(seconds=float(hold_seconds) + HOLD_SLACK_SECONDS)
window_end = _window_end(hold_seconds)
natural_next = _parse_aware(job.get("next_run_at"))
if natural_next is not None and natural_next >= window_end:
if natural_next is not None and not _instant_before(natural_next, window_end):
clear_state(job)
return False
if schedule.get("kind") == "interval":
@@ -100,7 +106,7 @@ def hold_notice(job: Dict[str, Any], hold_seconds: Optional[float]) -> str:
"""Line appended to the ONE failure alert delivered on entering the hold, else ""."""
if not hold_seconds or (job.get("schedule") or {}).get("kind") not in {"cron", "interval"}:
return ""
window_end = _hermes_now() + timedelta(seconds=float(hold_seconds) + HOLD_SLACK_SECONDS)
window_end = _window_end(hold_seconds)
hours = float(hold_seconds) / 3600.0
return (
f"\nThe provider's usage window is closed for about {hours:.1f}h. This job is held "

View File

@@ -17,7 +17,7 @@ import threading
import time
import uuid
from dataclasses import dataclass
from datetime import datetime, timedelta, timezone
from datetime import datetime, timezone
# fcntl is Unix-only; Windows uses msvcrt
try:
@@ -375,12 +375,12 @@ def _repeat_alert_withheld(incident: dict) -> bool:
if not alerted_at:
return False
try:
from cron.jobs import _ensure_aware
from cron.jobs import _elapsed_seconds, _ensure_aware
last = _ensure_aware(datetime.fromisoformat(str(alerted_at)))
except (TypeError, ValueError):
return False
return _hermes_now() - last < timedelta(hours=hours)
return _elapsed_seconds(_hermes_now(), last) < hours * 3600
def _upsert_incident_for_failure(

View File

@@ -20,7 +20,6 @@ run that reaches the model — success or not — resets the ladder. Disable wit
from __future__ import annotations
import logging
from datetime import timedelta
from typing import Any, Dict, Optional
from hermes_time import now as _hermes_now
@@ -103,11 +102,12 @@ def plan_retry(job: Dict[str, Any]) -> bool:
job.get("name", job.get("id", "?")), attempt, job.get("next_run_at"))
return False
delay = RETRY_DELAYS_SECONDS[attempt]
retry_dt = _hermes_now() + timedelta(seconds=delay)
from cron.jobs import _parse_aware # late: jobs imports this module's helpers
# late: jobs imports this module's helpers
from cron.jobs import _instant_at_or_before, _parse_aware, _seconds_after
retry_dt = _seconds_after(_hermes_now(), delay)
natural_next = _parse_aware(job.get("next_run_at"))
if natural_next is not None and natural_next <= retry_dt:
if natural_next is not None and _instant_at_or_before(natural_next, retry_dt):
# The schedule fires again sooner than the ladder would — no point consuming an
# attempt; the natural occurrence IS the retry.
clear_state(job)