fix(tradein/scraper): drain detached run-children on SIGTERM, not just coordinator (#1182 Phase 2)
Pre-push review: _await_scheduler ждал только scheduler_loop COORDINATOR, но вся scrape-работа крутится в detached asyncio.create_task детях (каждый trigger_* делал `task = create_task(_run())` без join). На SIGTERM coordinator выходил из while True и завершался → asyncio.run() teardown хард-кансельил ещё бегущих детей mid-await = ровно #1182 failure mode. Кооперативный checkpoint спасал ребёнка лишь когда его residual-время случайно перекрывало drain — вероятностно, не гарантированно. Fix — coordinator теперь дренажит детей перед выходом: - _spawn_tracked(coro): централизованный detached-spawn, кладёт задачу в module-level registry _inflight_tasks (strong-ref = RUF006 keep-alive) + done-callback ретривит exception и убирает из set'а. Заменил 26 одинаковых `task = create_task(_run()); task.add_done_callback(...)` сайтов. - _drain_inflight(): один asyncio.wait по живым детям с бюджетом _CHILD_DRAIN_TIMEOUT_S=80s (< scheduler_main 100s < docker grace 120s). Кооперативные дети (avito_detail_backfill, rosreestr-executor) дочекивают карточку/батч + mark_done и резолвятся; некооперативные упираются в timeout и падают на внешний hard-cancel. - scheduler_loop по выходу из tick-loop (только по SIGTERM) зовёт await _drain_inflight(). NB: raw asyncio.all_tasks()-minus-self здесь НЕЛЬЗЯ — в нашей топологии он захватывает _run parent-task (блокирован на wait_for(coordinator)) и shutdown_waiter → циклическое ожидание coordinator↔_run, всегда упирающееся в timeout. Точный registry это исключает. Tests: tests/test_scheduler.py — detached cooperative child drained-not-cancelled, non-cooperative child timeout→left for hard-cancel, no-op без детей, scheduler_loop→drain wiring. Обновил 3 source-inspection теста под новый _spawn_tracked паттерн.
This commit is contained in:
parent
6f76b565e5
commit
566b2f9617
5 changed files with 197 additions and 69 deletions
|
|
@ -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
|
||||||
|
|
||||||
|
|
@ -158,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,
|
||||||
|
|
@ -369,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
|
||||||
|
|
||||||
|
|
@ -407,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
|
||||||
|
|
||||||
|
|
@ -447,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
|
||||||
|
|
||||||
|
|
@ -507,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
|
||||||
|
|
||||||
|
|
@ -543,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
|
||||||
|
|
||||||
|
|
@ -579,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
|
||||||
|
|
||||||
|
|
@ -625,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
|
||||||
|
|
||||||
|
|
@ -669,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
|
||||||
|
|
||||||
|
|
@ -706,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
|
||||||
|
|
||||||
|
|
@ -810,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
|
||||||
|
|
||||||
|
|
@ -849,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
|
||||||
|
|
||||||
|
|
@ -874,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
|
||||||
|
|
||||||
|
|
@ -909,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
|
||||||
|
|
||||||
|
|
@ -954,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
|
||||||
|
|
||||||
|
|
@ -989,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
|
||||||
|
|
||||||
|
|
@ -1042,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
|
||||||
|
|
||||||
|
|
@ -1078,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
|
||||||
|
|
||||||
|
|
@ -1116,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
|
||||||
|
|
||||||
|
|
@ -1159,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
|
||||||
|
|
||||||
|
|
@ -1208,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
|
||||||
|
|
||||||
|
|
@ -1243,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
|
||||||
|
|
||||||
|
|
@ -1531,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
|
||||||
|
|
||||||
|
|
@ -1568,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
|
||||||
|
|
||||||
|
|
@ -1610,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
|
||||||
|
|
||||||
|
|
@ -1683,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
|
||||||
|
|
||||||
|
|
@ -1730,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
|
||||||
|
|
||||||
|
|
@ -1846,3 +1875,9 @@ async def scheduler_loop() -> None:
|
||||||
logger.info("scheduler: SIGTERM-drain — exiting tick loop after dispatch")
|
logger.info("scheduler: SIGTERM-drain — exiting tick loop after dispatch")
|
||||||
break
|
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)")
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
|
|
||||||
|
|
@ -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()
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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()
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue