refactor(cron): simplify quota-hold recovery plumbing
_recovery_worthwhile subtracted same-tz aware datetimes directly, which is wall-clock arithmetic and off by an hour when the hold-to-occurrence span crosses a DST transition (weekly jobs make that reachable); use the module's canonical _elapsed_seconds. Derive the schedule from the job and gate the "scheduled fire" boolean at the call site under one name (recover_consumed_fire) from scheduler -> mark_job_run -> plan_hold instead of three. Collapse plan_hold's interval/cron arms, which duplicated the not-blocked early return and the boundary park. Persist the expr fingerprint only for the off-lattice recovery fire: the coalesce branch parks on the lattice (STALE_CRON_MATCH, also under the croniter-missing fallback), so is_recovery_fire is never consulted there. Drop a comment restating the docstring.
This commit is contained in:
@@ -2378,7 +2378,7 @@ def mark_job_run(
|
||||
expected_fire_owner: Optional[str] = None,
|
||||
model_unreachable: bool = False,
|
||||
quota_hold_seconds: Optional[float] = None,
|
||||
quota_recover_occurrence: bool = False,
|
||||
recover_consumed_fire: bool = False,
|
||||
) -> bool:
|
||||
"""Mark a job as run: update last_run_at/last_status, bump completed, recompute next_run_at,
|
||||
and retire the record as a terminal completion when the repeat limit is reached.
|
||||
@@ -2395,7 +2395,7 @@ def mark_job_run(
|
||||
|
||||
``quota_hold_seconds``: the provider said it stays closed for this long (a quota 429 with
|
||||
``retry after <N>s``). Recurring jobs are parked through the window instead of re-firing into
|
||||
it on every tick. ``quota_recover_occurrence`` lets a scheduled sparse cron recover its
|
||||
it on every tick. ``recover_consumed_fire`` lets a scheduled sparse cron recover its
|
||||
consumed fire when the provider reopens; manual runs retain the natural schedule
|
||||
(cron/quota_hold.py, #89376).
|
||||
"""
|
||||
@@ -2419,8 +2419,7 @@ def mark_job_run(
|
||||
# Any run that reached the model (either outcome) resets the re-run ladder.
|
||||
clear_state(job)
|
||||
if not success and quota_hold_seconds and not is_terminal_job(job):
|
||||
quota_hold.plan_hold(
|
||||
job, quota_hold_seconds, recover_consumed_fire=quota_recover_occurrence)
|
||||
quota_hold.plan_hold(job, quota_hold_seconds, recover_consumed_fire=recover_consumed_fire)
|
||||
else:
|
||||
quota_hold.clear_state(job)
|
||||
save_jobs(jobs)
|
||||
@@ -2997,8 +2996,6 @@ def _reanchor_stale_cron(d: _DueJob) -> bool:
|
||||
if stale_class == STALE_CRON_EXPR_EDIT:
|
||||
from cron.quota_hold import is_recovery_fire
|
||||
if is_recovery_fire(d.job, d.next_run):
|
||||
# plan_hold deliberately creates one off-lattice recovery fire when a provider
|
||||
# reopens before a sparse cron's next natural occurrence.
|
||||
logger.info(
|
||||
"cron.quota_hold.recovery_fire job='%s' id=%s expr=%r at=%s",
|
||||
d.label, d.job.get("id"), d.schedule.get("expr"), d.next_run)
|
||||
|
||||
@@ -90,8 +90,7 @@ def _window_end(hold_seconds: float) -> datetime:
|
||||
|
||||
|
||||
def _recovery_worthwhile(
|
||||
job: Dict[str, Any], schedule: Dict[str, Any], natural_next: Optional[datetime],
|
||||
window_end: datetime, scheduled: bool,
|
||||
job: Dict[str, Any], natural_next: Optional[datetime], window_end: datetime,
|
||||
) -> bool:
|
||||
"""One off-lattice recovery fire, and only for a sparse schedule.
|
||||
|
||||
@@ -102,12 +101,12 @@ def _recovery_worthwhile(
|
||||
``cron.jobs._compute_grace_seconds`` — otherwise the recovery fire is a near-duplicate of the
|
||||
natural one (hourly job, hold ending at :58, would fire :58 AND :00).
|
||||
"""
|
||||
from cron.jobs import _schedule_cadence_seconds
|
||||
from cron.jobs import _elapsed_seconds, _schedule_cadence_seconds
|
||||
|
||||
if not scheduled or job.get(STATE_KEY) or natural_next is None:
|
||||
if job.get(STATE_KEY) or natural_next is None:
|
||||
return False
|
||||
cadence = _schedule_cadence_seconds(schedule)
|
||||
return bool(cadence) and (natural_next - window_end).total_seconds() >= cadence / 2
|
||||
cadence = _schedule_cadence_seconds(job.get("schedule") or {})
|
||||
return bool(cadence) and _elapsed_seconds(natural_next, window_end) >= cadence / 2
|
||||
|
||||
|
||||
def plan_hold(
|
||||
@@ -127,23 +126,21 @@ def plan_hold(
|
||||
window_end = _window_end(hold_seconds)
|
||||
natural_next = _parse_aware(job.get("next_run_at"))
|
||||
blocked = natural_next is None or _instant_before(natural_next, window_end)
|
||||
if kind == "interval":
|
||||
if not blocked:
|
||||
# No interval occurrence is blocked: its natural next run is already after recovery.
|
||||
clear_state(job)
|
||||
return False
|
||||
parked = window_end.isoformat()
|
||||
job.pop(SCHEDULE_EXPR_KEY, None)
|
||||
recover = (kind == "cron" and not blocked and recover_consumed_fire
|
||||
and _recovery_worthwhile(job, natural_next, window_end))
|
||||
if not blocked and not recover:
|
||||
clear_state(job)
|
||||
return False
|
||||
if kind == "cron" and blocked:
|
||||
# Coalesce cron occurrences inside the closed window to the first legal instant after it.
|
||||
parked = compute_next_run(schedule, window_end.isoformat()) or window_end.isoformat()
|
||||
else:
|
||||
if not blocked:
|
||||
if not _recovery_worthwhile(job, schedule, natural_next, window_end, recover_consumed_fire):
|
||||
clear_state(job)
|
||||
return False
|
||||
parked = window_end.isoformat()
|
||||
else:
|
||||
# Coalesce cron occurrences inside the closed window to the first legal instant after it.
|
||||
parked = compute_next_run(schedule, window_end.isoformat()) or window_end.isoformat()
|
||||
parked = window_end.isoformat()
|
||||
if recover:
|
||||
# Only the recovery fire is off-lattice; the coalesced instant is a legal occurrence.
|
||||
job[SCHEDULE_EXPR_KEY] = schedule.get("expr")
|
||||
else:
|
||||
job.pop(SCHEDULE_EXPR_KEY, None)
|
||||
job["next_run_at"] = parked
|
||||
job[STATE_KEY] = parked
|
||||
logger.warning(
|
||||
|
||||
@@ -3059,7 +3059,7 @@ def _finish_completed_run(d: _RunDelivery, fire_owner: Optional[str], execution_
|
||||
if not d.success and _hold_s:
|
||||
# Provider window closed for a known duration: park past it (cron/quota_hold.py, #89376).
|
||||
mark_kwargs["quota_hold_seconds"] = _hold_s
|
||||
mark_kwargs["quota_recover_occurrence"] = bool(job.get("_scheduled_instant"))
|
||||
mark_kwargs["recover_consumed_fire"] = bool(job.get("_scheduled_instant"))
|
||||
if d.success and not d.delivery_error and d.should_deliver and job.get("last_delivery_queued"):
|
||||
mark_kwargs["status"] = "delivery_queued"
|
||||
if fire_owner is not None:
|
||||
|
||||
@@ -66,7 +66,7 @@ def test_weekly_cron_retries_when_quota_recovers_before_next_occurrence(
|
||||
|
||||
assert mark_job_run(
|
||||
job["id"], False, QUOTA_MSG, quota_hold_seconds=20 * 60 * 60,
|
||||
quota_recover_occurrence=True,
|
||||
recover_consumed_fire=True,
|
||||
)
|
||||
|
||||
held = get_job(job["id"])
|
||||
|
||||
Reference in New Issue
Block a user