diff --git a/cron/jobs.py b/cron/jobs.py index 84bb8fefbf..20e3ae5a8f 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -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( diff --git a/cron/quota_hold.py b/cron/quota_hold.py index 0d031a4e76..f4c7a690d2 100644 --- a/cron/quota_hold.py +++ b/cron/quota_hold.py @@ -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 " diff --git a/cron/scheduler.py b/cron/scheduler.py index c70b868a53..f1c41948f6 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -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( diff --git a/cron/unreachable_retry.py b/cron/unreachable_retry.py index 924a826a52..bd4e02f999 100644 --- a/cron/unreachable_retry.py +++ b/cron/unreachable_retry.py @@ -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)