fix(scheduler): defer уважает interval_minutes источника
Some checks failed
CI Trade-In / changes (pull_request) Successful in 14s
CI / changes (pull_request) Successful in 16s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Failing after 5m49s

_defer_next_run_at (путь пропуска pre_claim — например протухшие куки Циана)
читал только interval_days, поэтому cian_detail_backfill с каденцией 360 мин
после пропуска получал случайный час завтрашних суток: от ~1 до ~47 часов
вместо шести. Ручное обновление кук не давало эффекта ещё сутки.

Разбор interval_minutes вынесен в _interval_minutes и переиспользован
post_claim-хуком reschedule_after_minutes — два входа в одну каденцию.

Нижний порог _MIN_DEFER_MINUTES = 3 тика планировщика сохраняет #1522:
defer обязан пережить несколько get_due_schedules, иначе pre-check снова
гоняется каждую минуту.

Refs #3312, #1522
This commit is contained in:
bot-backend 2026-09-05 22:50:58 +05:00
parent 99f112db0b
commit 63e2210abd
2 changed files with 122 additions and 6 deletions

View file

@ -0,0 +1,91 @@
"""Пропуск по pre_claim не теряет sub-hourly каденцию источника (#3312).
Баг: _defer_next_run_at читал только interval_days, поэтому cian_detail_backfill
(interval_minutes=360, окно 0..23, ключа interval_days нет) после пропуска по протухшей
куке получал случайный час ЗАВТРАШНИХ суток от ~1 до ~47 часов вместо шести. Ручное
обновление кук не давало эффекта ещё сутки.
Осторожно с #1522: defer обязан пережить несколько тиков планировщика, иначе pre-check
снова гоняется раз в минуту отсюда нижний порог _MIN_DEFER_MINUTES.
Без сети, без БД: _defer_next_run_at вызывается напрямую с mock-сессией, проверяется
bind-param next_at.
"""
from __future__ import annotations
import os
from datetime import UTC, datetime, timedelta
from typing import Any
from unittest.mock import MagicMock
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
from scraper_kit.orchestration.scheduler import (
SCHEDULER_TICK_SEC,
_defer_next_run_at,
compute_next_run_at,
)
# Порог #1522 читаем НЕ из фиксируемой константы, а выводим из такта планировщика: иначе
# откат правки даёт ImportError («возможности нет») вместо красного ПО ЗНАЧЕНИЮ.
_MIN_DEFER = timedelta(seconds=3 * SCHEDULER_TICK_SEC)
def _row(source: str, params: dict[str, Any]) -> dict[str, Any]:
return {
"source": source,
"window_start_hour": 0,
"window_end_hour": 23,
"default_params": params,
}
def _deferred_to(params: dict[str, Any], source: str = "cian_detail_backfill") -> datetime:
"""next_at из bind-params единственного UPDATE в _defer_next_run_at."""
db = MagicMock()
_defer_next_run_at(db, _row(source, params))
return db.execute.call_args_list[0].args[1]["next_at"]
def test_sub_hourly_source_defers_by_its_own_cadence_not_tomorrow() -> None:
"""interval_minutes=360 → next_run_at ≈ now+6ч. Падает на старом коде (завтра)."""
before = datetime.now(tz=UTC)
next_at = _deferred_to({"batch_size": 400, "interval_minutes": 360})
delta = next_at - before
assert timedelta(minutes=359) <= delta <= timedelta(minutes=361), (
f"ожидали ~360 мин, получили {delta}"
)
# Гарантия «не завтра»: суточный defer в окне 0..23 даёт минимум +1ч и в среднем сутки.
assert delta < timedelta(hours=7)
def test_daily_source_keeps_compute_next_run_at_behaviour() -> None:
"""Без interval_minutes — прежняя формула: та же дата и окно, что у compute_next_run_at."""
next_at = _deferred_to({"batch_size": 400}, source="avito_full_load")
expected = compute_next_run_at(0, 23, interval_days=1)
assert next_at.date() == expected.date()
assert 0 <= next_at.hour < 23
assert next_at.tzinfo is not None
def test_daily_source_honours_interval_days() -> None:
"""interval_days=7 (недельный источник) продолжает работать как раньше."""
next_at = _deferred_to({"interval_days": 7}, source="avito_full_load")
assert next_at.date() == compute_next_run_at(0, 23, interval_days=7).date()
def test_tiny_interval_clamped_to_min_defer_1522() -> None:
"""interval_minutes=1 → ровно порог #1522 (3 тика), а не 1 минута и не завтра."""
before = datetime.now(tz=UTC)
next_at = _deferred_to({"interval_minutes": 1})
delta = next_at - before
assert _MIN_DEFER <= delta <= _MIN_DEFER + timedelta(seconds=5), (
f"ожидали порог {_MIN_DEFER}, получили {delta}"
)

View file

@ -696,20 +696,43 @@ def _pick_resume(db: Session, run_id: int) -> int | None:
return resume_run_id
def _interval_minutes(params: dict[str, Any], *, param: str = "interval_minutes") -> int | None:
"""Sub-hourly каденс источника из default_params; None — ключа нет (суточный источник).
Общий разбор для двух входов в одну каденцию: post_claim-хука
(reschedule_after_minutes) и пути пропуска (_defer_next_run_at, #3312).
"""
raw = params.get(param)
return None if raw is None else int(raw)
# Нижний порог defer'а (#1522): пропуск обязан пережить несколько get_due_schedules,
# иначе pre-check снова гоняется каждый тик. 3 тика = минимум один промах даже при
# рассинхроне часов планировщика и БД на такт в любую сторону.
_MIN_DEFER_MINUTES = 3 * SCHEDULER_TICK_SEC // 60
def _defer_next_run_at(db: Session, schedule_row: dict[str, Any]) -> None:
"""Сдвинуть next_run_at на следующее окно БЕЗ создания run (#1522).
Используется когда pre_claim делает early-return (например cian: отсутствуют/протухли
cookies) ДО _claim_run. Без этого get_due_schedules переотбирает schedule на каждом
тике (SCHEDULER_TICK_SEC), и pre-check гоняется раз в минуту круглосуточно.
Каденцию берём ту же, что post_claim-хук: есть interval_minutes now + N минут
(не ниже _MIN_DEFER_MINUTES), иначе прежнее суточное окно по interval_days (#3312).
"""
source = schedule_row["source"]
params = schedule_row.get("default_params") or {}
next_at = compute_next_run_at(
schedule_row["window_start_hour"],
schedule_row["window_end_hour"],
interval_days=int(params.get("interval_days", 1)),
)
minutes = _interval_minutes(params)
if minutes is not None:
next_at = datetime.now(tz=UTC) + timedelta(minutes=max(minutes, _MIN_DEFER_MINUTES))
else:
next_at = compute_next_run_at(
schedule_row["window_start_hour"],
schedule_row["window_end_hour"],
interval_days=int(params.get("interval_days", 1)),
)
db.execute(
text(
"""
@ -742,7 +765,9 @@ def reschedule_after_minutes(*, param: str = "interval_minutes", default: int =
# source читаем из ЗАКЛЕЙМЛЕННОГО schedule — берём из params-owner через closure нельзя,
# поэтому апдейтим по run_id → source (schedule owner). Проще: UPDATE по source из
# scrape_runs текущего run_id, что эквивалентно WHERE source = <this source>.
minutes = int(params.get(param, default))
minutes = _interval_minutes(params, param=param)
if minutes is None:
minutes = default
db.execute(
text(
"""