feat(tradein/scraper): cooperative SIGTERM-drain for graceful deploy (#1182 Phase 2)
Деплой (docker recreate tradein-scraper) шлёт SIGTERM с stop_grace_period=120s (Phase 1). Раньше scheduler_main отвечал hard task.cancel() → бегущий scrape-unit получал CancelledError посреди браузер-карточки → карточка гибла без COMMIT. Теперь SIGTERM/SIGINT лишь выставляют кооперативный флаг; tick-loop и длинные задачи опрашивают его на СВОИХ существующих between-unit checkpoint'ах, докоммичивают текущий unit и выходят сами. - app/core/shutdown.py: новый standalone-модуль (asyncio.Event + request_shutdown / shutdown_requested / wait_for_shutdown), без app-зависимостей → нет циклов. - scheduler_main._run: SIGTERM → request_shutdown() вместо task.cancel(); ждём добровольного drain'а, safety-net wait_for(100s < docker grace) с fallback на hard-cancel для некооперирующей задачи. - scheduler_loop: проверка флага в начале тика и после каждого dispatch — не reap/claim новые run'ы во время shutdown, выходим из loop'а. - avito_detail_backfill: break на границе карточки рядом с budget-guard; snapshot pending-query идемпотентен → следующий run сам резюмит остаток. - import_rosreestr_dkp: расширен is_cancelled checkpoint — drain делает mark_done (partial), не mark_cancelled; user-cancel семантика без изменений. Без SIGTERM поведение идентично прежнему. reap_zombies не трогает drained-run'ы (они mark_done, не 'running'). Tests: tests/core/test_shutdown.py, scheduler_main drain+timeout-fallback, avito_detail_backfill partial-drain. 26 + 29 scheduler passed.
This commit is contained in:
parent
77b4f7e64c
commit
848e6d3bdd
8 changed files with 367 additions and 12 deletions
48
tradein-mvp/backend/app/core/shutdown.py
Normal file
48
tradein-mvp/backend/app/core/shutdown.py
Normal file
|
|
@ -0,0 +1,48 @@
|
||||||
|
"""Кооперативный shutdown-флаг для scraper-планировщика (#1182 Phase 2).
|
||||||
|
|
||||||
|
Зачем: деплой (docker recreate `tradein-scraper`) шлёт SIGTERM с
|
||||||
|
`stop_grace_period=120s`. Раньше `scheduler_main` отвечал hard `task.cancel()` →
|
||||||
|
бегущий scrape-unit получал CancelledError посреди браузер-карточки → карточка
|
||||||
|
гибла без COMMIT. Теперь SIGTERM лишь ВЫСТАВЛЯЕТ этот флаг; tick-loop планировщика
|
||||||
|
и длинные задачи опрашивают `shutdown_requested()` на СВОИХ существующих
|
||||||
|
between-unit checkpoint'ах, докоммичивают текущий unit и выходят сами.
|
||||||
|
|
||||||
|
Модуль намеренно standalone (зависит только от `asyncio`): планировщик импортирует
|
||||||
|
задачи, а задачи должны читать флаг — общий нижний слой без app-зависимостей
|
||||||
|
исключает циклический импорт.
|
||||||
|
|
||||||
|
Без SIGTERM `shutdown_requested()` всегда False → поведение идентично прежнему.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
# Module-level Event: создаётся на import'е. В Python 3.10+ asyncio.Event не
|
||||||
|
# привязывается к loop'у в __init__ (loop берётся лениво в .wait()), поэтому
|
||||||
|
# set()/is_set() безопасны даже до запуска event loop'а.
|
||||||
|
_shutdown_event = asyncio.Event()
|
||||||
|
|
||||||
|
|
||||||
|
def request_shutdown() -> None:
|
||||||
|
"""Запросить кооперативный drain. Идемпотентно — повторные вызовы безвредны."""
|
||||||
|
_shutdown_event.set()
|
||||||
|
|
||||||
|
|
||||||
|
def shutdown_requested() -> bool:
|
||||||
|
"""True после первого request_shutdown() (идёт кооперативный drain)."""
|
||||||
|
return _shutdown_event.is_set()
|
||||||
|
|
||||||
|
|
||||||
|
async def wait_for_shutdown() -> None:
|
||||||
|
"""Заблокироваться до request_shutdown(). Нужно safety-net drain'у в scheduler_main."""
|
||||||
|
await _shutdown_event.wait()
|
||||||
|
|
||||||
|
|
||||||
|
def reset_shutdown() -> None:
|
||||||
|
"""Сбросить флаг — ТОЛЬКО для тестов. Пересоздаёт Event (а не .clear()): asyncio.Event
|
||||||
|
кеширует event loop при первом .wait(), а pytest-asyncio даёт каждому тесту НОВЫЙ loop;
|
||||||
|
свежий объект избегает 'bound to a different event loop'. В проде не вызывается
|
||||||
|
(один asyncio.run → один loop), поэтому переприсвоение глобала здесь безопасно."""
|
||||||
|
global _shutdown_event
|
||||||
|
_shutdown_event = asyncio.Event()
|
||||||
|
|
@ -20,6 +20,7 @@ import sys
|
||||||
from contextlib import suppress
|
from contextlib import suppress
|
||||||
|
|
||||||
from app.core.config import settings
|
from app.core.config import settings
|
||||||
|
from app.core.shutdown import request_shutdown, shutdown_requested, wait_for_shutdown
|
||||||
|
|
||||||
logging.basicConfig(
|
logging.basicConfig(
|
||||||
level=logging.INFO,
|
level=logging.INFO,
|
||||||
|
|
@ -28,6 +29,12 @@ logging.basicConfig(
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
# Safety net: после запроса drain'а ждём добровольного выхода задачи не дольше
|
||||||
|
# этого порога, который ОБЯЗАН быть < docker stop_grace_period (120s). Если
|
||||||
|
# некооперативная задача не вышла за это время — hard-cancel, чтобы процесс не
|
||||||
|
# подвис до SIGKILL за пределами grace. Module-level → monkeypatch'абельно в тестах.
|
||||||
|
_DRAIN_TIMEOUT_S = 100.0
|
||||||
|
|
||||||
# Мониторинг ошибок — GlitchTip (Sentry-совместимый, #396).
|
# Мониторинг ошибок — GlitchTip (Sentry-совместимый, #396).
|
||||||
# Только integrations без Starlette/FastAPI (у нас нет ASGI-приложения здесь).
|
# Только integrations без Starlette/FastAPI (у нас нет ASGI-приложения здесь).
|
||||||
# SqlalchemyIntegration + HttpxIntegration + LoggingIntegration покрывают scraper-слой.
|
# SqlalchemyIntegration + HttpxIntegration + LoggingIntegration покрывают scraper-слой.
|
||||||
|
|
@ -60,29 +67,79 @@ def _should_run() -> bool:
|
||||||
return settings.scheduler_enable
|
return settings.scheduler_enable
|
||||||
|
|
||||||
|
|
||||||
|
async def _await_scheduler(task: asyncio.Task[None]) -> None:
|
||||||
|
"""Дождаться завершения scheduler-задачи с кооперативным SIGTERM-drain'ом.
|
||||||
|
|
||||||
|
- Нет shutdown: scheduler_loop бесконечен → задача никогда не завершается, ждём её
|
||||||
|
как есть (процесс просто работает).
|
||||||
|
- shutdown запрошен, задача ещё бежит: bounded `wait_for(_DRAIN_TIMEOUT_S)` —
|
||||||
|
кооперативная задача докоммитит текущий unit на ближайшем checkpoint'е и выйдет
|
||||||
|
сама. Если превысила grace → hard-cancel + suppress CancelledError, чтобы
|
||||||
|
некооперативная задача не подвесила процесс за пределами docker stop_grace_period.
|
||||||
|
"""
|
||||||
|
# Гонка «задача завершилась сама» против «пришёл SIGTERM»: отсчёт safety-net'а
|
||||||
|
# должен стартовать от МОМЕНТА запроса drain'а, а не от старта процесса.
|
||||||
|
shutdown_waiter = asyncio.create_task(wait_for_shutdown())
|
||||||
|
try:
|
||||||
|
await asyncio.wait({task, shutdown_waiter}, return_when=asyncio.FIRST_COMPLETED)
|
||||||
|
finally:
|
||||||
|
shutdown_waiter.cancel()
|
||||||
|
with suppress(asyncio.CancelledError):
|
||||||
|
await shutdown_waiter
|
||||||
|
|
||||||
|
if task.done():
|
||||||
|
# Задача вышла сама (drained на checkpoint'е; без сигнала сюда не попадаем).
|
||||||
|
logger.info("scheduler_main: scheduler task exited cleanly")
|
||||||
|
return
|
||||||
|
|
||||||
|
logger.info(
|
||||||
|
"scheduler_main: SIGTERM-drain — waiting up to %.0fs for in-flight unit to commit",
|
||||||
|
_DRAIN_TIMEOUT_S,
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
await asyncio.wait_for(task, timeout=_DRAIN_TIMEOUT_S)
|
||||||
|
logger.info("scheduler_main: scheduler drained and exited cleanly")
|
||||||
|
except TimeoutError:
|
||||||
|
logger.warning(
|
||||||
|
"scheduler_main: drain exceeded %.0fs grace — hard-cancelling scheduler task",
|
||||||
|
_DRAIN_TIMEOUT_S,
|
||||||
|
)
|
||||||
|
task.cancel()
|
||||||
|
with suppress(asyncio.CancelledError):
|
||||||
|
await task
|
||||||
|
|
||||||
|
|
||||||
async def _run() -> None:
|
async def _run() -> None:
|
||||||
"""Async entrypoint: запустить scheduler_loop() с чистой отменой по SIGTERM/SIGINT."""
|
"""Async entrypoint: scheduler_loop() с кооперативным SIGTERM/SIGINT-drain'ом.
|
||||||
|
|
||||||
|
SIGTERM/SIGINT → request_shutdown() (НЕ task.cancel()): бегущий scrape-unit
|
||||||
|
докоммитит текущую карточку и выйдет сам на ближайшем between-unit checkpoint'е
|
||||||
|
(scheduler_loop / avito_detail_backfill / import_rosreestr_dkp). Safety-net в
|
||||||
|
_await_scheduler гарантирует выход в пределах docker grace.
|
||||||
|
"""
|
||||||
from app.services.scheduler import scheduler_loop
|
from app.services.scheduler import scheduler_loop
|
||||||
|
|
||||||
task = asyncio.create_task(scheduler_loop())
|
task = asyncio.create_task(scheduler_loop())
|
||||||
|
|
||||||
loop = asyncio.get_running_loop()
|
loop = asyncio.get_running_loop()
|
||||||
|
|
||||||
def _cancel(_signum: int, _frame: object = None) -> None:
|
def _on_signal(signum: int) -> None:
|
||||||
logger.info("scheduler_main: signal received — cancelling scheduler task")
|
logger.info("scheduler_main: signal %d received — requesting cooperative drain", signum)
|
||||||
task.cancel()
|
request_shutdown()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
loop.add_signal_handler(signal.SIGTERM, lambda: _cancel(signal.SIGTERM))
|
loop.add_signal_handler(signal.SIGTERM, lambda: _on_signal(signal.SIGTERM))
|
||||||
loop.add_signal_handler(signal.SIGINT, lambda: _cancel(signal.SIGINT))
|
loop.add_signal_handler(signal.SIGINT, lambda: _on_signal(signal.SIGINT))
|
||||||
except NotImplementedError:
|
except NotImplementedError:
|
||||||
# Windows dev: signal handlers через loop не поддерживаются
|
# Windows dev: signal handlers через loop не поддерживаются
|
||||||
logger.warning("scheduler_main: loop.add_signal_handler not supported (Windows dev)")
|
logger.warning("scheduler_main: loop.add_signal_handler not supported (Windows dev)")
|
||||||
|
|
||||||
with suppress(asyncio.CancelledError):
|
await _await_scheduler(task)
|
||||||
await task
|
|
||||||
|
|
||||||
logger.info("scheduler_main: scheduler cancelled cleanly")
|
if shutdown_requested():
|
||||||
|
logger.info("scheduler_main: scheduler drained cleanly (SIGTERM)")
|
||||||
|
else:
|
||||||
|
logger.info("scheduler_main: scheduler task exited")
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
|
|
|
||||||
|
|
@ -140,6 +140,7 @@ from sqlalchemy import text
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
from app.core.db import SessionLocal
|
from app.core.db import SessionLocal
|
||||||
|
from app.core.shutdown import shutdown_requested
|
||||||
from app.services import scrape_runs as runs_mod
|
from app.services import scrape_runs as runs_mod
|
||||||
from app.services.scrape_pipeline import (
|
from app.services.scrape_pipeline import (
|
||||||
run_avito_city_sweep,
|
run_avito_city_sweep,
|
||||||
|
|
@ -1319,9 +1320,28 @@ def import_rosreestr_dkp(
|
||||||
|
|
||||||
try:
|
try:
|
||||||
while True:
|
while True:
|
||||||
if runs_mod.is_cancelled(db, run_id):
|
cancelled = runs_mod.is_cancelled(db, run_id)
|
||||||
logger.info("rosreestr_dkp_import run_id=%d: cancelled by user", run_id)
|
if cancelled or shutdown_requested():
|
||||||
runs_mod.mark_cancelled(db, run_id)
|
if cancelled:
|
||||||
|
# User-cancel: семантика без изменений — mark_cancelled.
|
||||||
|
logger.info("rosreestr_dkp_import run_id=%d: cancelled by user", run_id)
|
||||||
|
runs_mod.mark_cancelled(db, run_id)
|
||||||
|
return
|
||||||
|
# #1182 Phase 2: кооперативный SIGTERM-drain (деплой recreate scraper).
|
||||||
|
# Это НЕ user-cancel → mark_done (partial), не mark_cancelled. Курсор
|
||||||
|
# уже зафиксирован update_heartbeat'ом каждый батч; следующий run
|
||||||
|
# пере-сканирует с id=0 (INSERT ... ON CONFLICT DO NOTHING идемпотентен)
|
||||||
|
# → уже вставленные строки пропускаются, остаток до-импортируется.
|
||||||
|
# mark_done выводит run из 'running' → reap_zombies его не тронет.
|
||||||
|
logger.info(
|
||||||
|
"rosreestr_dkp_import run_id=%d: SIGTERM-drain — committing partial "
|
||||||
|
"(last_id=%d, batches=%d) and exiting",
|
||||||
|
run_id,
|
||||||
|
last_id,
|
||||||
|
total_batches,
|
||||||
|
)
|
||||||
|
runs_mod.update_heartbeat(db, run_id, counters)
|
||||||
|
runs_mod.mark_done(db, run_id, counters)
|
||||||
return
|
return
|
||||||
|
|
||||||
# Cursor-based pagination via foreign table gendesign_rosreestr_deals.
|
# Cursor-based pagination via foreign table gendesign_rosreestr_deals.
|
||||||
|
|
@ -1742,6 +1762,11 @@ async def scheduler_loop() -> None:
|
||||||
# Initial sleep 30s чтобы дать FastAPI startup завершиться
|
# Initial sleep 30s чтобы дать FastAPI startup завершиться
|
||||||
await asyncio.sleep(30)
|
await asyncio.sleep(30)
|
||||||
while True:
|
while True:
|
||||||
|
# #1182 Phase 2: кооперативный SIGTERM-drain. Не reap'аем и не claim'аем
|
||||||
|
# новые run'ы во время shutdown — даём текущему dispatch'у докатиться и выходим.
|
||||||
|
if shutdown_requested():
|
||||||
|
logger.info("scheduler: SIGTERM-drain — stop claiming new runs, exiting tick loop")
|
||||||
|
break
|
||||||
try:
|
try:
|
||||||
db = SessionLocal()
|
db = SessionLocal()
|
||||||
try:
|
try:
|
||||||
|
|
@ -1803,8 +1828,21 @@ async def scheduler_loop() -> None:
|
||||||
await trigger_house_dedup_merge_run(db, sch)
|
await trigger_house_dedup_merge_run(db, sch)
|
||||||
else:
|
else:
|
||||||
logger.warning("scheduler: unknown source=%s, skip", source)
|
logger.warning("scheduler: unknown source=%s, skip", source)
|
||||||
|
|
||||||
|
# #1182 Phase 2: после каждого dispatch'а проверяем drain — текущий
|
||||||
|
# run уже отпущен в свою asyncio-задачу (сам докоммитит/выйдет по
|
||||||
|
# своему checkpoint'у), а новые в этом тике не запускаем.
|
||||||
|
if shutdown_requested():
|
||||||
|
logger.info(
|
||||||
|
"scheduler: SIGTERM-drain — stop dispatch after source=%s", source
|
||||||
|
)
|
||||||
|
break
|
||||||
finally:
|
finally:
|
||||||
db.close()
|
db.close()
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("scheduler: tick failed")
|
logger.exception("scheduler: tick failed")
|
||||||
|
# Промежуточная проверка перед 60s-сном: выходим сразу, не ждём целый тик.
|
||||||
|
if shutdown_requested():
|
||||||
|
logger.info("scheduler: SIGTERM-drain — exiting tick loop after dispatch")
|
||||||
|
break
|
||||||
await asyncio.sleep(SCHEDULER_TICK_SEC)
|
await asyncio.sleep(SCHEDULER_TICK_SEC)
|
||||||
|
|
|
||||||
|
|
@ -26,6 +26,7 @@ from sqlalchemy import text
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
from app.core.config import settings
|
from app.core.config import settings
|
||||||
|
from app.core.shutdown import shutdown_requested
|
||||||
from app.services import scrape_runs as runs_mod
|
from app.services import scrape_runs as runs_mod
|
||||||
from app.services.scrape_pipeline import _CHROME_HEADERS, _avito_proxies
|
from app.services.scrape_pipeline import _CHROME_HEADERS, _avito_proxies
|
||||||
from app.services.scrapers.avito import AvitoScraper
|
from app.services.scrapers.avito import AvitoScraper
|
||||||
|
|
@ -212,6 +213,20 @@ async def run_avito_detail_backfill(
|
||||||
)
|
)
|
||||||
break
|
break
|
||||||
|
|
||||||
|
# #1182 Phase 2: кооперативный SIGTERM-drain (деплой recreate scraper).
|
||||||
|
# Чисто выходим на границе карточки — каждая карточка уже закоммичена,
|
||||||
|
# mark_done с текущими счётчиками вызовется после loop'а. Snapshot —
|
||||||
|
# pending-query (detail_enriched_at IS NULL), per-item идемпотентна →
|
||||||
|
# следующий run сам до-резюмит остаток, resume_cursor не нужен.
|
||||||
|
if shutdown_requested():
|
||||||
|
logger.info(
|
||||||
|
"avito_detail_backfill: run_id=%d SIGTERM-drain — stopping at #%d/%d",
|
||||||
|
run_id,
|
||||||
|
idx,
|
||||||
|
len(snapshot),
|
||||||
|
)
|
||||||
|
break
|
||||||
|
|
||||||
# Delay before each request except the first
|
# Delay before each request except the first
|
||||||
if do_sleep:
|
if do_sleep:
|
||||||
await asyncio.sleep(request_delay_sec)
|
await asyncio.sleep(request_delay_sec)
|
||||||
|
|
|
||||||
0
tradein-mvp/backend/tests/core/__init__.py
Normal file
0
tradein-mvp/backend/tests/core/__init__.py
Normal file
62
tradein-mvp/backend/tests/core/test_shutdown.py
Normal file
62
tradein-mvp/backend/tests/core/test_shutdown.py
Normal file
|
|
@ -0,0 +1,62 @@
|
||||||
|
"""Тесты для app/core/shutdown.py — кооперативный SIGTERM-флаг (#1182 Phase 2)."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from app.core import shutdown as sd
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(autouse=True)
|
||||||
|
def _reset() -> None:
|
||||||
|
"""Event — module-global: сбрасываем до и после каждого теста для изоляции."""
|
||||||
|
sd.reset_shutdown()
|
||||||
|
yield
|
||||||
|
sd.reset_shutdown()
|
||||||
|
|
||||||
|
|
||||||
|
def test_default_not_requested() -> None:
|
||||||
|
"""По умолчанию shutdown не запрошен."""
|
||||||
|
assert sd.shutdown_requested() is False
|
||||||
|
|
||||||
|
|
||||||
|
def test_request_sets_flag() -> None:
|
||||||
|
"""request_shutdown() выставляет флаг."""
|
||||||
|
sd.request_shutdown()
|
||||||
|
assert sd.shutdown_requested() is True
|
||||||
|
|
||||||
|
|
||||||
|
def test_request_is_idempotent() -> None:
|
||||||
|
"""Повторный request_shutdown() безвреден — флаг остаётся True."""
|
||||||
|
sd.request_shutdown()
|
||||||
|
sd.request_shutdown()
|
||||||
|
assert sd.shutdown_requested() is True
|
||||||
|
|
||||||
|
|
||||||
|
def test_reset_clears_flag() -> None:
|
||||||
|
"""reset_shutdown() (test-only) снимает флаг."""
|
||||||
|
sd.request_shutdown()
|
||||||
|
assert sd.shutdown_requested() is True
|
||||||
|
sd.reset_shutdown()
|
||||||
|
assert sd.shutdown_requested() is False
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_wait_for_shutdown_blocks_until_requested() -> None:
|
||||||
|
"""wait_for_shutdown() висит, пока не вызван request_shutdown()."""
|
||||||
|
waiter = asyncio.create_task(sd.wait_for_shutdown())
|
||||||
|
await asyncio.sleep(0.02)
|
||||||
|
assert not waiter.done(), "не должен завершиться до request_shutdown()"
|
||||||
|
|
||||||
|
sd.request_shutdown()
|
||||||
|
await asyncio.wait_for(waiter, timeout=1.0)
|
||||||
|
assert waiter.done()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_wait_for_shutdown_returns_immediately_if_already_set() -> None:
|
||||||
|
"""Если флаг уже выставлен, wait_for_shutdown() возвращается сразу."""
|
||||||
|
sd.request_shutdown()
|
||||||
|
await asyncio.wait_for(sd.wait_for_shutdown(), timeout=1.0)
|
||||||
|
|
@ -12,11 +12,21 @@ sys.modules.setdefault("weasyprint", _wp_mock)
|
||||||
|
|
||||||
import pytest # noqa: E402
|
import pytest # noqa: E402
|
||||||
|
|
||||||
|
from app.core import shutdown as _sd # noqa: E402
|
||||||
from app.tasks.avito_detail_backfill import ( # noqa: E402
|
from app.tasks.avito_detail_backfill import ( # noqa: E402
|
||||||
AvitoDetailBackfillResult,
|
AvitoDetailBackfillResult,
|
||||||
run_avito_detail_backfill,
|
run_avito_detail_backfill,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(autouse=True)
|
||||||
|
def _reset_shutdown() -> None:
|
||||||
|
"""shutdown — module-global Event: чистим вокруг каждого теста (изоляция #1182)."""
|
||||||
|
_sd.reset_shutdown()
|
||||||
|
yield
|
||||||
|
_sd.reset_shutdown()
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
# Helpers
|
# Helpers
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
@ -42,6 +52,7 @@ _SLEEP = "app.tasks.avito_detail_backfill.asyncio.sleep"
|
||||||
_SESSION = "app.tasks.avito_detail_backfill.AsyncSession"
|
_SESSION = "app.tasks.avito_detail_backfill.AsyncSession"
|
||||||
_SCRAPER = "app.tasks.avito_detail_backfill.AvitoScraper"
|
_SCRAPER = "app.tasks.avito_detail_backfill.AvitoScraper"
|
||||||
_SETTINGS = "app.tasks.avito_detail_backfill.settings"
|
_SETTINGS = "app.tasks.avito_detail_backfill.settings"
|
||||||
|
_SHUTDOWN = "app.tasks.avito_detail_backfill.shutdown_requested"
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
# Tests
|
# Tests
|
||||||
|
|
@ -139,6 +150,48 @@ async def test_backfill_blocked_abort_after_max_consecutive() -> None:
|
||||||
assert mock_scraper.return_value._rotate_ip.call_count == 5
|
assert mock_scraper.return_value._rotate_ip.call_count == 5
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_backfill_sigterm_drain_breaks_and_marks_done_partial() -> None:
|
||||||
|
"""#1182 Phase 2: shutdown_requested() True → loop выходит на границе карточки,
|
||||||
|
mark_done вызывается с ЧАСТИЧНЫМИ счётчиками (не mark_failed, не mark_cancelled).
|
||||||
|
|
||||||
|
side_effect [False, True]: 1-я карточка обрабатывается (attempted=1, enriched=1),
|
||||||
|
перед 2-й приходит SIGTERM-drain → break. Snapshot — pending-query, остаток
|
||||||
|
до-резюмит следующий run (resume_cursor не нужен).
|
||||||
|
"""
|
||||||
|
snapshot = _make_snapshot(3)
|
||||||
|
db = _mock_db(snapshot)
|
||||||
|
runs = MagicMock()
|
||||||
|
mock_enrichment = MagicMock()
|
||||||
|
mock_fetch = AsyncMock(return_value=mock_enrichment)
|
||||||
|
mock_save = MagicMock(return_value=True)
|
||||||
|
fake_settings = MagicMock(scraper_fetch_mode="cffi")
|
||||||
|
with (
|
||||||
|
patch(_SETTINGS, fake_settings),
|
||||||
|
patch(_SESSION, return_value=AsyncMock()),
|
||||||
|
patch(_SCRAPER),
|
||||||
|
patch(_RUNS, runs),
|
||||||
|
patch(_FETCH, mock_fetch),
|
||||||
|
patch(_SAVE, mock_save),
|
||||||
|
patch(_SLEEP, new_callable=AsyncMock),
|
||||||
|
patch(_SHUTDOWN, side_effect=[False, True]),
|
||||||
|
):
|
||||||
|
result = await run_avito_detail_backfill(
|
||||||
|
db, run_id=13, params={"batch_size": 10, "budget_sec": 3600}
|
||||||
|
)
|
||||||
|
|
||||||
|
# Обработана только 1-я карточка, на 2-й — drain-break.
|
||||||
|
assert result.attempted == 1
|
||||||
|
assert result.enriched == 1
|
||||||
|
assert mock_fetch.call_count == 1
|
||||||
|
runs.mark_done.assert_called_once()
|
||||||
|
runs.mark_failed.assert_not_called()
|
||||||
|
runs.mark_cancelled.assert_not_called()
|
||||||
|
# mark_done получил ЧАСТИЧНЫЕ счётчики (attempted=1, а не весь snapshot=3).
|
||||||
|
done_counters = runs.mark_done.call_args.args[2]
|
||||||
|
assert done_counters["attempted"] == 1
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_backfill_budget_guard_stops_loop() -> None:
|
async def test_backfill_budget_guard_stops_loop() -> None:
|
||||||
"""Budget expired before first listing -> fetch_detail not called."""
|
"""Budget expired before first listing -> fetch_detail not called."""
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,8 @@
|
||||||
(a) модуль импортируется без реального DB/сети.
|
(a) модуль импортируется без реального DB/сети.
|
||||||
(b) _should_run() возвращает False при SCHEDULER_ENABLE=false.
|
(b) _should_run() возвращает False при SCHEDULER_ENABLE=false.
|
||||||
(c) _run() отменяется чисто без unhandled exceptions.
|
(c) _run() отменяется чисто без unhandled exceptions.
|
||||||
|
(d) #1182 Phase 2: SIGTERM → request_shutdown() (не task.cancel); задача
|
||||||
|
кооперативно drain'ится; safety-net wait_for-timeout добивает зависшую.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
@ -15,6 +17,15 @@ os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
import app.scheduler_main as sm
|
import app.scheduler_main as sm
|
||||||
|
from app.core import shutdown as sd
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(autouse=True)
|
||||||
|
def _reset_shutdown() -> None:
|
||||||
|
"""shutdown — module-global Event: сбрасываем вокруг каждого теста для изоляции."""
|
||||||
|
sd.reset_shutdown()
|
||||||
|
yield
|
||||||
|
sd.reset_shutdown()
|
||||||
|
|
||||||
|
|
||||||
# (a) Импорт модуля
|
# (a) Импорт модуля
|
||||||
|
|
@ -70,3 +81,74 @@ async def test_run_cancels_cleanly(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
|
||||||
# Тест прошёл если не было unhandled exception
|
# Тест прошёл если не было unhandled exception
|
||||||
assert task.done()
|
assert task.done()
|
||||||
|
|
||||||
|
|
||||||
|
# (d) #1182 Phase 2 — кооперативный SIGTERM-drain
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_await_scheduler_drains_on_shutdown() -> None:
|
||||||
|
"""_await_scheduler: кооперативная задача видит флаг и выходит сама — без cancel."""
|
||||||
|
drained = False
|
||||||
|
|
||||||
|
async def _coop_loop() -> None:
|
||||||
|
nonlocal drained
|
||||||
|
while not sd.shutdown_requested():
|
||||||
|
await asyncio.sleep(0.01)
|
||||||
|
drained = True # докоммитил unit и вышел
|
||||||
|
|
||||||
|
task = asyncio.create_task(_coop_loop())
|
||||||
|
runner = asyncio.create_task(sm._await_scheduler(task))
|
||||||
|
|
||||||
|
await asyncio.sleep(0.03)
|
||||||
|
assert not runner.done(), "без shutdown _await_scheduler не завершается"
|
||||||
|
|
||||||
|
sd.request_shutdown()
|
||||||
|
await asyncio.wait_for(runner, timeout=2.0)
|
||||||
|
|
||||||
|
assert drained is True
|
||||||
|
assert task.done()
|
||||||
|
assert not task.cancelled(), "кооперативный выход — НЕ hard-cancel"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_await_scheduler_timeout_hard_cancels(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
"""_await_scheduler: некооперативная задача (игнорит флаг) добивается cancel'ом
|
||||||
|
по истечении safety-net grace (_DRAIN_TIMEOUT_S)."""
|
||||||
|
monkeypatch.setattr(sm, "_DRAIN_TIMEOUT_S", 0.05)
|
||||||
|
|
||||||
|
async def _stuck_loop() -> None:
|
||||||
|
# Игнорирует shutdown_requested() — висит вечно.
|
||||||
|
await asyncio.Event().wait()
|
||||||
|
|
||||||
|
task = asyncio.create_task(_stuck_loop())
|
||||||
|
sd.request_shutdown() # drain уже запрошен
|
||||||
|
|
||||||
|
await asyncio.wait_for(sm._await_scheduler(task), timeout=2.0)
|
||||||
|
|
||||||
|
assert task.done()
|
||||||
|
assert task.cancelled(), "превысила grace → hard-cancel"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_run_uses_request_shutdown_not_cancel(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
"""_run(): когда выставлен shutdown, scheduler_loop сам drain'ится и _run выходит,
|
||||||
|
БЕЗ обращения к task.cancel() (кооперативный путь, не hard-cancel)."""
|
||||||
|
started = asyncio.Event()
|
||||||
|
|
||||||
|
async def _coop_loop() -> None:
|
||||||
|
started.set()
|
||||||
|
while not sd.shutdown_requested():
|
||||||
|
await asyncio.sleep(0.01)
|
||||||
|
|
||||||
|
import app.services.scheduler as _sched_mod
|
||||||
|
|
||||||
|
monkeypatch.setattr(_sched_mod, "scheduler_loop", _coop_loop)
|
||||||
|
|
||||||
|
run_task = asyncio.create_task(sm._run())
|
||||||
|
await asyncio.wait_for(started.wait(), timeout=2.0)
|
||||||
|
|
||||||
|
# Имитируем то, что делает SIGTERM-хендлер (на Windows add_signal_handler нет).
|
||||||
|
sm.request_shutdown()
|
||||||
|
|
||||||
|
await asyncio.wait_for(run_task, timeout=2.0)
|
||||||
|
assert run_task.done()
|
||||||
|
assert run_task.exception() is None
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue