fix(tradein/scheduler): SIGTERM-drain снимает с 'running' in-flight app-task'и (#3391)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 10s
CI / changes (pull_request) Successful in 12s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m55s

Прод 07.09 02:36 UTC, первый настоящий drain после #3363: hard-cancel из
scheduler_main оборвал дрейн, и два бэкфилла (cian_detail_backfill 6167,
cian_history_backfill 6173) остались в scrape_runs со статусом 'running' —
boot-reap следующего контейнера сделал их 'zombie' (boot_reaped=true), метки
interrupted не было. interrupted=1 при дрейне писали только kit-пайплайны и
DKP-импорт: у задач, чьё тело живёт в app, ставить её было некому.

Метка ставится в единственной точке, через которую проходит любая detached
run-задача — SchedulerContext.drain_inflight: и по истечении
_CHILD_DRAIN_TIMEOUT_S, и в обработчике CancelledError (тот самый прод-путь).
run_id берётся из нового реестра {task: run_id}, который заполняет _dispatch
сразу после claim'а; claim-логика не тронута. Статус строки перечитывается
перед записью, поэтому успевший финализироваться сам прогон не
перезаписывается, а counters читаются из строки и дописываются — app-копия
mark_done их ЗАМЕНЯЕТ (#3390), голая {"interrupted": 1} стёрла бы чекпоинт.

scheduler_main: _await_scheduler возвращает признак hard-cancel'а, и строка
«scheduler drained cleanly (SIGTERM)» больше не печатается сразу за WARNING'ом
о превышении grace — на проде эти две строки стояли подряд и противоречили
друг другу.

Запас времени на запись: docker stop_grace_period 120s − _DRAIN_TIMEOUT_S 100s
= 20 с после hard-cancel'а, запись синхронная (несколько statement'ов).
This commit is contained in:
bot-backend 2026-09-06 07:56:35 +05:00
parent 3fe2d310b2
commit 30e3bacc5e
3 changed files with 364 additions and 9 deletions

View file

@ -101,15 +101,23 @@ def _should_run() -> bool:
return settings.scheduler_enable
async def _await_scheduler(task: asyncio.Task[None]) -> None:
async def _await_scheduler(task: asyncio.Task[None]) -> bool:
"""Дождаться завершения scheduler-задачи с кооперативным SIGTERM-drain'ом.
Возвращает True, если задачу пришлось хард-кансельнуть (grace истёк), False если
она вышла сама. Вызывающий обязан различать эти исходы в логах (#3391): строка
«drained cleanly» после hard-cancel'а — ложь, а именно она печаталась на проде
07.09 сразу за WARNING'ом о превышении grace.
- Нет shutdown: scheduler_loop бесконечен задача никогда не завершается, ждём её
как есть (процесс просто работает).
- shutdown запрошен, задача ещё бежит: bounded `wait_for(_DRAIN_TIMEOUT_S)`
кооперативная задача докоммитит текущий unit на ближайшем checkpoint'е и выйдет
сама. Если превысила grace hard-cancel + suppress CancelledError, чтобы
некооперативная задача не подвесила процесс за пределами docker stop_grace_period.
Пометку `interrupted` своим in-flight прогонам kit-scheduler успевает поставить в
обработчике CancelledError (SchedulerContext.mark_inflight_interrupted): между
hard-cancel'ом и SIGKILL'ом остаётся 20 с docker-grace (120s 100s).
"""
# Гонка «задача завершилась сама» против «пришёл SIGTERM»: отсчёт safety-net'а
# должен стартовать от МОМЕНТА запроса drain'а, а не от старта процесса.
@ -127,7 +135,7 @@ async def _await_scheduler(task: asyncio.Task[None]) -> None:
# упал (сохраняем прежнюю loud-crash семантику `await task`), а не глушит его.
task.result()
logger.info("scheduler_main: scheduler task exited cleanly")
return
return False
logger.info(
"scheduler_main: SIGTERM-drain — waiting up to %.0fs for in-flight unit to commit",
@ -136,6 +144,7 @@ async def _await_scheduler(task: asyncio.Task[None]) -> None:
try:
await asyncio.wait_for(task, timeout=_DRAIN_TIMEOUT_S)
logger.info("scheduler_main: scheduler drained and exited cleanly")
return False
except TimeoutError:
logger.warning(
"scheduler_main: drain exceeded %.0fs grace — hard-cancelling scheduler task",
@ -144,6 +153,7 @@ async def _await_scheduler(task: asyncio.Task[None]) -> None:
task.cancel()
with suppress(asyncio.CancelledError):
await task
return True
async def _run_kit_scheduler() -> None:
@ -224,8 +234,12 @@ async def _run() -> None:
# Windows dev: signal handlers через loop не поддерживаются
logger.warning("scheduler_main: loop.add_signal_handler not supported (Windows dev)")
await _await_scheduler(task)
hard_cancelled = await _await_scheduler(task)
if hard_cancelled:
# Дрейн НЕ был чистым: задача не вышла сама, её сняли. WARNING об этом уже
# напечатан в _await_scheduler — второй строкой её не «переобъявляем».
return
if shutdown_requested():
logger.info("scheduler_main: scheduler drained cleanly (SIGTERM)")
else:

View file

@ -0,0 +1,266 @@
"""SIGTERM-drain снимает с 'running' прогоны app-task'ов, не доживших до финализации (#3391).
Прод 07.09 02:36 UTC (первый настоящий SIGTERM-drain после #3363): деплой пересоздал
`tradein-scraper`, kit-scheduler напечатал «draining 2 in-flight run task(s)», через 100 с
`scheduler_main` хард-кансельнул задачу и оба бэкфилла (cian_detail_backfill 6167,
cian_history_backfill 6173) остались в scrape_runs со статусом 'running' до boot-reap'а
следующего контейнера, где стали 'zombie' с `boot_reaped=true`. `interrupted=1` при дрейне
пишут только kit-пайплайны (#3363) и DKP-импорт: у задач, чьё тело живёт в `app`, метку
ставить некому.
Проверяется НЕ «функция вызвалась», а значение в строке прогона:
1. grace истёк строка застрявшей задачи 'done' + counters.interrupted=1, прежние
счётчики (в т.ч. чекпоинт) на месте;
2. соседняя задача, успевшая финализироваться сама, НЕ перезаписана (без interrupted);
3. hard-cancel из scheduler_main (CancelledError прямо в `asyncio.wait` дрейна) даёт тот
же результат это ровно прод-путь 07.09, ветка таймаута его не покрывает;
4. `scheduler_main` после hard-cancel'а больше не печатает «drained cleanly».
Фейковый `mark_done` counters ЗАМЕНЯЕТ, а не мержит так ведёт себя боевая app-копия
(`app.services.scrape_runs`, #3390), которая и инжектируется в прод-контекст. Голый
`{"interrupted": 1}` на этом фейке стёр бы чекпоинт, и (1) покраснел бы.
"""
from __future__ import annotations
import asyncio
import logging
import os
import types
from contextlib import suppress
from typing import Any
from unittest.mock import MagicMock
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
import pytest
from scraper_kit.orchestration import scheduler as kit_sched
from scraper_kit.orchestration.scheduler import Handler, SchedulerContext, _dispatch
import app.scheduler_main as sm
from app.core import shutdown as sd
class _Table:
"""Мини-`scrape_runs`: id → {status, counters}."""
def __init__(self) -> None:
self.rows: dict[int, dict[str, Any]] = {}
self._next_id = 6167
def claim(self, source: str) -> int:
run_id = self._next_id
self._next_id += 6
self.rows[run_id] = {"status": "running", "counters": {}, "source": source}
return run_id
class _FakeResult:
def __init__(self, *, fetchone: Any = None, scalar: Any = None) -> None:
self._fetchone = fetchone
self._scalar = scalar
def fetchone(self) -> Any:
return self._fetchone
def scalar(self) -> Any:
return self._scalar
class _FakeDb:
"""Session поверх `_Table`: claim-путь `_claim_run` + SELECT статуса дрейна."""
def __init__(self, table: _Table, *, running_states: list[bool] | None = None) -> None:
self.table = table
self._running = list(running_states or [])
self.closed = False
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
sql = str(stmt)
if "SELECT status, counters FROM scrape_runs" in sql:
row = self.table.rows.get(int((params or {})["id"]))
if row is None:
return _FakeResult(fetchone=None)
return _FakeResult(
fetchone=types.SimpleNamespace(status=row["status"], counters=row["counters"])
)
if "SELECT 1 FROM scrape_runs" in sql:
return _FakeResult(fetchone=(1,) if self._running.pop(0) else None)
if "pg_try_advisory_xact_lock" in sql:
return _FakeResult(scalar=True)
return _FakeResult()
def commit(self) -> None:
pass
def rollback(self) -> None:
pass
def close(self) -> None:
self.closed = True
class _FakeRuns:
"""ctx.runs: create_run + mark_done с боевым гейтом `WHERE status = 'running'`."""
def __init__(self, table: _Table) -> None:
self.table = table
def create_run(self, db: Any, *, source: str, params: dict[str, Any]) -> int:
return self.table.claim(source)
def mark_done(self, db: Any, run_id: int, counters: dict[str, Any]) -> None:
row = self.table.rows.get(run_id)
if row is None or row["status"] != "running":
return # боевой UPDATE ... WHERE status = 'running' — no-op
row["status"] = "done"
row["counters"] = dict(counters) # app-копия ЗАМЕНЯЕТ counters (#3390)
def _make_sched(source: str) -> dict[str, Any]:
return {
"id": 1,
"source": source,
"enabled": True,
"window_start_hour": 2,
"window_end_hour": 5,
"default_params": {},
}
def _ctx(table: _Table) -> SchedulerContext:
return SchedulerContext(
config=MagicMock(),
matcher=MagicMock(),
enrichment=MagicMock(),
session_factory=lambda: _FakeDb(table),
runs=_FakeRuns(table),
)
async def _dispatch_stuck(ctx: SchedulerContext, table: _Table, started: asyncio.Event) -> int:
"""app-task, который на дрейн не реагирует (долгий await вместо checkpoint'а)."""
async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
started.set()
await asyncio.sleep(10)
run_id = await _dispatch(
Handler(_job, "cian_detail_backfill"),
_FakeDb(table, running_states=[False, False]),
_make_sched("cian_detail_backfill"),
ctx,
)
assert run_id is not None
await asyncio.wait_for(started.wait(), timeout=2.0)
# Пульс до дрейна: собранное этим прогоном и его чекпоинт.
table.rows[run_id]["counters"] = {"lots_fetched": 120, "done_buckets": ["1:0:5"]}
return run_id
async def _cancel_inflight(ctx: SchedulerContext) -> None:
for task in list(ctx._inflight_tasks):
task.cancel()
with suppress(asyncio.CancelledError):
await task
# ── 1-2. grace истёк: застрявший помечен, самофинализировавшийся не тронут ───────
async def test_drain_grace_expiry_marks_stuck_run_interrupted(
monkeypatch: pytest.MonkeyPatch,
) -> None:
monkeypatch.setattr(kit_sched, "_CHILD_DRAIN_TIMEOUT_S", 0.05)
table = _Table()
ctx = _ctx(table)
stuck_id = await _dispatch_stuck(ctx, table, asyncio.Event())
finished = asyncio.Event()
async def _quick(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
c.runs.mark_done(db, run_id, {"lots_fetched": 5})
finished.set()
quick_id = await _dispatch(
Handler(_quick, "cian_history_backfill"),
_FakeDb(table, running_states=[False, False]),
_make_sched("cian_history_backfill"),
ctx,
)
await asyncio.wait_for(finished.wait(), timeout=2.0)
await ctx.drain_inflight()
assert table.rows[stuck_id]["status"] == "done"
assert table.rows[stuck_id]["counters"]["interrupted"] == 1
# Чекпоинт и собранное пережили пометку — иначе резюм подхватывать нечего.
assert table.rows[stuck_id]["counters"]["done_buckets"] == ["1:0:5"]
assert table.rows[stuck_id]["counters"]["lots_fetched"] == 120
# Финализировавшийся сам — не перезаписан: ни статуса, ни лишней метки.
assert table.rows[quick_id]["status"] == "done"
assert table.rows[quick_id]["counters"] == {"lots_fetched": 5}
await _cancel_inflight(ctx)
# ── 3. hard-cancel из scheduler_main — прод-путь 07.09 ──────────────────────────
async def test_drain_hard_cancel_marks_stuck_run_interrupted(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""CancelledError прилетает прямо в `asyncio.wait` дрейна (grace scheduler_main истёк)."""
monkeypatch.setattr(kit_sched, "_CHILD_DRAIN_TIMEOUT_S", 30.0)
table = _Table()
ctx = _ctx(table)
stuck_id = await _dispatch_stuck(ctx, table, asyncio.Event())
drain = asyncio.create_task(ctx.drain_inflight())
await asyncio.sleep(0) # дать дрейну дойти до asyncio.wait
await asyncio.sleep(0)
drain.cancel()
with suppress(asyncio.CancelledError):
await drain
assert table.rows[stuck_id]["status"] == "done"
assert table.rows[stuck_id]["counters"]["interrupted"] == 1
await _cancel_inflight(ctx)
# ── 4. scheduler_main: после hard-cancel'а «drained cleanly» не печатается ───────
@pytest.fixture()
def _reset_shutdown() -> Any:
sd.reset_shutdown()
yield
sd.reset_shutdown()
async def test_hard_cancel_does_not_log_drained_cleanly(
monkeypatch: pytest.MonkeyPatch,
caplog: pytest.LogCaptureFixture,
_reset_shutdown: Any,
) -> None:
"""Прод 07.09: WARNING о hard-cancel'е и INFO «drained cleanly» стояли подряд."""
caplog.set_level(logging.INFO)
monkeypatch.setattr(sm, "_DRAIN_TIMEOUT_S", 0.05)
started = asyncio.Event()
async def _stub_loop() -> None:
started.set()
await asyncio.sleep(30) # некооперативная задача: shutdown не смотрит
monkeypatch.setattr(sm, "_run_kit_scheduler", _stub_loop)
task = asyncio.create_task(sm._run())
await asyncio.wait_for(started.wait(), timeout=2.0)
sd.request_shutdown()
await asyncio.wait_for(task, timeout=5.0)
assert "hard-cancelling scheduler task" in caplog.text
assert "drained cleanly" not in caplog.text
assert "drained and exited cleanly" not in caplog.text

View file

@ -32,7 +32,7 @@ from __future__ import annotations
import asyncio
import logging
import random
from collections.abc import Callable, Coroutine, Mapping
from collections.abc import Callable, Coroutine, Iterable, Mapping
from dataclasses import dataclass, field
from datetime import UTC, datetime, time, timedelta
from typing import TYPE_CHECKING, Any
@ -348,42 +348,117 @@ class SchedulerContext:
runs: Any = _kit_runs
# #1182 P2: strong-ref'ы detached run-задач для graceful SIGTERM-drain'а.
_inflight_tasks: set[asyncio.Task[None]] = field(default_factory=set)
# #3391: task → run_id заклеймленного прогона. Без него дрейн знает, что задача
# не дожила до финализации, но не знает, КАКУЮ строку scrape_runs снимать с
# 'running' (claim-логика run_id никуда не публиковала).
_inflight_run_ids: dict[asyncio.Task[None], int] = field(default_factory=dict)
def spawn_tracked(self, coro: Coroutine[Any, Any, None]) -> asyncio.Task[None]:
def spawn_tracked(
self, coro: Coroutine[Any, Any, None], *, run_id: int | None = None
) -> asyncio.Task[None]:
"""create_task + регистрация в _inflight_tasks для graceful-drain'а (#1182 P2).
strong-ref в set'е держит задачу до завершения (RUF006) и даёт scheduler_loop'у
дождаться её на SIGTERM-drain'е; done-callback ретривит exception и убирает задачу.
`run_id` (#3391) — строка scrape_runs этой задачи; None у задач без прогона.
"""
task = asyncio.create_task(coro)
self._inflight_tasks.add(task)
if run_id is not None:
self._inflight_run_ids[task] = run_id
def _on_done(t: asyncio.Task[None]) -> None:
self._inflight_tasks.discard(t)
self._inflight_run_ids.pop(t, None)
if not t.cancelled():
t.exception()
task.add_done_callback(_on_done)
return task
def mark_inflight_interrupted(self, tasks: Iterable[asyncio.Task[None]]) -> int:
"""Снять с 'running' прогоны задач, не доживших до собственной финализации (#3391).
Прод 07.09 (первый настоящий SIGTERM-drain после #3363): два app-task бэкфилла
(cian_detail_backfill 6167, cian_history_backfill 6173) не уложились в grace,
scheduler_main их хард-кансельнул и строки остались 'running' до boot-reap'а
следующего контейнера, где стали 'zombie' с `boot_reaped=true`. `interrupted=1`
при дрейне пишут только kit-пайплайны (#3363) и DKP-импорт: у задач, чьё тело
живёт в app и о дрейне не знает, метку ставить некому. Ставим её ЗДЕСЬ в
единственной точке, через которую проходит любая detached run-задача.
Синхронно (sync-сессия, ни одного await): вызывается в том числе из except
CancelledError, где лишний await мог бы не дожить до конца.
Статус строки перечитывается перед записью: задача, успевшая финализироваться
сама (done/failed/banned), не перезаписывается. Гонку добивает сам `mark_done`
(`WHERE status = 'running'`) но тогда счётчики уже прочитаны, и их merge был бы
холостым. counters берём из строки и дописываем `interrupted`, а не отдаём
`{"interrupted": 1}` голым: app-копия `mark_done` counters ЗАМЕНЯЕТ, а не мержит
(#3390) — голый словарь стёр бы всю бухгалтерию прогона, включая чекпоинт.
"""
run_ids = [rid for t in tasks if (rid := self._inflight_run_ids.get(t)) is not None]
if not run_ids:
return 0
marked = 0
db = self.session_factory()
try:
for run_id in run_ids:
try:
row = db.execute(
text("SELECT status, counters FROM scrape_runs WHERE id = :id"),
{"id": run_id},
).fetchone()
if row is None or row.status != "running":
continue
counters = dict(row.counters or {})
counters["interrupted"] = 1
self.runs.mark_done(db, run_id, counters)
marked += 1
except Exception:
# Один непроходимый прогон не должен утащить остальные: дрейну
# осталось секунды до SIGKILL, а каждая незакрытая строка — это
# ещё один 'zombie' у следующего контейнера.
logger.exception("scheduler: drain — не удалось пометить run_id=%d", run_id)
finally:
db.close()
if marked:
logger.warning(
"scheduler: drain — %d прогон(ов) сняты с 'running' как interrupted: %s",
marked,
run_ids,
)
return marked
async def drain_inflight(self) -> None:
"""Дождаться завершения detached run-задач перед teardown'ом (#1182 P2).
Кооперативные дети дочекивают текущий unit + mark_done и резолвятся;
некооперативные упираются в _CHILD_DRAIN_TIMEOUT_S и остаются на внешний
hard-cancel fallback. Не busy-spin один asyncio.wait.
некооперативные упираются в _CHILD_DRAIN_TIMEOUT_S либо во внешний hard-cancel
из scheduler_main и в обоих случаях их строки помечаются `interrupted`
(#3391), иначе они доживают в 'running' до boot-reap'а. Не busy-spin — один
asyncio.wait.
"""
pending = [t for t in self._inflight_tasks if not t.done()]
if not pending:
return
logger.info("scheduler: draining %d in-flight run task(s) on shutdown", len(pending))
_done, still = await asyncio.wait(pending, timeout=_CHILD_DRAIN_TIMEOUT_S)
try:
_done, still = await asyncio.wait(pending, timeout=_CHILD_DRAIN_TIMEOUT_S)
except asyncio.CancelledError:
# Hard-cancel из scheduler_main (grace 100s истёк раньше нашего дрейна —
# такт tick-сна + _CHILD_DRAIN_TIMEOUT_S могут превысить его). Пишем метку
# синхронно и пробрасываем отмену дальше.
self.mark_inflight_interrupted([t for t in pending if not t.done()])
raise
if still:
logger.warning(
"scheduler: %d task(s) did not drain in %.0fs — leaving for hard-cancel",
len(still),
_CHILD_DRAIN_TIMEOUT_S,
)
self.mark_inflight_interrupted(still)
else:
logger.info("scheduler: all in-flight run task(s) drained cleanly")
@ -897,7 +972,7 @@ async def _dispatch(
finally:
run_db.close()
ctx.spawn_tracked(_run())
ctx.spawn_tracked(_run(), run_id=run_id)
logger.info("scheduler: triggered %s run_id=%d", handler.log_name, run_id)
return run_id