From 63e2210abdb9f478a0be0e59704b90d1d37ab05e Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 5 Sep 2026 22:50:58 +0500 Subject: [PATCH] =?UTF-8?q?fix(scheduler):=20defer=20=D1=83=D0=B2=D0=B0?= =?UTF-8?q?=D0=B6=D0=B0=D0=B5=D1=82=20interval=5Fminutes=20=D0=B8=D1=81?= =?UTF-8?q?=D1=82=D0=BE=D1=87=D0=BD=D0=B8=D0=BA=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit _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 --- ...st_defer_respects_interval_minutes_3312.py | 91 +++++++++++++++++++ .../scraper_kit/orchestration/scheduler.py | 37 ++++++-- 2 files changed, 122 insertions(+), 6 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_defer_respects_interval_minutes_3312.py diff --git a/tradein-mvp/backend/tests/test_defer_respects_interval_minutes_3312.py b/tradein-mvp/backend/tests/test_defer_respects_interval_minutes_3312.py new file mode 100644 index 00000000..f4c3c265 --- /dev/null +++ b/tradein-mvp/backend/tests/test_defer_respects_interval_minutes_3312.py @@ -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}" + ) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py index d49712c7..7454d657 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py @@ -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 = . - minutes = int(params.get(param, default)) + minutes = _interval_minutes(params, param=param) + if minutes is None: + minutes = default db.execute( text( """