diff --git a/tradein-mvp/backend/app/scheduler_main.py b/tradein-mvp/backend/app/scheduler_main.py index 59254465..57e17551 100644 --- a/tradein-mvp/backend/app/scheduler_main.py +++ b/tradein-mvp/backend/app/scheduler_main.py @@ -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: diff --git a/tradein-mvp/backend/app/services/scrape_runs.py b/tradein-mvp/backend/app/services/scrape_runs.py index 4475253c..ba5d65d8 100644 --- a/tradein-mvp/backend/app/services/scrape_runs.py +++ b/tradein-mvp/backend/app/services/scrape_runs.py @@ -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() diff --git a/tradein-mvp/backend/tests/test_3391_drain_marks_inflight_app_tasks.py b/tradein-mvp/backend/tests/test_3391_drain_marks_inflight_app_tasks.py new file mode 100644 index 00000000..631c3043 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3391_drain_marks_inflight_app_tasks.py @@ -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 diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py index 3ebebdbd..85fcacf9 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py @@ -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() diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py index 711ddd0e..a30ea2a1 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py @@ -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).