fix(scraper-kit/scheduler): SIGTERM-drain помечает in-flight прогоны interrupted=1 при hard-cancel; честный лог drain'а (#3391) #3392

Merged
bot-backend merged 2 commits from fix/drain-marks-inflight-app-tasks into main 2026-09-06 04:14:07 +00:00
5 changed files with 716 additions and 19 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

@ -599,9 +599,23 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
COALESCE: если ключа нет в counters старое значение колонки сохраняется.
#3391: гейт по статусу — тот же, что у mark_done/mark_failed. Пульс по УЖЕ
финализированной строке не просто холостой: эта копия counters ЗАМЕНЯЕТ (#3390),
поэтому задача, помеченная дрейном как `interrupted`, но ещё живая (ветка таймаута
drain_inflight отдаёт её на внешний hard-cancel несколько итераций спустя),
следующим же пульсом стирала метку, и оборванный прогон снова читался как полный
проход.
'cancelled' в гейте НАМЕРЕННО, это не «running и всё финализированное»: user-cancel
финализирует строку, но задача останавливается лишь на ближайшей границе якоря, и её
последний пульс ЕДИНСТВЕННЫЙ писатель чекпоинта в этот момент (pipeline.py:1308 /
2368 / 2986 / 4488; mark_done там уже no-op по своему гейту, а 'cancelled' входит в
_RESUME_STATUSES). Сузить гейт до 'running' молча потерять точку возобновления
у каждой отмены.
"""
total_seen, new_count = _column_counts(counters)
db.execute(
row = db.execute(
text(
"""
UPDATE scrape_runs
@ -609,7 +623,8 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None
counters = CAST(:counters AS jsonb),
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
new_count = COALESCE(CAST(:new_count AS int), new_count)
WHERE id = :run_id
WHERE id = :run_id AND status IN ('running', 'cancelled')
RETURNING id
"""
),
{
@ -618,7 +633,12 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None
"total_seen": total_seen,
"new_count": new_count,
},
)
).first()
if row is None:
# Не рутина: задача пережила собственную финализацию и продолжает слать прогресс.
logger.warning(
"update_heartbeat no-op: run_id=%d not in 'running'/'cancelled' state", run_id
)
db.commit()

View file

@ -0,0 +1,529 @@
"""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»;
5. пульс НЕ пишет по финализированной строке (обе копии `update_heartbeat`) иначе
помеченная, но ещё живая задача стирает метку следующим же ударом;
6. отмена, прилетевшая в тело тика (не в дрейн), помечает in-flight так же;
7. отказ SQL на одном run_id не уносит остальные, а WARNING перечисляет ровно
помеченных.
Фейковый `mark_done` counters ЗАМЕНЯЕТ, а не мержит так ведёт себя боевая app-копия
(`app.services.scrape_runs`, #3390), которая и инжектируется в прод-контекст. Голый
`{"interrupted": 1}` на этом фейке стёр бы чекпоинт, и (1) покраснел бы.
"""
from __future__ import annotations
import asyncio
import json
import logging
import os
import re
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 runs as kit_runs
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
from app.services import scrape_runs as app_runs
# Обе копии runs-модуля: у app counters ЗАМЕНЯЮТСЯ, у kit мержатся (#3390) — гейт по
# статусу нужен обеим, и проверяется на обеих одним и тем же телом теста.
_RUNS_MODULES = {"kit": kit_runs, "app": app_runs}
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, first: Any = None) -> None:
self._fetchone = fetchone
self._scalar = scalar
self._first = first
def fetchone(self) -> Any:
return self._fetchone
def scalar(self) -> Any:
return self._scalar
def first(self) -> Any:
return self._first
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
# ── 5. Пульс не пишет по финализированной строке (обе копии update_heartbeat) ────
class _RunRowDb:
"""Мини-Postgres на одну строку scrape_runs: UPDATE применяется, только если строка
проходит WHERE из ТЕКСТА самого statement'а.
Допустимые статусы вычитываются из SQL (`status = 'x'` / `status IN ('x', 'y')`), а не
зашиты ожиданием теста: на коде без гейта (WHERE только по id) апдейт проходит по
строке ЛЮБОГО статуса и тест краснеет по значению, а не по отсутствию подстроки.
Мерж jsonb (`||`) против замены (`CAST(:counters AS jsonb)`) тоже по тексту: у
kit- и app-копии он разный, а гейт нужен обеим.
SELECT'ы alert-хуков (`_alert_on_run_id`) не моделируются — они best-effort и
возвращают пусто.
"""
def __init__(self, *, status: str = "running", counters: dict[str, Any] | None = None) -> None:
self.row: dict[str, Any] = {
"status": status,
"counters": dict(counters or {}),
"heartbeat_at": None,
}
self.clock = 0.0
@staticmethod
def _allowed_statuses(sql: str) -> set[str] | None:
"""Статусы из WHERE. None — гейта нет, UPDATE бьёт по строке любого статуса."""
where = sql.rsplit("WHERE", 1)[-1]
in_list = re.search(r"status\s+IN\s*\(([^)]*)\)", where)
if in_list is not None:
return {s.strip().strip("'") for s in in_list.group(1).split(",")}
eq = re.search(r"status\s*=\s*'(\w+)'", where)
return {eq.group(1)} if eq is not None else None
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
sql = " ".join(str(stmt).split())
if not sql.startswith("UPDATE scrape_runs"):
return _FakeResult()
allowed = self._allowed_statuses(sql)
if allowed is not None and self.row["status"] not in allowed:
return _FakeResult(first=None) # WHERE не пропустил — 0 строк, RETURNING пуст
payload = json.loads((params or {})["counters"])
self.row["counters"] = {**self.row["counters"], **payload} if "||" in sql else dict(payload)
self.clock += 1.0
self.row["heartbeat_at"] = self.clock
new_status = re.search(r"SET status = '(\w+)'", sql)
if new_status is not None:
self.row["status"] = new_status.group(1)
return _FakeResult(first=(1,))
def commit(self) -> None:
pass
def rollback(self) -> None:
pass
@pytest.mark.parametrize("name", list(_RUNS_MODULES))
def test_heartbeat_does_not_erase_drain_mark_of_finalized_run(name: str) -> None:
"""Прогон помечен дрейном ('done' + interrupted=1), но задача ещё жива и бьёт пульс.
Ровно эта последовательность достижима на проде: ветка таймаута `drain_inflight`
оставляет задачу внешнему hard-cancel'у, и до него она успевает несколько итераций
(`app/services/scheduler.py:141` пульс на каждый батч). Без гейта по статусу
app-копия ЗАМЕНЯЛА counters и стирала метку: оборванный прогон снова читался как
полный проход, а резюм его не подхватывал.
"""
mod = _RUNS_MODULES[name]
db = _RunRowDb(counters={"lots_fetched": 120, "done_buckets": ["1:0:5"]})
mod.mark_done(db, 6167, {"lots_fetched": 120, "done_buckets": ["1:0:5"], "interrupted": 1})
assert db.row["status"] == "done"
finalized_at = db.row["heartbeat_at"]
mod.update_heartbeat(db, 6167, {"listings_processed": 900})
assert db.row["counters"]["interrupted"] == 1, "пульс стёр метку дрейна"
assert "listings_processed" not in db.row["counters"], "пульс дописал прогресс в финал"
assert db.row["heartbeat_at"] == finalized_at, "пульс сдвинул heartbeat финализированной"
@pytest.mark.parametrize("name", list(_RUNS_MODULES))
@pytest.mark.parametrize("status", ["running", "cancelled"])
def test_heartbeat_still_writes_live_and_cancelled_rows(name: str, status: str) -> None:
"""Гейт не должен глушить живой пульс — и в частности пульс по ОТМЕНЁННОЙ строке.
'cancelled' финализирует строку, но задача останавливается лишь на ближайшей границе
якоря, и её последний пульс единственный, кто персистит чекпоинт в этот момент
(`pipeline.py:1308/2368/2986/4488`: mark_done там уже no-op по своему гейту).
'cancelled' входит в _RESUME_STATUSES, так что потеря этой записи стоила бы точки
возобновления у КАЖДОЙ отмены поэтому гейт `IN ('running', 'cancelled')`, а не
`= 'running'`.
"""
mod = _RUNS_MODULES[name]
db = _RunRowDb(status=status, counters={"done_buckets": ["1:0:5"]})
mod.update_heartbeat(db, 6167, {"done_buckets": ["1:0:5", "1:0:6"]})
assert db.row["counters"]["done_buckets"] == ["1:0:5", "1:0:6"]
assert db.row["heartbeat_at"] is not None
# ── 6. Отмена в теле тика (не в дрейне) помечает in-flight ──────────────────────
async def test_cancel_inside_tick_body_marks_inflight(
caplog: pytest.LogCaptureFixture,
) -> None:
"""Hard-cancel застаёт `scheduler_loop` в await'е ТЕЛА тика, а не в `drain_inflight`.
Прод-достижимо: grace scheduler_main отсчитывается от SIGTERM, а не от нашего дрейна,
и вполне может истечь, пока тик сидит в reap / stale-digest / `_dispatch` (там сеть).
`except Exception` тика CancelledError не ловит (BaseException), а до
`await ctx.drain_inflight()` выполнение уже не доходит без пометки in-flight строки
остаются 'running' и становятся 'zombie' у следующего контейнера, как 07.09.
"""
caplog.set_level(logging.INFO)
table = _Table()
ctx = _ctx(table)
stuck_id = await _dispatch_stuck(ctx, table, asyncio.Event())
loop_task = asyncio.create_task(kit_sched.scheduler_loop(ctx, {}))
for _ in range(3): # дать лупу дойти до первого await ВНУТРИ тика
await asyncio.sleep(0)
assert "scheduler: started" in caplog.text, "отмена обязана застать луп внутри тика"
loop_task.cancel()
with suppress(asyncio.CancelledError):
await loop_task
assert table.rows[stuck_id]["status"] == "done"
assert table.rows[stuck_id]["counters"]["interrupted"] == 1
assert table.rows[stuck_id]["counters"]["lots_fetched"] == 120 # чекпоинт не стёрт
await _cancel_inflight(ctx)
# ── 7. Отказ SQL на одном прогоне не уносит остальные; лог перечисляет помеченных ──
class _AbortingDb:
"""Session, ведущая себя как Postgres после отказавшего statement'а.
До `rollback()` КАЖДЫЙ следующий execute падает («current transaction is aborted»),
поэтому один упавший run_id без сброса сессии утаскивает за собой все остальные.
Моделируется поведение СУБД, а не ожидание теста.
"""
def __init__(self, table: _Table, fail_ids: set[int]) -> None:
self.table = table
self.fail_ids = fail_ids
self.aborted = False
self.closed = False
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
if self.aborted:
raise RuntimeError("current transaction is aborted, commands ignored until rollback")
run_id = int((params or {}).get("id", 0))
if run_id in self.fail_ids:
self.aborted = True
raise RuntimeError("deadlock detected")
row = self.table.rows.get(run_id)
if row is None:
return _FakeResult(fetchone=None)
return _FakeResult(
fetchone=types.SimpleNamespace(status=row["status"], counters=row["counters"])
)
def commit(self) -> None:
pass
def rollback(self) -> None:
self.aborted = False
def close(self) -> None:
self.closed = True
async def test_failed_run_does_not_take_the_others_down(
caplog: pytest.LogCaptureFixture,
) -> None:
"""Первый run_id падает на SELECT — второй всё равно помечен, и в WARNING только он."""
caplog.set_level(logging.INFO)
table = _Table()
first = table.claim("cian_detail_backfill")
second = table.claim("cian_history_backfill")
db = _AbortingDb(table, fail_ids={first})
ctx = SchedulerContext(
config=MagicMock(),
matcher=MagicMock(),
enrichment=MagicMock(),
session_factory=lambda: db,
runs=_FakeRuns(table),
)
tasks = [ctx.spawn_tracked(asyncio.sleep(10), run_id=rid) for rid in (first, second)]
marked = ctx.mark_inflight_interrupted(tasks)
assert marked == 1
assert table.rows[second]["counters"]["interrupted"] == 1
assert table.rows[first]["status"] == "running", "упавший — как был, его снимет boot-reap"
assert f"run_id={first}" in caplog.text, "отказ обязан быть в логе, а не проглочен"
# Честный список: в строке «сняты как interrupted» — только реально помеченные.
marked_line = [r for r in caplog.records if "как interrupted" in r.getMessage()]
assert len(marked_line) == 1
assert str(second) in marked_line[0].getMessage()
assert str(first) not in marked_line[0].getMessage()
assert db.closed
for task in tasks:
task.cancel()
with suppress(asyncio.CancelledError):
await task
async def test_marking_survives_unavailable_db() -> None:
"""БД недоступна: пометка не бросает и не подменяет собой CancelledError.
Вызов приходит из-под `except CancelledError`; исключение отсюда заменило бы отмену
собой, а `suppress(CancelledError)` в scheduler_main такое не глушит процесс ушёл
бы с трейсбеком вместо чистого drain-выхода.
"""
def _boom() -> Any:
raise RuntimeError("could not connect to server")
table = _Table()
run_id = table.claim("cian_detail_backfill")
ctx = SchedulerContext(
config=MagicMock(),
matcher=MagicMock(),
enrichment=MagicMock(),
session_factory=_boom,
runs=_FakeRuns(table),
)
task = ctx.spawn_tracked(asyncio.sleep(10), run_id=run_id)
assert ctx.mark_inflight_interrupted([task]) == 0
assert table.rows[run_id]["status"] == "running"
task.cancel()
with suppress(asyncio.CancelledError):
await task

View file

@ -698,9 +698,19 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, Any]) -> None
точка в БД жила лишь от сохранения бакета до ближайшего тика. Прод-след: 15
оборванных прогонов с доказанной работой (35 706 + 9 222 fetched) и БЕЗ ключа
вообще. Мерж делает точку монотонной для любого писателя, а не только для знающих.
#3391: гейт по статусу — тот же, что у mark_done/mark_failed. Пульс по УЖЕ
финализированной строке холостой по построению, а у app-копии (counters ЗАМЕНЯЕТ,
#3390) ещё и вредный: он стирал метку `interrupted`, которую дрейн поставил живой
ещё задаче, и оборванный прогон снова читался как полный проход.
'cancelled' в гейте НАМЕРЕННО: user-cancel финализирует строку, но задача
останавливается лишь на ближайшей границе якоря, и её последний пульс ЕДИНСТВЕННЫЙ
писатель чекпоинта в этот момент (pipeline.py:1308 / 2368 / 2986 / 4488; mark_done
там уже no-op по своему гейту, а 'cancelled' входит в _RESUME_STATUSES).
"""
total_seen, new_count = _column_counts(counters)
db.execute(
row = db.execute(
text(
"""
UPDATE scrape_runs
@ -708,7 +718,8 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, Any]) -> None
counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb),
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
new_count = COALESCE(CAST(:new_count AS int), new_count)
WHERE id = :run_id
WHERE id = :run_id AND status IN ('running', 'cancelled')
RETURNING id
"""
),
{
@ -717,7 +728,12 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, Any]) -> None
"total_seen": total_seen,
"new_count": new_count,
},
)
).first()
if row is None:
# Не рутина: задача пережила собственную финализацию и продолжает слать прогресс.
logger.warning(
"update_heartbeat no-op: run_id=%d not in 'running'/'cancelled' state", run_id
)
db.commit()

View file

@ -32,7 +32,8 @@ 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 contextlib import suppress
from dataclasses import dataclass, field
from datetime import UTC, datetime, time, timedelta
from typing import TYPE_CHECKING, Any
@ -348,42 +349,141 @@ 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) — голый словарь стёр бы всю бухгалтерию прогона, включая чекпоинт.
Потолок (осознанный, #3391): SELECT синхронный, а у движка нет ни connect-, ни
statement-таймаута (`app/core/db.py:8-19`). Недоступная БД блокирует луп до
SIGKILL'а — 20 с docker-grace (120 s stop_grace_period 100 s _DRAIN_TIMEOUT_S).
Данные при этом НЕ хуже прежних: строки просто остаются 'running' и их снимет
boot-reap следующего контейнера ровно то, что было до этой пометки. Поднимать
до отдельного таймаута есть смысл только вместе с таймаутами на самом движке.
"""
run_ids = [rid for t in tasks if (rid := self._inflight_run_ids.get(t)) is not None]
if not run_ids:
return 0
marked_ids: list[int] = []
db = None
# Открытие сессии и её закрытие — ВНУТРИ try: вызов приходит из-под
# `except CancelledError`, и исключение отсюда ЗАМЕНИЛО бы отмену собой.
# `suppress(CancelledError)` в scheduler_main такое не глушит → процесс уходит
# с трейсбеком вместо чистого drain-выхода. CancelledError (BaseException)
# через `except Exception` проходит насквозь и пробрасывается как есть.
try:
db = self.session_factory()
for run_id in run_ids:
try:
row = db.execute(
text("SELECT status, counters FROM scrape_runs WHERE id = :id"),
{"id": run_id},
).fetchone()
if row is None or row.status != "running":
continue
counters = dict(row.counters or {})
counters["interrupted"] = 1
self.runs.mark_done(db, run_id, counters)
marked_ids.append(run_id)
except Exception:
# Один непроходимый прогон не должен утащить остальные: дрейну
# осталось секунды до SIGKILL, а каждая незакрытая строка — это
# ещё один 'zombie' у следующего контейнера.
logger.exception("scheduler: drain — не удалось пометить run_id=%d", run_id)
# Отказавший statement оставляет сессию в aborted-tx, и КАЖДЫЙ
# следующий run_id падал бы на ровном месте (образец — defensive
# rollback в mark_failed/mark_banned, app/services/scrape_runs.py).
with suppress(Exception):
db.rollback()
except Exception:
logger.exception("scheduler: drain — пометка interrupted не выполнена")
finally:
if db is not None:
with suppress(Exception):
db.close()
if marked_ids:
# Именно marked_ids, а не весь run_ids: в списке не должно быть прогонов,
# пропущенных по статусу, — иначе лог приписывает пометку тем, кого не трогал.
logger.warning(
"scheduler: drain — %d прогон(ов) сняты с 'running' как interrupted: %s",
len(marked_ids),
marked_ids,
)
return len(marked_ids)
async def drain_inflight(self) -> None:
"""Дождаться завершения detached run-задач перед teardown'ом (#1182 P2).
Кооперативные дети дочекивают текущий unit + mark_done и резолвятся;
некооперативные упираются в _CHILD_DRAIN_TIMEOUT_S и остаются на внешний
hard-cancel 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 +997,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
@ -1193,11 +1293,11 @@ def get_due_schedules(db: Session) -> list[dict[str, Any]]:
return [dict(r) for r in rows]
async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler]) -> None:
"""Бесконечный async loop — tick каждые SCHEDULER_TICK_SEC секунд.
async def _tick_loop(ctx: SchedulerContext, registry: Mapping[str, Handler]) -> None:
"""Тело tick-лупа: reap → get_due → dispatch, до выхода по SIGTERM-drain'у.
Структура (reap get_due dispatch SIGTERM-drain) идентична боевому
scheduler'у; 27-веточный if/elif заменён единственным `resolve_handler` + `_dispatch`.
Вынесено из `scheduler_loop` (#3391) ровно затем, чтобы отмену, прилетевшую в ТЕЛО
тика, можно было поймать одним `except` см. вызывающего.
"""
logger.info("scheduler: started (tick=%ds)", SCHEDULER_TICK_SEC)
boot_time = datetime.now(UTC) # граница «моих» прогонов для boot-reap (#3122)
@ -1258,6 +1358,24 @@ async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler])
break
await asyncio.sleep(SCHEDULER_TICK_SEC)
async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler]) -> None:
"""Бесконечный async loop — tick каждые SCHEDULER_TICK_SEC секунд.
Структура (reap get_due dispatch SIGTERM-drain) идентична боевому
scheduler'у; 27-веточный if/elif заменён единственным `resolve_handler` + `_dispatch`.
"""
try:
await _tick_loop(ctx, registry)
except asyncio.CancelledError:
# #3391: hard-cancel из scheduler_main приходит по расписанию grace'а, а не по
# нашему — он вполне может застать ТЕЛО тика (reap / stale-digest / `_dispatch`
# с его сетевым pre_claim), а не `drain_inflight`. `except Exception` тика
# CancelledError не ловит (BaseException), до `await ctx.drain_inflight()` ниже
# выполнение не доходит — без этой ветки in-flight строки остались бы 'running'
# ровно так же, как 07.09 на проде.
ctx.mark_inflight_interrupted([t for t in ctx._inflight_tasks if not t.done()])
raise
# Tick-loop вышел только по SIGTERM-drain'у (иначе while True бесконечен): дожидаемся
# detached run-задач (spawn_tracked) — пусть докоммитят текущий unit и сделают
# mark_done, а не дадим asyncio.run() teardown'у хард-кансельнуть их mid-await (#1182).