fix(scraper-kit): пропуск по куке двигает next_run_at на interval_minutes, а не на случайный час завтра #3351
2 changed files with 122 additions and 6 deletions
|
|
@ -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}"
|
||||
)
|
||||
|
|
@ -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(
|
||||
"""
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue