Merge pull request 'feat(tradein/scheduler): cooperative SIGTERM-drain — in-flight scrape unit commits before exit (phase 2)' (#2070) from feat/tradein-scheduler-graceful-drain into main
All checks were successful
Deploy Trade-In / changes (push) Successful in 10s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 1m27s
Deploy Trade-In / build-backend (push) Successful in 1m38s
Deploy Trade-In / deploy (push) Successful in 51s

This commit is contained in:
lekss361 2026-06-28 17:20:11 +00:00
commit 03ca46fe1d
12 changed files with 567 additions and 81 deletions

View 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()

View file

@ -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,82 @@ 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'е; без сигнала сюда не попадаем).
# .result() пробрасывает реальное исключение, если scheduler_loop неожиданно
# упал (сохраняем прежнюю loud-crash семантику `await task`), а не глушит его.
task.result()
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__":

View file

@ -132,6 +132,7 @@ from __future__ import annotations
import asyncio import asyncio
import logging import logging
import random import random
from collections.abc import Coroutine
from datetime import UTC, datetime, time, timedelta from datetime import UTC, datetime, time, timedelta
from typing import Any from typing import Any
@ -140,6 +141,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,
@ -157,6 +159,63 @@ logger = logging.getLogger(__name__)
SCHEDULER_TICK_SEC = 60 SCHEDULER_TICK_SEC = 60
ZOMBIE_THRESHOLD_HOURS = 6 ZOMBIE_THRESHOLD_HOURS = 6
# #1182 Phase 2: graceful SIGTERM-drain. Scrape-работа крутится в detached
# asyncio.create_task (см. _spawn_tracked) — coordinator (scheduler_loop) их
# НЕ join'ит. Без явного дренажа asyncio.run() teardown хард-кансельнул бы их
# mid-await (тот самый #1182 failure mode). Регистр strong-ref'ов даёт coordinator'у
# дождаться детей перед выходом из tick-loop'а.
_inflight_tasks: set[asyncio.Task[None]] = set()
# Бюджет дренажа детей < scheduler_main._DRAIN_TIMEOUT_S (100s) < docker grace (120s):
# coordinator успевает дождаться кооперативных детей (avito_detail_backfill,
# rosreestr-executor дочекивают карточку/батч + mark_done и резолвятся) ВНУТРИ
# внешнего wait_for, а некооперативные упираются в этот timeout и падают на внешний
# hard-cancel fallback.
_CHILD_DRAIN_TIMEOUT_S = 80.0
def _spawn_tracked(coro: Coroutine[Any, Any, None]) -> asyncio.Task[None]:
"""create_task + регистрация задачи в _inflight_tasks для graceful-drain'а (#1182 P2).
Заменяет прежний `task = create_task(_run()); task.add_done_callback(...)`:
strong-ref в set'е держит задачу до завершения (RUF006) и даёт scheduler_loop'у
дождаться её на SIGTERM-drain'е; done-callback ретривит exception (чтобы не было
'Task exception was never retrieved') и убирает задачу из set'а.
"""
task = asyncio.create_task(coro)
_inflight_tasks.add(task)
def _on_done(t: asyncio.Task[None]) -> None:
_inflight_tasks.discard(t)
if not t.cancelled():
t.exception()
task.add_done_callback(_on_done)
return task
async def _drain_inflight() -> None:
"""Дождаться завершения detached run-задач перед teardown'ом процесса (#1182 P2).
Зовётся из scheduler_loop при выходе из tick-loop'а по SIGTERM-drain'у. Кооперативные
дети дочекивают текущий unit (карточку/батч), делают mark_done и резолвятся;
некооперативные упираются в _CHILD_DRAIN_TIMEOUT_S и остаются на внешний hard-cancel
fallback (scheduler_main 100s + docker 120s). Не busy-spin один asyncio.wait.
"""
pending = [t for t in _inflight_tasks if not t.done()]
if not pending:
return
logger.info("scheduler: draining %d in-flight run task(s) on shutdown", len(pending))
_done, still = await asyncio.wait(pending, timeout=_CHILD_DRAIN_TIMEOUT_S)
if still:
logger.warning(
"scheduler: %d task(s) did not drain in %.0fs — leaving for hard-cancel",
len(still),
_CHILD_DRAIN_TIMEOUT_S,
)
else:
logger.info("scheduler: all in-flight run task(s) drained cleanly")
def compute_next_run_at( def compute_next_run_at(
window_start_hour: int, window_start_hour: int,
@ -368,9 +427,7 @@ async def trigger_avito_city_sweep_run(db: Session, schedule_row: dict[str, Any]
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
# Keep reference to avoid GC before task completes (RUF006)
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered city_sweep run_id=%d", run_id) logger.info("scheduler: triggered city_sweep run_id=%d", run_id)
return run_id return run_id
@ -406,9 +463,7 @@ async def trigger_avito_newbuilding_sweep_run(
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
# Keep reference to avoid GC before task completes (RUF006)
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered newbuilding_sweep run_id=%d", run_id) logger.info("scheduler: triggered newbuilding_sweep run_id=%d", run_id)
return run_id return run_id
@ -446,9 +501,7 @@ async def trigger_yandex_city_sweep_run(db: Session, schedule_row: dict[str, Any
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
# Keep reference to avoid GC before task completes (RUF006)
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered yandex_city_sweep run_id=%d", run_id) logger.info("scheduler: triggered yandex_city_sweep run_id=%d", run_id)
return run_id return run_id
@ -506,8 +559,7 @@ async def trigger_cian_backfill_run(db: Session, schedule_row: dict[str, Any]) -
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered cian_history_backfill run_id=%d", run_id) logger.info("scheduler: triggered cian_history_backfill run_id=%d", run_id)
return run_id return run_id
@ -542,8 +594,7 @@ async def trigger_cian_city_sweep_run(db: Session, schedule_row: dict[str, Any])
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered cian_city_sweep run_id=%d", run_id) logger.info("scheduler: triggered cian_city_sweep run_id=%d", run_id)
return run_id return run_id
@ -578,8 +629,7 @@ async def trigger_cian_full_load_run(db: Session, schedule_row: dict[str, Any])
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered cian_full_load run_id=%d", run_id) logger.info("scheduler: triggered cian_full_load run_id=%d", run_id)
return run_id return run_id
@ -624,8 +674,7 @@ async def trigger_avito_full_load_run(db: Session, schedule_row: dict[str, Any])
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered avito_full_load run_id=%d", run_id) logger.info("scheduler: triggered avito_full_load run_id=%d", run_id)
return run_id return run_id
@ -668,8 +717,7 @@ async def trigger_avito_full_load_exhaustive_run(
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered avito_full_load_exhaustive run_id=%d", run_id) logger.info("scheduler: triggered avito_full_load_exhaustive run_id=%d", run_id)
return run_id return run_id
@ -705,8 +753,7 @@ async def trigger_domclick_city_sweep_run(db: Session, schedule_row: dict[str, A
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered domclick_city_sweep run_id=%d", run_id) logger.info("scheduler: triggered domclick_city_sweep run_id=%d", run_id)
return run_id return run_id
@ -809,8 +856,7 @@ async def trigger_yandex_address_backfill_run(
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered yandex_address_backfill run_id=%d", run_id) logger.info("scheduler: triggered yandex_address_backfill run_id=%d", run_id)
return run_id return run_id
@ -848,8 +894,7 @@ async def trigger_geocode_missing_listings_run(
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered geocode_missing_listings run_id=%d", run_id) logger.info("scheduler: triggered geocode_missing_listings run_id=%d", run_id)
return run_id return run_id
@ -873,8 +918,7 @@ async def trigger_rosreestr_dkp_run(db: Session, schedule_row: dict[str, Any]) -
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered rosreestr_dkp_import run_id=%d", run_id) logger.info("scheduler: triggered rosreestr_dkp_import run_id=%d", run_id)
return run_id return run_id
@ -908,8 +952,7 @@ async def trigger_listing_source_snapshot_run(
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered listing_source_snapshot run_id=%d", run_id) logger.info("scheduler: triggered listing_source_snapshot run_id=%d", run_id)
return run_id return run_id
@ -953,8 +996,7 @@ async def trigger_refresh_search_matview_run(
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered refresh_search_matview run_id=%d", run_id) logger.info("scheduler: triggered refresh_search_matview run_id=%d", run_id)
return run_id return run_id
@ -988,8 +1030,7 @@ async def trigger_deactivate_stale_avito_run(
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered deactivate_stale_avito run_id=%d", run_id) logger.info("scheduler: triggered deactivate_stale_avito run_id=%d", run_id)
return run_id return run_id
@ -1041,8 +1082,7 @@ async def trigger_deactivate_stale_run(db: Session, schedule_row: dict[str, Any]
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered deactivate_stale source=%s run_id=%d", listing_source, run_id) logger.info("scheduler: triggered deactivate_stale source=%s run_id=%d", listing_source, run_id)
return run_id return run_id
@ -1077,8 +1117,7 @@ async def trigger_sber_index_pull_run(db: Session, schedule_row: dict[str, Any])
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered sber_index_pull run_id=%d", run_id) logger.info("scheduler: triggered sber_index_pull run_id=%d", run_id)
return run_id return run_id
@ -1115,8 +1154,7 @@ async def trigger_rosreestr_quarter_poll_run(
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered rosreestr_quarter_poll run_id=%d", run_id) logger.info("scheduler: triggered rosreestr_quarter_poll run_id=%d", run_id)
return run_id return run_id
@ -1158,8 +1196,7 @@ async def trigger_newbuilding_enrich_run(db: Session, schedule_row: dict[str, An
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered newbuilding_enrich run_id=%d", run_id) logger.info("scheduler: triggered newbuilding_enrich run_id=%d", run_id)
return run_id return run_id
@ -1207,8 +1244,7 @@ async def trigger_yandex_newbuilding_sweep_run(
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered yandex_newbuilding_sweep run_id=%d", run_id) logger.info("scheduler: triggered yandex_newbuilding_sweep run_id=%d", run_id)
return run_id return run_id
@ -1242,8 +1278,7 @@ async def trigger_asking_to_sold_ratio_run(db: Session, schedule_row: dict[str,
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered asking_to_sold_ratio_refresh run_id=%d", run_id) logger.info("scheduler: triggered asking_to_sold_ratio_refresh run_id=%d", run_id)
return run_id return run_id
@ -1319,9 +1354,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.
@ -1511,8 +1565,7 @@ async def trigger_avito_detail_backfill_run(
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered avito_detail_backfill run_id=%d", run_id) logger.info("scheduler: triggered avito_detail_backfill run_id=%d", run_id)
return run_id return run_id
@ -1548,8 +1601,7 @@ async def trigger_yandex_detail_backfill_run(
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered yandex_detail_backfill run_id=%d", run_id) logger.info("scheduler: triggered yandex_detail_backfill run_id=%d", run_id)
return run_id return run_id
@ -1590,8 +1642,7 @@ async def trigger_cadastral_geo_match_run(db: Session, schedule_row: dict[str, A
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered cadastral_geo_match run_id=%d", run_id) logger.info("scheduler: triggered cadastral_geo_match run_id=%d", run_id)
return run_id return run_id
@ -1663,8 +1714,7 @@ async def trigger_house_imv_backfill_run(db: Session, schedule_row: dict[str, An
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered house_imv_backfill run_id=%d", run_id) logger.info("scheduler: triggered house_imv_backfill run_id=%d", run_id)
return run_id return run_id
@ -1710,8 +1760,7 @@ async def trigger_house_dedup_merge_run(db: Session, schedule_row: dict[str, Any
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) _spawn_tracked(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered house_dedup_merge run_id=%d", run_id) logger.info("scheduler: triggered house_dedup_merge run_id=%d", run_id)
return run_id return run_id
@ -1742,6 +1791,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 +1857,27 @@ 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)
# Tick-loop вышел только по SIGTERM-drain'у (иначе while True бесконечен): дожидаемся
# detached run-задач (_spawn_tracked) — пусть докоммитят текущий unit и сделают
# mark_done, а не дадим asyncio.run() teardown'у хард-кансельнуть их mid-await (#1182).
await _drain_inflight()
logger.info("scheduler: tick loop exited (in-flight drain complete)")

View file

@ -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)

View 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)

View file

@ -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."""

View file

@ -4,6 +4,7 @@ No real network, no real Postgres. The Cian fetch is mocked; the DB session is a
lightweight in-memory fake that models just enough of the three target tables to lightweight in-memory fake that models just enough of the three target tables to
assert (a) rows land and (b) an idempotent re-run does not duplicate them. assert (a) rows land and (b) an idempotent re-run does not duplicate them.
""" """
from __future__ import annotations from __future__ import annotations
import os import os
@ -636,7 +637,9 @@ def test_trigger_claims_run_and_delegates_to_task() -> None:
trig_src = inspect.getsource(scheduler.trigger_newbuilding_enrich_run) trig_src = inspect.getsource(scheduler.trigger_newbuilding_enrich_run)
assert "_claim_run(db, schedule_row)" in trig_src assert "_claim_run(db, schedule_row)" in trig_src
assert "run_newbuilding_enrich" in trig_src assert "run_newbuilding_enrich" in trig_src
assert "asyncio.create_task(" in trig_src # #1182 P2: detached spawn централизован в _spawn_tracked (внутри — asyncio.create_task
# + registry strong-ref для graceful SIGTERM-drain).
assert "_spawn_tracked(_run())" in trig_src
@pytest.mark.asyncio @pytest.mark.asyncio

View file

@ -396,9 +396,10 @@ def test_trigger_claims_run_and_runs_in_executor() -> None:
# sync DB-only task → run_in_executor (mirror cadastral_geo_match, not a bare await). # sync DB-only task → run_in_executor (mirror cadastral_geo_match, not a bare await).
assert "run_in_executor" in src assert "run_in_executor" in src
assert "run_house_dedup_merge" in src assert "run_house_dedup_merge" in src
# fresh session inside the spawned task + RUF006 keep-alive callback. # fresh session inside the spawned task + #1182 P2 detached spawn via _spawn_tracked
# (registry strong-ref = RUF006 keep-alive + graceful SIGTERM-drain).
assert "run_db = SessionLocal()" in src assert "run_db = SessionLocal()" in src
assert "task.add_done_callback(" in src assert "_spawn_tracked(_run())" in src
assert "run_db.close()" in src assert "run_db.close()" in src
@ -535,9 +536,7 @@ def test_real_merge_repoints_dedups_deletes_and_is_idempotent() -> None:
).all() ).all()
assert [r.id for r in survivors] == [900001] assert [r.id for r in survivors] == [900001]
# loser's listing re-pointed to keeper. # loser's listing re-pointed to keeper.
repointed = db.execute( repointed = db.execute(_t("SELECT house_id_fk FROM listings WHERE id = 910002")).scalar()
_t("SELECT house_id_fk FROM listings WHERE id = 910002")
).scalar()
assert repointed == 900001 assert repointed == 900001
# price_dynamics deduped — exactly one row for the keeper at that 6-col key (no dup), # price_dynamics deduped — exactly one row for the keeper at that 6-col key (no dup),
# and it is the KEEPER's own row (price_per_sqm=100000), not the loser's (999999): # and it is the KEEPER's own row (price_per_sqm=100000), not the loser's (999999):
@ -580,10 +579,7 @@ def test_real_merge_repoints_dedups_deletes_and_is_idempotent() -> None:
_t("DELETE FROM house_sources WHERE ext_id IN ('SRC-KEEP','SRC-LOSE','EXT-KEEP')") _t("DELETE FROM house_sources WHERE ext_id IN ('SRC-KEEP','SRC-LOSE','EXT-KEEP')")
) )
db.execute( db.execute(
_t( _t("DELETE FROM house_address_aliases WHERE normalized_address = 'тестдом 1772, 1'")
"DELETE FROM house_address_aliases "
"WHERE normalized_address = 'тестдом 1772, 1'"
)
) )
db.execute(_t("DELETE FROM houses WHERE id IN (900001,900002)")) db.execute(_t("DELETE FROM houses WHERE id IN (900001,900002)"))
db.commit() db.commit()

View file

@ -54,9 +54,10 @@ def test_trigger_mirrors_backfill_sibling_pattern() -> None:
assert "backfill_house_imv" in _TRIGGER_SRC assert "backfill_house_imv" in _TRIGGER_SRC
# async service → bare await (NOT run_in_executor, unlike the sync DB-only triggers). # async service → bare await (NOT run_in_executor, unlike the sync DB-only triggers).
assert "run_in_executor" not in _TRIGGER_SRC assert "run_in_executor" not in _TRIGGER_SRC
# RUF006: keep a reference + done-callback so the task is not GC'd mid-flight. # #1182 P2: detached spawn через _spawn_tracked — strong-ref в registry (RUF006 +
assert "task = asyncio.create_task(_run())" in _TRIGGER_SRC # graceful SIGTERM-drain), done-callback ретривит exception. Заменил прежний
assert "task.add_done_callback(" in _TRIGGER_SRC # `task = asyncio.create_task(_run()); task.add_done_callback(...)`.
assert "_spawn_tracked(_run())" in _TRIGGER_SRC
assert "run_db.close()" in _TRIGGER_SRC assert "run_db.close()" in _TRIGGER_SRC

View file

@ -2,7 +2,7 @@
import asyncio import asyncio
import os import os
from unittest.mock import MagicMock from unittest.mock import AsyncMock, MagicMock
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
@ -10,9 +10,22 @@ from datetime import UTC, datetime
import pytest import pytest
import app.services.scheduler as sched
from app.core import shutdown as sd
from app.services.scheduler import ZOMBIE_THRESHOLD_HOURS, compute_next_run_at from app.services.scheduler import ZOMBIE_THRESHOLD_HOURS, compute_next_run_at
@pytest.fixture(autouse=True)
def _reset_drain_state() -> None:
"""#1182 P2: shutdown-флаг и registry детей — module-global; чистим вокруг каждого
теста (изоляция от detached-задач, что триггеры спавнят в _inflight_tasks)."""
sd.reset_shutdown()
sched._inflight_tasks.clear()
yield
sched._inflight_tasks.clear()
sd.reset_shutdown()
def test_compute_next_run_in_normal_window() -> None: def test_compute_next_run_in_normal_window() -> None:
now = datetime(2026, 5, 23, 12, 0, tzinfo=UTC) now = datetime(2026, 5, 23, 12, 0, tzinfo=UTC)
next_at = compute_next_run_at(2, 5, now=now) next_at = compute_next_run_at(2, 5, now=now)
@ -730,3 +743,83 @@ async def test_scheduler_dispatch_routes_avito_full_load_exhaustive(
await _sched.trigger_avito_full_load_exhaustive_run(db, sch) await _sched.trigger_avito_full_load_exhaustive_run(db, sch)
assert "avito_full_load_exhaustive" in triggered assert "avito_full_load_exhaustive" in triggered
# ---------------------------------------------------------------------------
# #1182 Phase 2 — graceful SIGTERM-drain детей (detached run-задач)
# ---------------------------------------------------------------------------
@pytest.mark.asyncio
async def test_drain_inflight_awaits_detached_cooperative_child() -> None:
"""Гвоздь Phase 2: _drain_inflight ДОЖИДАЕТСЯ detached-ребёнка (как scrape-run,
отвязанный через _spawn_tracked) до его clean-exit, а НЕ хард-кансельит его.
Регрессия gap'а: coordinator раньше ждал только себя → asyncio.run teardown
хард-кансельил живых детей mid-await (#1182 failure mode).
"""
clean_exit = asyncio.Event() # == mark_done-эквивалент ребёнка
async def _child() -> None:
# Кооперативный ребёнок: крутится, пока не запрошен shutdown, потом выходит чисто.
while not sd.shutdown_requested():
await asyncio.sleep(0.01)
clean_exit.set()
child = sched._spawn_tracked(_child()) # detached, попадает в _inflight_tasks
await asyncio.sleep(0.02)
assert not child.done(), "ребёнок ещё бежит (shutdown не запрошен)"
sd.request_shutdown()
await asyncio.wait_for(sched._drain_inflight(), timeout=2.0)
assert clean_exit.is_set(), "drain дождался clean-exit ребёнка"
assert child.done()
assert not child.cancelled(), "ребёнок вышел сам, НЕ хард-cancel"
@pytest.mark.asyncio
async def test_drain_inflight_times_out_on_noncooperative_child(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""Некооперативный ребёнок (игнорит shutdown) упирается в _CHILD_DRAIN_TIMEOUT_S и
остаётся бежать отдаётся внешнему hard-cancel fallback (scheduler_main/docker)."""
monkeypatch.setattr(sched, "_CHILD_DRAIN_TIMEOUT_S", 0.05)
async def _stuck() -> None:
await asyncio.Event().wait() # никогда не проверяет shutdown_requested()
child = sched._spawn_tracked(_stuck())
sd.request_shutdown()
await asyncio.wait_for(sched._drain_inflight(), timeout=2.0)
assert not child.done(), "drain НЕ повесил процесс — оставил ребёнка под hard-cancel"
# cleanup
child.cancel()
with pytest.raises(asyncio.CancelledError):
await child
@pytest.mark.asyncio
async def test_drain_inflight_noop_when_no_children() -> None:
"""Без in-flight детей _drain_inflight() — мгновенный no-op (без сна/ожидания)."""
assert not sched._inflight_tasks
await asyncio.wait_for(sched._drain_inflight(), timeout=1.0)
@pytest.mark.asyncio
async def test_scheduler_loop_calls_drain_on_shutdown(monkeypatch: pytest.MonkeyPatch) -> None:
"""scheduler_loop при выходе по SIGTERM-drain'у вызывает _drain_inflight (проводка
coordinator дренаж детей). Initial sleep(30) и tick-sleep занулены, чтобы тест
был детерминированным и быстрым."""
drain_spy = AsyncMock()
monkeypatch.setattr(sched, "_drain_inflight", drain_spy)
monkeypatch.setattr(sched.asyncio, "sleep", AsyncMock()) # без реальных 30s/60s
sd.request_shutdown() # первый же top-of-tick check → break
await asyncio.wait_for(sched.scheduler_loop(), timeout=2.0)
drain_spy.assert_awaited_once()

View file

@ -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