gendesign/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py
bot-backend a7362bc5fa
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 9s
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) Successful in 4m56s
fix(#3391): пульс не пишет по финализированной строке; отмена в теле тика тоже помечает; rollback/try/честный лог
Пять замечаний deep-ревью к PR #3392, ровно они.

1. Пульс стирал метку дрейна. `update_heartbeat` обеих копий бил `WHERE id = :run_id`
   без гейта по статусу, а app-копия counters ЗАМЕНЯЕТ (#3390): задача, помеченная
   `interrupted`, но ещё живая (ветка таймаута drain_inflight отдаёт её внешнему
   hard-cancel'у — несколько итераций спустя, пульс на каждый батч —
   app/services/scheduler.py:141), следующим же ударом стирала метку, и оборванный
   прогон снова читался как полный проход. Гейт — `IN ('running', 'cancelled')`, а не
   `= 'running'`: 'cancelled' финализирует строку, но задача встаёт лишь на ближайшей
   границе якоря, и её последний пульс — ЕДИНСТВЕННЫЙ писатель чекпоинта в этот момент
   (pipeline.py:1308/2368/2986/4488, mark_done там уже no-op по своему гейту), а
   'cancelled' входит в _RESUME_STATUSES — сужение до 'running' молча съело бы точку
   возобновления у каждой отмены. Возвращаемое значение update_heartbeat не читает
   никто (обе копии -> None, ни одного присваивания на 130 сайтах вызова), так что
   «0 строк обновлено» ломать нечего; no-op логируется WARNING'ом, как у mark_done.

2. Отмена вне drain_inflight. Hard-cancel приходит по расписанию grace'а
   scheduler_main, а не по нашему, и может застать ТЕЛО тика (reap / stale-digest /
   `_dispatch` с сетевым pre_claim). `except Exception` тика CancelledError не ловит,
   до `await ctx.drain_inflight()` дело не доходит — строки оставались 'running'.
   Тело вынесено в `_tick_loop`, `scheduler_loop` ловит CancelledError, помечает
   in-flight и пробрасывает отмену.

3. rollback в except пометки: отказавший statement оставляет сессию в aborted-tx, и
   первый же непроходимый run_id утаскивал все следующие (образец — defensive rollback
   в mark_failed/mark_banned).

4. session_factory()/db.close() втянуты в try: исключение оттуда ЗАМЕНИЛО бы собой
   CancelledError, а suppress(CancelledError) в scheduler_main его не глушит — процесс
   уходил бы с трейсбеком вместо чистого drain-выхода.

5. WARNING перечисляет marked_ids, а не весь run_ids (там были и пропущенные по
   статусу). В докстринге назван потолок: SELECT синхронный, у движка нет ни connect-,
   ни statement-таймаута (app/core/db.py:8-19) — недоступная БД блокирует луп до
   SIGKILL'а через 20 с docker-grace; данные при этом не хуже прежних (строки остаются
   'running' → boot-reap).

Тесты — по значению, не по факту вызова; на исходниках 30e3bacc краснеют все пять:
'listings_processed' дописан в финализированную строку (kit), KeyError: 'interrupted'
(app), assert 'running' == 'done' (отмена в теле тика), assert 0 == 1 (второй run_id
не помечен после отказа первого), RuntimeError наружу (недоступная БД).
2026-09-06 09:03:56 +05:00

1403 lines
79 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""In-app scheduler — strangler-копия `app.services.scheduler` в scraper_kit (#2136).
Разница с боевым `app.services.scheduler`:
1. Полная развязка от `app.*` — все продуктовые зависимости инжектируются через
`SchedulerContext` (config/matcher/enrichment/session_factory/runs/shutdown),
описанные Protocol'ами в `scraper_kit.contracts`. `grep "from app"` — пусто.
2. **Registry-dispatch вместо copy-paste.** Боевой модуль держит 25+ почти
одинаковых `trigger_*_run` (claim → new-session → create_task → close) плюс
27-веточный `if/elif` в `scheduler_loop`. Здесь — единственный параметризованный
путь `_dispatch()` + `HANDLERS: dict[str, Handler]`. Каждый source — одна запись
реестра; общий claim/session/drain-boilerplate вынесен в `_dispatch`.
Kit-native (sweep-оркестраторы, уже перенесённые в `scraper_kit.orchestration.pipeline`)
регистрируются встроенно (`_default_kit_handlers`). Продуктовые джобы, тело которых
осталось в `app` (rosreestr_dkp / sber_index / deactivate_stale / *_backfill / …),
инжектируются извне как `Handler` через `build_registry(product_handlers=...)`.
Боевой рантайм (scraper-контейнер, `python -m app.scheduler_main`) крутит ИМЕННО ЭТОТ
loop: #2397 Part C удалил legacy-ветку `app.services.scheduler.scheduler_loop`, и
`_run_kit_scheduler` остался единственным путём. Строка «по-прежнему крутит старый
app.services.scheduler» жила здесь после того, как перестала быть правдой, и посылала
правку сторожей не в тот файл (#2670).
Критичная concurrency-логика (`_claim_run` advisory-lock + double-check, `reap_zombies`
порог, heartbeat, SIGTERM-drain) перенесена ДОСЛОВНО — тот же SQL, то же ветвление.
"""
from __future__ import annotations
import asyncio
import logging
import random
from collections.abc import Callable, Coroutine, Iterable, Mapping
from contextlib import suppress
from dataclasses import dataclass, field
from datetime import UTC, datetime, time, timedelta
from typing import TYPE_CHECKING, Any
from sqlalchemy import text
try: # sentry опционален — kit standalone-импортируем, sentry-sdk не в его зависимостях
import sentry_sdk
except ImportError: # pragma: no cover - в проде backend-env sentry_sdk присутствует
sentry_sdk = None # type: ignore[assignment]
from scraper_kit.orchestration import runs as _kit_runs
from scraper_kit.orchestration.pipeline import (
get_city_anchors,
run_avito_city_sweep,
run_avito_full_load,
run_avito_newbuilding_sweep,
run_cian_city_sweep,
run_cian_full_load,
run_domclick_city_sweep,
run_yandex_city_sweep,
)
if TYPE_CHECKING:
from sqlalchemy.orm import Session
from scraper_kit.contracts import (
EnrichmentJobs,
HouseMatcher,
ProxyProvider,
ScraperConfig,
SessionFactory,
)
logger = logging.getLogger(__name__)
# Loop interval — check каждую минуту
SCHEDULER_TICK_SEC = 60
ZOMBIE_THRESHOLD_HOURS = 6
# Бюджет дренажа детей < scheduler_main._DRAIN_TIMEOUT_S (100s) < docker grace (120s)
# (см. #1182 P2 — идентично боевому scheduler'у).
_CHILD_DRAIN_TIMEOUT_S = 80.0
# Машиночитаемые причины пропуска расписания (#2658) — пишутся в scrape_runs.error
# строки со status='skipped'. Слаг, а не человеческий текст: по нему «уже бежит»
# отличается от «нет кук» (продуктовые причины — в app.services.product_handlers)
# запросом, а не грепом логов.
SKIP_ALREADY_RUNNING = "already_running"
SKIP_CONCURRENT_CLAIM = "concurrent_claim"
SKIP_RUNNING_UNDER_LOCK = "running_appeared_under_lock"
SKIP_UNKNOWN_SOURCE = "unknown_source"
# ── сводка «что сейчас не собирает» (#2670, второй пункт задачи) ─────────────
# Лестница напоминаний из #2720 считает ПОДРЯД ИДУЩИЕ неудачные ПРОГОНЫ. Прод
# 2026-08-10: шесть источников не имели успешного прогона дольше 3× своего такта, и
# трое худших из них лестнице недоступны ПО ПОСТРОЕНИЮ, а не из-за редких вех:
#
# cian_history_backfill 42.1 сут без успеха, стрик 0 — с 30.06 прогонов нет
# вовсе (сегодняшний единственный — 'skipped',
# cian_cookies_expired), а лестница шагает только по
# завершённым failed/banned;
# avito_full_load_exhaustive 49.5 сут без успеха, стрик 0 — 5 банов подряд обнулил
# один 'cancelled' 09.08 (деплой убил бегущий прогон);
# avito_full_load 37.7 сут без успеха, стрик 31 — веха 48 при такте
# interval_days=7 наступит через 17 прогонов ≈ 119 суток.
#
# Уплотнение вех чинит только третий случай: у первых двух стрик равен нулю, уплотнять
# нечего. Поэтому сводка не «ещё один сторож помельче», а ЕДИНСТВЕННЫЙ ответ на вопрос
# «что сломано сейчас»: она считает КАЛЕНДАРНЫЙ возраст последнего успеха, поэтому
# видит и молчащий источник, и обнулённый стрик, и редкую веху. Лестница остаётся как
# была — она отвечает на другой вопрос («что сломалось только что») и стоит дёшево.
STALE_DIGEST_INTERVAL_FACTOR = 3
STALE_DIGEST_PERIOD_H = 24
# Завершённые прогоны каждого включённого расписания — «принёс ли прогон данные»
# решается в Python (`run_brought_data`), а не статусом в WHERE (#3172).
#
# Раньше здесь стоял `status = 'done'`, и свежесть источника равнялась возрасту
# последнего прогона с этим статусом. Но прогон, поймавший блок, честно финализируется
# как 'banned' (#2657) — И ПРИ ЭТОМ ВСТАВЛЯЕТ СТРОКИ: у domclick_city_sweep 886 строк
# 25.08 и 128 строк 23.08, оба прогона 'banned'. Источник, который регулярно ловит блок
# и столь же регулярно приносит данные, числился мёртвым навсегда → ложный P1 #3118
# («домклик не собирается с 5 августа», хотя сбор шёл).
#
# ponytail: полный проход по завершённым прогонам включённых источников (порядка 10k
# строк на проде) — раз в STALE_DIGEST_PERIOD_H часов. Понадобится дешевле — оконный
# фильтр по finished_at, но тогда never_ok перестанет быть честным («не собирал НИ
# РАЗУ» превратится в «не собирал в окне»).
_STALE_SOURCES_SQL = text("""
SELECT sch.source,
sch.default_params->>'interval_days' AS interval_days,
sch.created_at,
r.finished_at,
r.status,
r.counters
FROM scrape_schedules sch
LEFT JOIN scrape_runs r
ON r.source = sch.source AND r.finished_at IS NOT NULL
WHERE sch.enabled
""")
# ponytail: последний выпуск сводки помнится В ПАМЯТИ процесса, поэтому рестарт
# scheduler'а (деплой) даёт лишний выпуск. Осознанный размен: альтернатива — таблица
# состояния (миграция) ради анти-спама у механизма, который и заводится ПРОТИВ
# молчания. Понадобится точность — переносить в scrape_runs строкой своего source'а.
_last_stale_digest_at: datetime | None = None
@dataclass(frozen=True)
class StaleSource:
"""Источник, не собиравший дольше STALE_DIGEST_INTERVAL_FACTOR× своего такта."""
source: str
interval_days: int
age_days: float
never_ok: bool
def _schedule_interval_days(raw: Any) -> int:
"""default_params.interval_days → такт в сутках; всё непонятное → 1 (как у claim'а).
Тот же дефолт, что у `compute_next_run_at` (interval_days=1 == daily): порог сводки
обязан считаться из ТОГО ЖЕ числа, которым расписание себя двигает, иначе «просрочен»
будет мерить не тот такт. `"interval_days": null` в jsonb приезжает сюда None.
"""
try:
return max(1, int(raw))
except (TypeError, ValueError):
return 1
def run_brought_data(status: str | None, counters: Mapping[str, Any] | None) -> bool:
"""Дал ли завершённый прогон данные — мера свежести источника (#3172).
Статус НЕ участвует, пока результат измерим: прогон с блоком ('banned', #2657)
успевает записать выдачу так же, как штатный 'done'. Меряем тем же результатным
словарём, которым уже судит сторож нулевого результата (`runs._RESULT_COUNTER_KEYS`,
#2703) — одна мера «сколько объявлений отдала выдача» на оба механизма, а не второй
список ключей, который разъедется с первым.
НЕ `lots_inserted`: это НОВИЗНА, а не наличие данных. Здоровый sweep, у которого вся
выдача уже в базе, вставляет ноль строк — по такой мере живой источник читался бы
мёртвым, то есть ровно дефект #3118 с другой стороны.
`_run_result_count() is None` — прогон результат НЕ СООБЩИЛ (28 источников на проде
не имеют результатного ключа вовсе: refresh_search_matview, deactivate_stale_*,
мониторы). Судить нечем → зачитываем прежнюю меру, успешный статус. Иначе «не
измерено» схлопнулось бы с «измерено, ноль» и все они разом стали бы просроченными
навсегда — та же ложная тревога, только оптом.
"""
result = _kit_runs._run_result_count(counters)
if result is None:
return status == "done"
return result > 0
@dataclass(frozen=True)
class FreshnessRow:
"""Свёртка прогонов одного расписания: когда источник ПОСЛЕДНИЙ РАЗ дал данные."""
source: str
interval_days: Any
since: datetime
never_ok: bool
def freshness_rows(rows: list[Any]) -> list[FreshnessRow]:
"""Строки `_STALE_SOURCES_SQL` (прогон на строку) → свежесть на источник.
since = последний прогон, ПРИНЁСШИЙ ДАННЫЕ (любой статус); нет такого — created_at
расписания, как и раньше. never_ok считается ТОЙ ЖЕ мерой: иначе «никогда не
собирал» соврал бы в другую сторону — источник с одними лишь 'banned'-прогонами с
данными числился бы ни разу не собиравшим.
"""
meta: dict[str, tuple[Any, datetime]] = {}
last_data: dict[str, datetime] = {}
for row in rows:
meta.setdefault(row.source, (row.interval_days, row.created_at))
if row.finished_at is None or not run_brought_data(row.status, row.counters):
continue
prev = last_data.get(row.source)
if prev is None or row.finished_at > prev:
last_data[row.source] = row.finished_at
return [
FreshnessRow(
source=source,
interval_days=interval_days,
since=last_data.get(source) or created_at,
never_ok=source not in last_data,
)
for source, (interval_days, created_at) in meta.items()
]
def stale_sources(rows: list[Any], now: datetime) -> list[StaleSource]:
"""Чистая часть сводки: какие расписания просрочены и на сколько (свежие — внизу).
Просрочка меряется в ТАКТАХ, а не в сутках: у rosreestr_quarter_poll такт 28 суток,
и 24 суток без сбора для него норма, а для суточного domclick_city_sweep — авария.
Сортировка по числу пропущенных тактов, а не по календарю, по той же причине.
"""
stale: list[StaleSource] = []
for row in rows:
interval = _schedule_interval_days(row.interval_days)
age_days = (now - row.since).total_seconds() / 86400.0
if age_days > STALE_DIGEST_INTERVAL_FACTOR * interval:
stale.append(
StaleSource(
source=row.source,
interval_days=interval,
age_days=age_days,
never_ok=bool(row.never_ok),
)
)
stale.sort(key=lambda s: s.age_days / s.interval_days, reverse=True)
return stale
def emit_stale_digest(db: Session, *, now: datetime | None = None) -> list[StaleSource]:
"""Раз в STALE_DIGEST_PERIOD_H часов — одно событие «что сейчас не собирает».
Возвращает список просроченных источников (пустой — либо всё свежо, либо выпуск ещё
не подошёл по времени). Best-effort, как и оба сторожа в runs.py: сводка не имеет
права уронить тик планировщика.
"""
global _last_stale_digest_at
now = now or datetime.now(UTC)
if _last_stale_digest_at is not None and now - _last_stale_digest_at < timedelta(
hours=STALE_DIGEST_PERIOD_H
):
return []
try:
run_rows = list(db.execute(_STALE_SOURCES_SQL).fetchall())
stale = stale_sources(freshness_rows(run_rows), now)
_last_stale_digest_at = now
if not stale:
logger.info("scheduler: stale-digest — просроченных источников нет (#2670)")
return []
details = ", ".join(
f"{s.source} {s.age_days:.1f}d/{s.interval_days}d"
+ (" (успеха не было ни разу)" if s.never_ok else "")
for s in stale
)
logger.error(
"scheduler: %d источников не собирают дольше %d× своего такта — %s (#2670)",
len(stale),
STALE_DIGEST_INTERVAL_FACTOR,
details,
)
if sentry_sdk is not None:
sentry_sdk.capture_message(
f"{len(stale)} scraper sources are stale (no successful run for more than "
f"{STALE_DIGEST_INTERVAL_FACTOR}× their schedule interval): {details}",
level="error",
)
return stale
except Exception:
logger.exception("scheduler: stale-digest failed")
return []
# ── типы job/handler ─────────────────────────────────────────────────────────
# Job получает свежую сессию (открыта `_dispatch`), run_id, params и весь контекст
# (config/matcher/enrichment/runs) — чтобы иметь доступ к инжектированным зависимостям.
Job = Callable[["Session", int, dict[str, Any], "SchedulerContext"], Coroutine[Any, Any, None]]
# Pre-claim gate: вернуть False → пропустить claim (напр. cian: нет/протухли cookies).
PreClaim = Callable[["Session", dict[str, Any], "SchedulerContext"], Coroutine[Any, Any, bool]]
# Post-claim hook: сразу после успешного claim (напр. proxy_healthcheck sub-hourly reschedule).
PostClaim = Callable[["Session", int, dict[str, Any], "SchedulerContext"], None]
@dataclass(frozen=True)
class Handler:
"""Одна запись реестра source→обработчик.
Инкапсулирует всё, что раньше размазывалось по индивидуальному `trigger_*_run`:
- `job` — асинхронное тело (SERP/backfill/enrichment). Общий boilerplate
(claim, свежая сессия, create_task, close, log) делает `_dispatch`.
- `pre_claim` — опциональный гейт ДО claim (cookie-проверка cian → defer).
- `post_claim` — опциональный хук сразу ПОСЛЕ claim (sub-hourly reschedule).
- `log_name` — имя источника в логах ("scheduler: triggered <log_name> run_id=…").
"""
job: Job
log_name: str
pre_claim: PreClaim | None = None
post_claim: PostClaim | None = None
@dataclass
class SchedulerContext:
"""Инжектируемые зависимости + рантайм-состояние планировщика.
Заменяет прямые импорты `app.core.config.settings` / `app.services.matching` /
`app.services.house_imv_backfill` / `app.core.db.SessionLocal` /
`app.services.scrape_runs` / `app.core.shutdown.shutdown_requested`.
"""
config: ScraperConfig
matcher: HouseMatcher
enrichment: EnrichmentJobs
session_factory: SessionFactory
shutdown_requested: Callable[[], bool] = lambda: False
# Proxy-пул (#2163 curl / #2164 P4 browser). None → пул выключен, curl/browser
# берут прокси из env (ship-dark). RealProxyProvider инжектируется scheduler_main,
# когда включён любой из флагов use_proxy_pool_curl / use_proxy_pool_browser.
# Прокидывается в sweep-пайплайны, откуда попадает в провайдеры (curl_proxy_url) и
# BrowserFetcher (POST /fetch{proxy}).
proxy_provider: ProxyProvider | None = None
# runs-модуль (create_run/update_heartbeat/mark_done/… ). По умолчанию — kit-копия
# scraper_kit.orchestration.runs; в тестах подменяется recorder'ом.
runs: Any = _kit_runs
# #1182 P2: strong-ref'ы detached run-задач для graceful SIGTERM-drain'а.
_inflight_tasks: set[asyncio.Task[None]] = field(default_factory=set)
# #3391: task → run_id заклеймленного прогона. Без него дрейн знает, что задача
# не дожила до финализации, но не знает, КАКУЮ строку scrape_runs снимать с
# 'running' (claim-логика run_id никуда не публиковала).
_inflight_run_ids: dict[asyncio.Task[None], int] = field(default_factory=dict)
def spawn_tracked(
self, coro: Coroutine[Any, Any, None], *, run_id: int | None = None
) -> asyncio.Task[None]:
"""create_task + регистрация в _inflight_tasks для graceful-drain'а (#1182 P2).
strong-ref в set'е держит задачу до завершения (RUF006) и даёт scheduler_loop'у
дождаться её на SIGTERM-drain'е; done-callback ретривит exception и убирает задачу.
`run_id` (#3391) — строка scrape_runs этой задачи; None у задач без прогона.
"""
task = asyncio.create_task(coro)
self._inflight_tasks.add(task)
if run_id is not None:
self._inflight_run_ids[task] = run_id
def _on_done(t: asyncio.Task[None]) -> None:
self._inflight_tasks.discard(t)
self._inflight_run_ids.pop(t, None)
if not t.cancelled():
t.exception()
task.add_done_callback(_on_done)
return task
def mark_inflight_interrupted(self, tasks: Iterable[asyncio.Task[None]]) -> int:
"""Снять с 'running' прогоны задач, не доживших до собственной финализации (#3391).
Прод 07.09 (первый настоящий SIGTERM-drain после #3363): два app-task бэкфилла
(cian_detail_backfill 6167, cian_history_backfill 6173) не уложились в grace,
scheduler_main их хард-кансельнул — и строки остались 'running' до boot-reap'а
следующего контейнера, где стали 'zombie' с `boot_reaped=true`. `interrupted=1`
при дрейне пишут только kit-пайплайны (#3363) и DKP-импорт: у задач, чьё тело
живёт в app и о дрейне не знает, метку ставить некому. Ставим её ЗДЕСЬ — в
единственной точке, через которую проходит любая detached run-задача.
Синхронно (sync-сессия, ни одного await): вызывается в том числе из except
CancelledError, где лишний await мог бы не дожить до конца.
Статус строки перечитывается перед записью: задача, успевшая финализироваться
сама (done/failed/banned), не перезаписывается. Гонку добивает сам `mark_done`
(`WHERE status = 'running'`) — но тогда счётчики уже прочитаны, и их merge был бы
холостым. counters берём из строки и дописываем `interrupted`, а не отдаём
`{"interrupted": 1}` голым: app-копия `mark_done` counters ЗАМЕНЯЕТ, а не мержит
(#3390) — голый словарь стёр бы всю бухгалтерию прогона, включая чекпоинт.
Потолок (осознанный, #3391): SELECT синхронный, а у движка нет ни connect-, ни
statement-таймаута (`app/core/db.py:8-19`). Недоступная БД блокирует луп до
SIGKILL'а — 20 с docker-grace (120 s stop_grace_period 100 s _DRAIN_TIMEOUT_S).
Данные при этом НЕ хуже прежних: строки просто остаются 'running' и их снимет
boot-reap следующего контейнера — ровно то, что было до этой пометки. Поднимать
до отдельного таймаута есть смысл только вместе с таймаутами на самом движке.
"""
run_ids = [rid for t in tasks if (rid := self._inflight_run_ids.get(t)) is not None]
if not run_ids:
return 0
marked_ids: list[int] = []
db = None
# Открытие сессии и её закрытие — ВНУТРИ try: вызов приходит из-под
# `except CancelledError`, и исключение отсюда ЗАМЕНИЛО бы отмену собой.
# `suppress(CancelledError)` в scheduler_main такое не глушит → процесс уходит
# с трейсбеком вместо чистого drain-выхода. CancelledError (BaseException)
# через `except Exception` проходит насквозь и пробрасывается как есть.
try:
db = self.session_factory()
for run_id in run_ids:
try:
row = db.execute(
text("SELECT status, counters FROM scrape_runs WHERE id = :id"),
{"id": run_id},
).fetchone()
if row is None or row.status != "running":
continue
counters = dict(row.counters or {})
counters["interrupted"] = 1
self.runs.mark_done(db, run_id, counters)
marked_ids.append(run_id)
except Exception:
# Один непроходимый прогон не должен утащить остальные: дрейну
# осталось секунды до SIGKILL, а каждая незакрытая строка — это
# ещё один 'zombie' у следующего контейнера.
logger.exception("scheduler: drain — не удалось пометить run_id=%d", run_id)
# Отказавший statement оставляет сессию в aborted-tx, и КАЖДЫЙ
# следующий run_id падал бы на ровном месте (образец — defensive
# rollback в mark_failed/mark_banned, app/services/scrape_runs.py).
with suppress(Exception):
db.rollback()
except Exception:
logger.exception("scheduler: drain — пометка interrupted не выполнена")
finally:
if db is not None:
with suppress(Exception):
db.close()
if marked_ids:
# Именно marked_ids, а не весь run_ids: в списке не должно быть прогонов,
# пропущенных по статусу, — иначе лог приписывает пометку тем, кого не трогал.
logger.warning(
"scheduler: drain — %d прогон(ов) сняты с 'running' как interrupted: %s",
len(marked_ids),
marked_ids,
)
return len(marked_ids)
async def drain_inflight(self) -> None:
"""Дождаться завершения detached run-задач перед teardown'ом (#1182 P2).
Кооперативные дети дочекивают текущий unit + mark_done и резолвятся;
некооперативные упираются в _CHILD_DRAIN_TIMEOUT_S либо во внешний hard-cancel
из scheduler_main — и в обоих случаях их строки помечаются `interrupted`
(#3391), иначе они доживают в 'running' до boot-reap'а. Не busy-spin — один
asyncio.wait.
"""
pending = [t for t in self._inflight_tasks if not t.done()]
if not pending:
return
logger.info("scheduler: draining %d in-flight run task(s) on shutdown", len(pending))
try:
_done, still = await asyncio.wait(pending, timeout=_CHILD_DRAIN_TIMEOUT_S)
except asyncio.CancelledError:
# Hard-cancel из scheduler_main (grace 100s истёк раньше нашего дрейна —
# такт tick-сна + _CHILD_DRAIN_TIMEOUT_S могут превысить его). Пишем метку
# синхронно и пробрасываем отмену дальше.
self.mark_inflight_interrupted([t for t in pending if not t.done()])
raise
if still:
logger.warning(
"scheduler: %d task(s) did not drain in %.0fs — leaving for hard-cancel",
len(still),
_CHILD_DRAIN_TIMEOUT_S,
)
self.mark_inflight_interrupted(still)
else:
logger.info("scheduler: all in-flight run task(s) drained cleanly")
def compute_next_run_at(
window_start_hour: int,
window_end_hour: int,
*,
now: datetime | None = None,
interval_days: int = 1,
) -> datetime:
"""Pick random datetime в window [start, end) UTC, через interval_days суток после now.
interval_days задаёт каденс источника: 1 (default) = daily (back-compat), 7 = weekly.
Берётся из schedule.default_params["interval_days"] вызывающим кодом (_claim_run /
_defer_next_run_at); отсутствие ключа → 1 → прежнее ежедневное поведение.
Если window_end_hour <= window_start_hour → cross-midnight window
(например 22→3 → окно 22:00-23:59 ИЛИ 00:00-02:59).
"""
now = now or datetime.now(tz=UTC)
interval_days = max(1, int(interval_days))
# Целевая дата = now + interval_days суток (interval_days=1 → завтра, как раньше).
target = (now + timedelta(days=interval_days)).date()
if window_end_hour > window_start_hour:
# Обычное окно (например 2..5 → 02:00-04:59)
start_seconds = window_start_hour * 3600
end_seconds = window_end_hour * 3600
rand_seconds = random.randint(start_seconds, end_seconds - 1)
return datetime.combine(target, time(0, 0), tzinfo=UTC) + timedelta(seconds=rand_seconds)
else:
# Cross-midnight (22..3 → 22:00-23:59 + 00:00-02:59)
# Длина окна = (24-start) + end часов
total_seconds = ((24 - window_start_hour) + window_end_hour) * 3600
rand_seconds = random.randint(0, total_seconds - 1)
# Если rand попадает в первую часть (start..24)
first_half = (24 - window_start_hour) * 3600
if rand_seconds < first_half:
# interval_days=1: текущая дата (если окно ещё не наступило сегодня) или next day.
# interval_days>1: всегда целевая дата (стаггер на N суток вперёд).
today_ok = interval_days == 1 and now.hour < window_start_hour
base_date = now.date() if today_ok else target
return datetime.combine(base_date, time(0, 0), tzinfo=UTC) + timedelta(
seconds=window_start_hour * 3600 + rand_seconds
)
else:
# Во второй части (0..end), целевого дня
offset = rand_seconds - first_half
return datetime.combine(target, time(0, 0), tzinfo=UTC) + timedelta(seconds=offset)
def has_running_run(db: Session, source: str) -> bool:
"""Есть ли активный run для source (status='running')."""
row = db.execute(
text(
"""
SELECT 1 FROM scrape_runs
WHERE source = :source AND status = 'running'
LIMIT 1
"""
),
{"source": source},
).fetchone()
return row is not None
def reap_zombies(db: Session) -> int:
"""Mark scrape_runs as 'zombie' если heartbeat не обновлялся > ZOMBIE_THRESHOLD_HOURS hours.
clock_timestamp(), а не now() (#2702): критерий сравнивает ЗАПИСАННЫЙ heartbeat со
временем «сейчас», и обе стороны сравнения должны быть настоящим временем. `now()`
замерзает на старте транзакции — у писавшей heartbeat стороны это давало отставание
на весь возраст открытой рабочей транзакции (см. docstring runs.py), у читающей
стороны — на возраст тика. На проде это уже стоило ложных срабатываний: у всех 6
прогонов cian_history_backfill, помеченных 'zombie', записанный heartbeat так и
остался на отметке старта (max advance 0.0 с) — при том что нормальный прогон этого
источника длится до 5.06 ч (прогон 346) и обязан был двигать heartbeat.
"""
zombie_interval = f"{ZOMBIE_THRESHOLD_HOURS} hours"
result = db.execute(
text(
"""
UPDATE scrape_runs
SET status = 'zombie', finished_at = clock_timestamp()
WHERE status = 'running'
AND (heartbeat_at IS NULL
OR heartbeat_at < clock_timestamp() - CAST(:interval AS interval))
RETURNING id
"""
),
{"interval": zombie_interval},
)
rows = result.fetchall()
db.commit()
if rows:
logger.warning("scheduler: reaped %d zombie runs: %s", len(rows), [r.id for r in rows])
return len(rows)
def reap_boot_zombies(db: Session, boot_time: datetime) -> int:
"""Однократный boot-reap (#3122): прогон, стартовавший раньше старта СОБСТВЕННОГО
процесса планировщика, мёртв по построению — его сборщик жил в предыдущем
контейнере и умер вместе с ним. Пульс в критерии не участвует вовсе, поэтому
ложные срабатывания класса #2702 (редкий пульс у 5-часовых прогонов)
невозможны: живой прогон этого процесса не может быть старше самого процесса.
Прод-цена без этого (27.08, после ночного офлайна #3119): 4 зомби блокировали
свои источники через has_running_run до 6-часового порогового reap'а — до пяти
часов слепоты на источник ровно после простоя, когда догон нужнее всего.
Минута запаса на часовой дрейф между хостом БД и приложением. Маркер
`boot_reaped: true` в counters отличает этот исход от порогового 'zombie':
у порогового процесс МОЖЕТ быть жив (reap не убивает его — см. комментарий
у _RESUME_STATUSES), у boot-зомби — гарантированно мёртв, поэтому его
чекпоинт безопасен для подхвата (_resume_decision пускает 'zombie' только
с этим маркером).
"""
result = db.execute(
text(
"""
UPDATE scrape_runs
SET status = 'zombie',
finished_at = clock_timestamp(),
counters = COALESCE(counters, CAST('{}' AS jsonb))
|| CAST('{"boot_reaped": true}' AS jsonb)
WHERE status = 'running'
AND started_at < CAST(:boot AS timestamptz) - CAST('60 seconds' AS interval)
RETURNING id, source
"""
),
{"boot": boot_time},
)
rows = result.fetchall()
db.commit()
if rows:
logger.warning(
"scheduler: boot-reap — %d прогонов предыдущего контейнера сняты с 'running': %s",
len(rows),
[(r.id, r.source) for r in rows],
)
else:
# Ноль — штатный исход (деплоевский startup-reap обычно успевает раньше),
# но исполнение обязано быть видимым: молчащий при нуле механизм
# неотличим от неподключённого (первая же живая приёмка #3122 это
# показала — оборванную проводку было бы не отличить от чистого прода).
logger.info("scheduler: boot-reap — прогонов предыдущего контейнера нет (0)")
return len(rows)
def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext) -> int | None:
"""Claim run: INSERT scrape_runs + UPDATE scrape_schedules next_run_at.
Returns run_id или None если уже есть running run для этого source.
Общий helper для всех source-handler'ов. Advisory-lock concurrency-логика перенесена
ДОСЛОВНО из боевого scheduler'а (создание run — через инжектированный ctx.runs).
#2658: каждая ветка «пропустили окно» пишет строку scrape_runs(status='skipped')
с машиночитаемой причиной — раньше пропуск был виден только в docker-логах, которые
теряются при редеплое. Пишем через kit-копию runs (у инжектированной app-копии
mark_skipped нет — строки-пропуски создаёт только планировщик).
"""
source = schedule_row["source"]
if has_running_run(db, source):
logger.info("scheduler: skip — already running for source=%s", source)
_kit_runs.mark_skipped(
db,
source=source,
reason=SKIP_ALREADY_RUNNING,
details="предыдущий прогон ещё идёт",
)
return None
# #750: атомарный claim через transaction-scoped advisory lock. Сериализует
# конкурентные _claim_run одного source (redeploy-overlap старого+нового контейнера
# / 2 реплики), иначе оба пройдут has_running_run-pre-check и запустят ДВОЙНОЙ sweep
# → 2× rate к Avito/Cian → бан. pg_try_advisory_xact_lock НЕ блокирует: lock занят
# другим тиком → false → return None. Лок авто-снимается на commit/rollback этой
# транзакции (has_running_run выше — быстрый pre-check, обычно без взятия лока).
got_lock = db.execute(
text("SELECT pg_try_advisory_xact_lock(hashtext(:source))"),
{"source": source},
).scalar()
if not got_lock:
logger.info("scheduler: skip — concurrent claim in progress for source=%s", source)
_kit_runs.mark_skipped(
db,
source=source,
reason=SKIP_CONCURRENT_CLAIM,
details="конкурентный тик уже клеймит этот source",
)
return None
# Double-checked под локом: конкурентный тик мог закоммитить running-run МЕЖДУ
# pre-check и взятием лока (затем освободить lock на commit). Перепроверяем под
# локом — иначе словим двойной run на узком окне.
if has_running_run(db, source):
logger.info("scheduler: skip — running appeared under lock for source=%s", source)
db.rollback() # освобождаем advisory lock (claim не состоялся)
# mark_skipped ПОСЛЕ rollback'а: его INSERT + commit иначе улетели бы в откат.
_kit_runs.mark_skipped(
db,
source=source,
reason=SKIP_RUNNING_UNDER_LOCK,
details="running-прогон появился под локом",
)
return None
params = schedule_row.get("default_params") or {}
run_id = ctx.runs.create_run(db, source=source, params=params)
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(
"""
UPDATE scrape_schedules
SET last_run_id = :run_id, last_run_at = NOW(),
next_run_at = :next_at, updated_at = NOW()
WHERE source = :source
"""
),
{"run_id": run_id, "next_at": next_at, "source": source},
)
db.commit()
logger.info(
"scheduler: claimed run_id=%d source=%s next_run_at=%s",
run_id,
source,
next_at.isoformat(),
)
return run_id
# ── подхват чекпоинта оборванного прогона (#930, вторая половина) ────────────
# #930 сделал чекпоинт (`counters.done_buckets`) и приёмную сторону
# (`run_*_full_load(resume_run_id=...)`), но единственным входом оставил админку. У
# avito full-load её нет вовсе, поэтому 299 из 433 «впустую перебранных» корзин за 90
# суток не имели НИКАКОГО пути возобновления, даже ручного. Планировщик передавал
# resume_run_id=None литералом.
#
# Точка берётся только когда выполнены ВСЕ условия ниже; иначе прогон честно начинает с
# нуля, а ПРИЧИНА пишется в его counters (молчаливый отказ неотличим от отсутствия
# правки — см. _resume_decision).
# 'zombie' НЕ в списке НАМЕРЕННО. reap_zombies снимает пометку 'running', но НЕ убивает
# процесс (прямо задокументировано в app/tasks/listing_source_snapshot.py) — а
# has_running_run гейтит claim именно по статусу. То есть после reap'а старый сборщик
# может продолжать писать в ту же строку: подхват читал бы ДВИЖУЩУЮСЯ точку и запускал
# второй сборщик на ту же площадку. 147 корзин в 10 zombie-прогонах за 90 суток
# остаются несобранными сознательно — это цена, а не недосмотр.
# 'failed' в списке: его чекпоинт больше не стирается финализатором (runs.py, мерж
# counters), а причина отказа («наш баг») ничего не говорит о полноте УЖЕ записанных
# корзин — они записаны тем же heartbeat'ом, что и у banned.
_RESUME_STATUSES = frozenset({"banned", "cancelled", "failed"})
# 'zombie' допускается к подхвату ТОЛЬКО с маркером counters.boot_reaped (#3122):
# у boot-зомби процесс гарантированно мёртв (жил в предыдущем контейнере), значит
# «движущейся точки» — причины, по которой 'zombie' исключён выше, — быть не может.
# Пороговый zombie без маркера по-прежнему отвергается.
# Длина цепочки возобновлений. Не круглое число: STALE_DIGEST_INTERVAL_FACTOR (=3) —
# уже существующий в этом файле порог «источник не собирал дольше 3× своего такта =
# сломан». Цепочка не имеет права отодвинуть полный обход дальше этой же черты,
# поэтому подряд идущих подхватов допускается на один меньше: прогоны 1 и 2 могут
# продолжать предшественника, третий обязан пойти с нуля. При такте avito 7 суток это
# гарантирует попытку полного обхода не реже, чем раз в 21 сутки — ровно в тот момент,
# когда сводка объявляет источник просроченным.
_MAX_RESUME_CHAIN = STALE_DIGEST_INTERVAL_FACTOR - 1
# Кандидат — ПОСЛЕДНИЙ прогон источника, а не последний подходящий: если после обрыва
# уже прошёл полный ('done') прогон, дерево обойдено и возобновлять нечего. Строки
# 'skipped' — бухгалтерия планировщика, а не прогоны, поэтому не в счёт.
_RESUME_CANDIDATE_SQL = text("""
WITH cur AS (
SELECT source, params FROM scrape_runs WHERE id = CAST(:rid AS bigint)
)
SELECT r.id AS prev_id,
r.status AS prev_status,
r.counters AS prev_counters,
(r.params IS NOT DISTINCT FROM cur.params) AS same_params,
EXTRACT(EPOCH FROM (clock_timestamp() - r.heartbeat_at)) / 3600.0 AS age_h,
cur.params ->> 'interval_days' AS interval_days
FROM scrape_runs r, cur
WHERE r.source = cur.source
AND r.id <> CAST(:rid AS bigint)
AND r.status <> 'skipped'
ORDER BY r.started_at DESC
LIMIT 1
""")
def _resume_decision(row: Any) -> tuple[int | None, dict[str, Any]]:
"""Решение «подхватывать ли точку» + счётчики-объяснение. Чистая функция.
Возвращает (resume_run_id | None, counters-заготовка нового прогона). Причина
отказа — машиночитаемый слаг в `resume_reason`, по нему «предыдущего прогона не
было» отличается от «параметры разъехались» ЗАПРОСОМ, а не грепом логов.
Срок годности точки — такт источника ПЛЮС сутки (`interval_days` из его же params).
Ни одно из слагаемых не выбрано произвольно. Такт — объявленный самим расписанием
срок, в течение которого собранное считается свежим; точка старше него пережила
цикл, в котором источник обязан был обойти дерево целиком. Сутки — сетка запуска:
`compute_next_run_at` выбирает ДЕНЬ (сегодня+interval_days) и случайное время внутри
окна, поэтому два соседних запуска отстоят друг от друга на interval_days ± меньше
суток. Без этого слагаемого точку отвергал бы jitter расписания, а не устаревание:
прогон 3547 убит деплоем 09.08 16:53, расписание 139 подхватит его 16.08 13:37 —
164.6 ч при такте 168 ч, запас 3.4 ч при ширине окна 2 ч. Пропущенный цикл в окно
всё равно не влезает: для avito это 13 суток против порога 8.
"""
if row is None:
return None, {"resume_from": None, "resume_reason": "no_prev_run", "resume_chain": 0}
prev_counters = row.prev_counters if isinstance(row.prev_counters, dict) else {}
done_buckets = prev_counters.get("done_buckets")
done_n = len(done_buckets) if isinstance(done_buckets, list) else 0
chain_raw = prev_counters.get("resume_chain")
prev_chain = chain_raw if isinstance(chain_raw, int) else 0
verdict: dict[str, Any] = {
"resume_from": None,
"resume_candidate": int(row.prev_id),
"resume_buckets": done_n,
"resume_chain": 0,
}
_boot_reaped_zombie = row.prev_status == "zombie" and prev_counters.get("boot_reaped") is True
# 'done' с counters.interrupted=1 — SIGTERM-drain (#3319): статус штатный, но обход
# оборван на границе корзины, часть дерева не собрана. Сам статус в _RESUME_STATUSES
# не добавлен НАМЕРЕННО: чистое 'done' — полный проход, резюмить у него нечего, а
# подхват такой точки означал бы, что источник больше никогда не обходится целиком.
# Метка — та же, что у rosreestr_dkp-дрейна (app/services/scheduler.py), поэтому ни
# один читатель статуса не меняется.
_drained_done = row.prev_status == "done" and bool(prev_counters.get("interrupted"))
if row.prev_status not in _RESUME_STATUSES and not _boot_reaped_zombie and not _drained_done:
verdict["resume_reason"] = f"status_{row.prev_status}"
elif not row.same_params:
verdict["resume_reason"] = "params_changed"
elif done_n == 0:
verdict["resume_reason"] = "no_checkpoint"
elif row.age_h is None or float(row.age_h) > 24.0 * (
_schedule_interval_days(row.interval_days) + 1
):
verdict["resume_reason"] = "checkpoint_stale"
elif prev_chain >= _MAX_RESUME_CHAIN:
verdict["resume_reason"] = "chain_limit"
else:
verdict["resume_from"] = int(row.prev_id)
verdict["resume_reason"] = "ok"
verdict["resume_chain"] = prev_chain + 1
# Наследуем чекпоинт в counters нового прогона ПРЯМО ПРИ CLAIM (#3074).
# Прод-факт (run 4707, 23.08): возобновлённый прогон, убитый деплоем на
# 26-й минуте — ДО завершения первой новой корзины, — не успел ни разу
# написать heartbeat с done_buckets. Его собственный чекпоинт остался
# пуст, кандидат следующей субботы увидел no_checkpoint, и 42 корзины
# предшественника (4117) пропали. Пайплайн сеет `done = set(skip_set)`
# у себя в памяти, но до первого _on_bucket это знание нигде не
# персистится; heartbeat мержит jsonb (`counters || :counters`), так что
# первый же настоящий bucket-heartbeat перезапишет ключ тем же множеством
# плюс новое — двойной записи не возникает.
verdict["done_buckets"] = sorted(done_buckets) if done_buckets else []
return int(row.prev_id), verdict
return None, verdict
def _pick_resume(db: Session, run_id: int) -> int | None:
"""Чекпоинт какого прогона наследует `run_id` (или None) + запись вердикта.
Тождество задания сверяется РОВНО по тем полям, которыми задание задаётся: source
(кандидат ищется в пределах одного source) и params целиком, побайтово. Кандидат
сравнивается с ТЕКУЩИМ прогоном, а не с расписанием, потому что именно params
прогона поехали в pipeline. На проде за 90 суток 89 корзин из 433 (21%) лежат в
прогонах, чьи params отличаются от следующего — по ним пропуск был бы неверным:
`incremental_days` меняет СМЫСЛ ключа (дочитано до watermark ≠ бакет перебран), а
`price_cap_per_bucket` меняет само дерево бисекции, то есть какие ключи существуют.
"""
row = db.execute(_RESUME_CANDIDATE_SQL, {"rid": run_id}).fetchone()
resume_run_id, verdict = _resume_decision(row)
# Вердикт кладём в counters НОВОГО прогона: `update_heartbeat` мержит jsonb, поэтому
# последующие heartbeat'ы пайплайна его не затрут и он доживёт до финализатора.
_kit_runs.update_heartbeat(db, run_id, verdict)
logger.info("scheduler: resume run_id=%d%s", run_id, verdict)
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 {}
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(
"""
UPDATE scrape_schedules
SET next_run_at = :next_at, updated_at = NOW()
WHERE source = :source
"""
),
{"next_at": next_at, "source": source},
)
db.commit()
logger.info(
"scheduler: deferred next_run_at source=%s next_run_at=%s (no run claimed)",
source,
next_at.isoformat(),
)
def reschedule_after_minutes(*, param: str = "interval_minutes", default: int = 30) -> PostClaim:
"""Фабрика post_claim-хука: sub-hourly re-schedule сразу после claim (#2162).
Боевой scheduler держит суточную гранулярность в compute_next_run_at, а
proxy_healthcheck хочет каденс раз в N минут. Хук СРАЗУ после claim переопределяет
next_run_at на now() + interval_minutes (source-специфично, не трогая shared
compute_next_run_at). Если run ещё идёт на следующем тике — has_running_run в
_claim_run вернёт None (skip) и next_run_at не сбросится → авто-throttle по факту.
"""
def _hook(db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext) -> None:
# source читаем из ЗАКЛЕЙМЛЕННОГО schedule — берём из params-owner через closure нельзя,
# поэтому апдейтим по run_id → source (schedule owner). Проще: UPDATE по source из
# scrape_runs текущего run_id, что эквивалентно WHERE source = <this source>.
minutes = _interval_minutes(params, param=param)
if minutes is None:
minutes = default
db.execute(
text(
"""
UPDATE scrape_schedules
SET next_run_at = now() + make_interval(mins => CAST(:mins AS integer)),
updated_at = now()
WHERE source = (SELECT source FROM scrape_runs WHERE id = CAST(:rid AS bigint))
"""
),
{"mins": minutes, "rid": run_id},
)
db.commit()
return _hook
# ── единственный параметризованный dispatch (заменяет 25 trigger_* + 27 if/elif) ──
async def _dispatch(
handler: Handler,
db: Session,
schedule_row: dict[str, Any],
ctx: SchedulerContext,
) -> int | None:
"""Общий путь claim→spawn для ЛЮБОГО source (весь boilerplate прежних trigger_*_run).
1. Опциональный pre_claim-гейт (False → skip без claim).
2. _claim_run (advisory-lock; None → уже running / concurrent claim).
3. Опциональный post_claim-хук (sub-hourly reschedule).
4. spawn detached run-задачи: свежая сессия → handler.job → close (+ exception-log).
"""
if handler.pre_claim is not None:
proceed = await handler.pre_claim(db, schedule_row, ctx)
if not proceed:
return None
run_id = _claim_run(db, schedule_row, ctx)
if run_id is None:
return None
params = schedule_row.get("default_params") or {}
if handler.post_claim is not None:
handler.post_claim(db, run_id, params, ctx)
async def _run() -> None:
run_db = ctx.session_factory()
try:
await handler.job(run_db, run_id, params, ctx)
except Exception:
logger.exception("scheduler: %s crashed run_id=%d", handler.log_name, run_id)
finally:
run_db.close()
ctx.spawn_tracked(_run(), run_id=run_id)
logger.info("scheduler: triggered %s run_id=%d", handler.log_name, run_id)
return run_id
def resolve_handler(source: str, registry: Mapping[str, Handler]) -> Handler | None:
"""Найти handler по source: exact-match, затем wildcard-префикс ("deactivate_stale_*").
Единственная точка ветвления вместо 27 if/elif. Wildcard-ключ (кончается на "*")
матчит по префиксу — покрывает семейство deactivate_stale_avito/yandex/cian одной
записью реестра (как боевой `source.startswith("deactivate_stale_")`).
"""
handler = registry.get(source)
if handler is not None:
return handler
for key, h in registry.items():
if key.endswith("*") and source.startswith(key[:-1]):
return h
return None
# ── kit-native sweep-обработчики (тело в scraper_kit.orchestration.pipeline) ──────
# Param-чтение идентично боевым trigger_*_run — та же семантика окон/дефолтов.
async def _job_avito_city_sweep(
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
) -> None:
# #B1 oblast rollout: default_params["city"] (slug, e.g. "nizhniy_tagil") → anchors
# города вместо EKB_ANCHORS. Отсутствует/неизвестен → get_city_anchors вернёт None →
# run_avito_city_sweep сам падает на EKB_ANCHORS (прежнее поведение без city).
city = params.get("city")
anchors = get_city_anchors(city) if city else None
await run_avito_city_sweep(
db,
run_id=run_id,
config=ctx.config,
matcher=ctx.matcher,
enrichment=ctx.enrichment,
shutdown_requested=ctx.shutdown_requested,
proxy_provider=ctx.proxy_provider,
anchors=anchors,
# #2487: слаг города-цели → SERP-фильтр _parse_html оставляет карточки
# этого города (иначе oblast-sweep дропал 100% из-за хардкода /ekaterinburg/).
city_slug=city,
pages_per_anchor=int(params.get("pages_per_anchor", 3)),
detail_top_n=int(params.get("detail_top_n", 20)),
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
enrich_houses=bool(params.get("enrich_houses", True)),
radius_m=int(params.get("radius_m", 1500)),
# #3074: подхват якорей у оборванного предшественника. Ключ чекпоинта —
# имя якоря, оно не зависит от количества якорей, поэтому в отличие от
# combo-чекпоинта yandex-свипа гарда по числу якорей здесь не требуется.
resume_run_id=_pick_resume(db, run_id),
)
async def _job_avito_newbuilding_sweep(
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
) -> None:
await run_avito_newbuilding_sweep(
db,
run_id=run_id,
config=ctx.config,
matcher=ctx.matcher,
shutdown_requested=ctx.shutdown_requested,
proxy_provider=ctx.proxy_provider,
pages=int(params.get("pages", 20)),
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
# #3074: подхват страниц у оборванного предшественника — см.
# _job_avito_city_sweep выше.
resume_run_id=_pick_resume(db, run_id),
)
async def _job_avito_full_load(
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
) -> None:
_incremental_days = params.get("incremental_days")
incremental_days = int(_incremental_days) if _incremental_days is not None else None
# Окно ретроспективы не может быть уже такта (#2674). incremental_days задаёт
# since=today-N, а прогон повторяется раз в interval_days суток → при N < interval_days
# объявления, поднятые в дни (D, D+interval_days-N), не попадают НИ в один прогон.
# Ровно это и случилось: миграция 206 перевела источник на interval_days=7, оставив
# incremental_days=2 из миграции 129 → 3 календарных дня из 7 в поле зрения.
# Два независимых литерала, которые обязаны совпадать, однажды уже разъехались —
# поэтому расхождение чиним здесь, а не только данными.
# None-safe так же, как incremental_days выше: `"interval_days": null` в jsonb
# приезжает сюда как None, а int(None) — TypeError.
_interval_days = params.get("interval_days")
interval_days = max(1, int(_interval_days)) if _interval_days is not None else 1
if incremental_days is not None and incremental_days < interval_days:
logger.warning(
"avito_full_load: окно ретроспективы incremental_days=%d уже такта "
"interval_days=%d — расширяю до %d, иначе %d суток объявлений не видит "
"ни один прогон",
incremental_days,
interval_days,
interval_days,
interval_days - incremental_days,
)
incremental_days = interval_days
await run_avito_full_load(
db,
run_id=run_id,
config=ctx.config,
matcher=ctx.matcher,
shutdown_requested=ctx.shutdown_requested,
proxy_provider=ctx.proxy_provider,
price_cap_per_bucket=int(params.get("price_cap_per_bucket", 1400)),
concurrency=int(params.get("concurrency", 5)),
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
secondary_only=bool(params.get("secondary_only", True)),
resume_run_id=_pick_resume(db, run_id),
incremental_days=incremental_days,
)
async def _job_avito_full_load_exhaustive(
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
) -> None:
# incremental_days форсится в None → полный exhaustive обход (weekly cadence).
await run_avito_full_load(
db,
run_id=run_id,
config=ctx.config,
matcher=ctx.matcher,
shutdown_requested=ctx.shutdown_requested,
proxy_provider=ctx.proxy_provider,
price_cap_per_bucket=int(params.get("price_cap_per_bucket", 1400)),
concurrency=int(params.get("concurrency", 5)),
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
secondary_only=bool(params.get("secondary_only", True)),
resume_run_id=_pick_resume(db, run_id),
incremental_days=None,
)
async def _job_yandex_city_sweep(
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
) -> None:
# #B1 oblast rollout — см. _job_avito_city_sweep. Для yandex city None-anchors
# означает "единственный центральный ЕКБ anchor combos-mode" (run_yandex_city_sweep
# само подставляет (56.8400, 60.6050, "Центр (combos)") при anchors=None); заданный
# city заменяет его на anchors города (та же итерация anchors×combos, 1 anchor).
city = params.get("city")
anchors = get_city_anchors(city) if city else None
await run_yandex_city_sweep(
db,
run_id=run_id,
config=ctx.config,
matcher=ctx.matcher,
enrichment=ctx.enrichment,
shutdown_requested=ctx.shutdown_requested,
proxy_provider=ctx.proxy_provider,
anchors=anchors,
city_slug=city,
pages_per_anchor=int(params.get("pages_per_anchor", 2)),
request_delay_sec=float(params.get("request_delay_sec", 9.0)),
radius_m=int(params.get("radius_m", 1500)),
enrich_address=bool(params.get("enrich_address", True)),
segments=params.get("segments"),
resume_run_id=_pick_resume(db, run_id),
)
async def _job_cian_city_sweep(
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
) -> None:
# #B1 oblast rollout — см. _job_avito_city_sweep.
city = params.get("city")
anchors = get_city_anchors(city) if city else None
await run_cian_city_sweep(
db,
run_id=run_id,
config=ctx.config,
matcher=ctx.matcher,
shutdown_requested=ctx.shutdown_requested,
proxy_provider=ctx.proxy_provider,
anchors=anchors,
city_slug=city,
pages_per_anchor=int(params.get("pages_per_anchor", 3)),
request_delay_sec=float(params.get("request_delay_sec", 5.0)),
radius_m=int(params.get("radius_m", 1500)),
detail_top_n=int(params.get("detail_top_n", 10)),
enrich_houses=bool(params.get("enrich_houses", True)),
newbuilding_only=bool(params.get("newbuilding_only", True)),
# #3074: подхват якорей у оборванного предшественника. Ключ чекпоинта —
# имя якоря, от их количества не зависит, поэтому гарда как у combo-
# чекпоинта яндекса здесь не требуется.
resume_run_id=_pick_resume(db, run_id),
)
async def _job_cian_full_load(
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
) -> None:
await run_cian_full_load(
db,
run_id=run_id,
config=ctx.config,
matcher=ctx.matcher,
shutdown_requested=ctx.shutdown_requested,
proxy_provider=ctx.proxy_provider,
price_cap_per_bucket=int(params.get("price_cap_per_bucket", 1400)),
concurrency=int(params.get("concurrency", 5)),
request_delay_sec=float(params.get("request_delay_sec", 4.0)),
enrich_detail=bool(params.get("enrich_detail", False)),
detail_top_n=int(params.get("detail_top_n", 0)),
resume_run_id=_pick_resume(db, run_id),
)
async def _job_domclick_city_sweep(
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
) -> None:
await run_domclick_city_sweep(
db,
run_id=run_id,
config=ctx.config,
matcher=ctx.matcher,
shutdown_requested=ctx.shutdown_requested,
proxy_provider=ctx.proxy_provider,
city_id=int(params.get("city_id", 4)),
rooms=params.get("rooms"),
pages=int(params.get("pages_per_anchor", 5)),
request_delay_sec=float(params.get("request_delay_sec", 6.0)),
resume_run_id=_pick_resume(db, run_id),
)
def _default_kit_handlers() -> dict[str, Handler]:
"""Встроенные kit-native sweep-обработчики (тело в scraper_kit.orchestration.pipeline).
"*_city_sweep_*" — wildcard-семья per-city schedule'ов вне ЕКБ (B1 oblast rollout,
Свердловская область region 66). `scrape_schedules.source` UNIQUE не даёт нескольких
строк с одним и тем же source ("avito_city_sweep") на разные города — поэтому каждый
oblast-город получает СВОЙ source ("avito_city_sweep_nizhniy_tagil" и т.п.,
data/sql/179_scrape_schedules_seed_oblast_city_sweeps.sql), а `resolve_handler`
резолвит их через wildcard-префикс (тот же механизм, что "deactivate_stale_*") на
ТОТ ЖЕ handler, что и ЕКБ-source — job читает default_params["city"] и подставляет
anchors города (см. _job_avito_city_sweep/_job_cian_city_sweep/_job_yandex_city_sweep).
"""
return {
"avito_city_sweep": Handler(_job_avito_city_sweep, "avito_city_sweep"),
"avito_city_sweep_*": Handler(_job_avito_city_sweep, "avito_city_sweep_*"),
"avito_newbuilding_sweep": Handler(_job_avito_newbuilding_sweep, "avito_newbuilding_sweep"),
"avito_full_load": Handler(_job_avito_full_load, "avito_full_load"),
"avito_full_load_exhaustive": Handler(
_job_avito_full_load_exhaustive, "avito_full_load_exhaustive"
),
"yandex_city_sweep": Handler(_job_yandex_city_sweep, "yandex_city_sweep"),
"yandex_city_sweep_*": Handler(_job_yandex_city_sweep, "yandex_city_sweep_*"),
"cian_city_sweep": Handler(_job_cian_city_sweep, "cian_city_sweep"),
"cian_city_sweep_*": Handler(_job_cian_city_sweep, "cian_city_sweep_*"),
"cian_full_load": Handler(_job_cian_full_load, "cian_full_load"),
"domclick_city_sweep": Handler(_job_domclick_city_sweep, "domclick_city_sweep"),
}
def build_registry(product_handlers: Mapping[str, Handler] | None = None) -> dict[str, Handler]:
"""Собрать реестр source→Handler: kit-native sweeps + инжектированные продуктовые джобы.
`product_handlers` — джобы, тело которых осталось в `app` (rosreestr_dkp / sber_index /
deactivate_stale_* / *_backfill / cadastral_geo_match / house_* / proxy_healthcheck / …).
Их строит app-side wiring и передаёт сюда. Продуктовые ключи МОГУТ переопределять
kit-native (последнее слово за продуктом).
Wildcard-ключи ("deactivate_stale_*") резолвятся `resolve_handler` по префиксу.
"""
registry = dict(_default_kit_handlers())
if product_handlers:
registry.update(product_handlers)
return registry
def get_due_schedules(db: Session) -> list[dict[str, Any]]:
"""SELECT scrape_schedules WHERE enabled AND (next_run_at IS NULL OR next_run_at <= NOW())."""
rows = (
db.execute(
text(
"""
SELECT id, source, enabled, window_start_hour, window_end_hour,
default_params, last_run_id, last_run_at, next_run_at
FROM scrape_schedules
WHERE enabled = true
AND (next_run_at IS NULL OR next_run_at <= NOW())
"""
),
)
.mappings()
.all()
)
return [dict(r) for r in rows]
async def _tick_loop(ctx: SchedulerContext, registry: Mapping[str, Handler]) -> None:
"""Тело tick-лупа: reap → get_due → dispatch, до выхода по SIGTERM-drain'у.
Вынесено из `scheduler_loop` (#3391) ровно затем, чтобы отмену, прилетевшую в ТЕЛО
тика, можно было поймать одним `except` — см. вызывающего.
"""
logger.info("scheduler: started (tick=%ds)", SCHEDULER_TICK_SEC)
boot_time = datetime.now(UTC) # граница «моих» прогонов для boot-reap (#3122)
# Initial sleep 30s чтобы дать FastAPI startup завершиться
await asyncio.sleep(30)
_boot_reap_done = False
while True:
# #1182 Phase 2: кооперативный SIGTERM-drain. Не reap'аем и не claim'аем
# новые run'ы во время shutdown — даём текущему dispatch'у докатиться и выходим.
if ctx.shutdown_requested():
logger.info("scheduler: SIGTERM-drain — stop claiming new runs, exiting tick loop")
break
try:
db = ctx.session_factory()
try:
# #3122: однократно на старте — снять прогоны предыдущего контейнера.
if not _boot_reap_done:
reap_boot_zombies(db, boot_time)
_boot_reap_done = True
# Reap zombies first
reap_zombies(db)
# #2670: раз в сутки — сводка «что сейчас не собирает». Календарная, а
# не по стрику: два самых залежавшихся источника прода имеют стрик 0.
emit_stale_digest(db)
# Process due schedules
due = get_due_schedules(db)
for sch in due:
source = sch["source"]
handler = resolve_handler(source, registry)
if handler is None:
logger.warning("scheduler: unknown source=%s, skip", source)
# #2658: enabled-расписание без handler'а молча не выполнялось
# бы вечно (next_run_at не двигается → warning каждый тик).
_kit_runs.mark_skipped(
db,
source=source,
reason=SKIP_UNKNOWN_SOURCE,
details="нет handler'а в реестре",
)
else:
await _dispatch(handler, db, sch, ctx)
# #1182 Phase 2: после каждого dispatch'а проверяем drain — текущий
# run уже отпущен в свою asyncio-задачу (сам докоммитит/выйдет по
# своему checkpoint'у), а новые в этом тике не запускаем.
if ctx.shutdown_requested():
logger.info(
"scheduler: SIGTERM-drain — stop dispatch after source=%s", source
)
break
finally:
db.close()
except Exception:
logger.exception("scheduler: tick failed")
# Промежуточная проверка перед 60s-сном: выходим сразу, не ждём целый тик.
if ctx.shutdown_requested():
logger.info("scheduler: SIGTERM-drain — exiting tick loop after dispatch")
break
await asyncio.sleep(SCHEDULER_TICK_SEC)
async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler]) -> None:
"""Бесконечный async loop — tick каждые SCHEDULER_TICK_SEC секунд.
Структура (reap → get_due → dispatch → SIGTERM-drain) идентична боевому
scheduler'у; 27-веточный if/elif заменён единственным `resolve_handler` + `_dispatch`.
"""
try:
await _tick_loop(ctx, registry)
except asyncio.CancelledError:
# #3391: hard-cancel из scheduler_main приходит по расписанию grace'а, а не по
# нашему — он вполне может застать ТЕЛО тика (reap / stale-digest / `_dispatch`
# с его сетевым pre_claim), а не `drain_inflight`. `except Exception` тика
# CancelledError не ловит (BaseException), до `await ctx.drain_inflight()` ниже
# выполнение не доходит — без этой ветки in-flight строки остались бы 'running'
# ровно так же, как 07.09 на проде.
ctx.mark_inflight_interrupted([t for t in ctx._inflight_tasks if not t.done()])
raise
# Tick-loop вышел только по SIGTERM-drain'у (иначе while True бесконечен): дожидаемся
# detached run-задач (spawn_tracked) — пусть докоммитят текущий unit и сделают
# mark_done, а не дадим asyncio.run() teardown'у хард-кансельнуть их mid-await (#1182).
await ctx.drain_inflight()
logger.info("scheduler: tick loop exited (in-flight drain complete)")
__all__ = [
"SCHEDULER_TICK_SEC",
"SKIP_ALREADY_RUNNING",
"SKIP_CONCURRENT_CLAIM",
"SKIP_RUNNING_UNDER_LOCK",
"SKIP_UNKNOWN_SOURCE",
"ZOMBIE_THRESHOLD_HOURS",
"Handler",
"SchedulerContext",
"build_registry",
"compute_next_run_at",
"get_due_schedules",
"has_running_run",
"reap_zombies",
"reschedule_after_minutes",
"resolve_handler",
"scheduler_loop",
]