fix(scraper-kit/scheduler): SIGTERM-drain помечает in-flight прогоны interrupted=1 при hard-cancel; честный лог drain'а (#3391) #3392
5 changed files with 716 additions and 19 deletions
|
|
@ -101,15 +101,23 @@ def _should_run() -> bool:
|
||||||
return settings.scheduler_enable
|
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'ом.
|
"""Дождаться завершения scheduler-задачи с кооперативным SIGTERM-drain'ом.
|
||||||
|
|
||||||
|
Возвращает True, если задачу пришлось хард-кансельнуть (grace истёк), False — если
|
||||||
|
она вышла сама. Вызывающий обязан различать эти исходы в логах (#3391): строка
|
||||||
|
«drained cleanly» после hard-cancel'а — ложь, а именно она печаталась на проде
|
||||||
|
07.09 сразу за WARNING'ом о превышении grace.
|
||||||
|
|
||||||
- Нет shutdown: scheduler_loop бесконечен → задача никогда не завершается, ждём её
|
- Нет shutdown: scheduler_loop бесконечен → задача никогда не завершается, ждём её
|
||||||
как есть (процесс просто работает).
|
как есть (процесс просто работает).
|
||||||
- shutdown запрошен, задача ещё бежит: bounded `wait_for(_DRAIN_TIMEOUT_S)` —
|
- shutdown запрошен, задача ещё бежит: bounded `wait_for(_DRAIN_TIMEOUT_S)` —
|
||||||
кооперативная задача докоммитит текущий unit на ближайшем checkpoint'е и выйдет
|
кооперативная задача докоммитит текущий unit на ближайшем checkpoint'е и выйдет
|
||||||
сама. Если превысила grace → hard-cancel + suppress CancelledError, чтобы
|
сама. Если превысила grace → hard-cancel + suppress CancelledError, чтобы
|
||||||
некооперативная задача не подвесила процесс за пределами docker stop_grace_period.
|
некооперативная задача не подвесила процесс за пределами 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'а
|
# Гонка «задача завершилась сама» против «пришёл SIGTERM»: отсчёт safety-net'а
|
||||||
# должен стартовать от МОМЕНТА запроса drain'а, а не от старта процесса.
|
# должен стартовать от МОМЕНТА запроса drain'а, а не от старта процесса.
|
||||||
|
|
@ -127,7 +135,7 @@ async def _await_scheduler(task: asyncio.Task[None]) -> None:
|
||||||
# упал (сохраняем прежнюю loud-crash семантику `await task`), а не глушит его.
|
# упал (сохраняем прежнюю loud-crash семантику `await task`), а не глушит его.
|
||||||
task.result()
|
task.result()
|
||||||
logger.info("scheduler_main: scheduler task exited cleanly")
|
logger.info("scheduler_main: scheduler task exited cleanly")
|
||||||
return
|
return False
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"scheduler_main: SIGTERM-drain — waiting up to %.0fs for in-flight unit to commit",
|
"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:
|
try:
|
||||||
await asyncio.wait_for(task, timeout=_DRAIN_TIMEOUT_S)
|
await asyncio.wait_for(task, timeout=_DRAIN_TIMEOUT_S)
|
||||||
logger.info("scheduler_main: scheduler drained and exited cleanly")
|
logger.info("scheduler_main: scheduler drained and exited cleanly")
|
||||||
|
return False
|
||||||
except TimeoutError:
|
except TimeoutError:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"scheduler_main: drain exceeded %.0fs grace — hard-cancelling scheduler task",
|
"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()
|
task.cancel()
|
||||||
with suppress(asyncio.CancelledError):
|
with suppress(asyncio.CancelledError):
|
||||||
await task
|
await task
|
||||||
|
return True
|
||||||
|
|
||||||
|
|
||||||
async def _run_kit_scheduler() -> None:
|
async def _run_kit_scheduler() -> None:
|
||||||
|
|
@ -224,8 +234,12 @@ async def _run() -> None:
|
||||||
# Windows dev: signal handlers через loop не поддерживаются
|
# Windows dev: signal handlers через loop не поддерживаются
|
||||||
logger.warning("scheduler_main: loop.add_signal_handler not supported (Windows dev)")
|
logger.warning("scheduler_main: loop.add_signal_handler not supported (Windows dev)")
|
||||||
|
|
||||||
await _await_scheduler(task)
|
hard_cancelled = await _await_scheduler(task)
|
||||||
|
|
||||||
|
if hard_cancelled:
|
||||||
|
# Дрейн НЕ был чистым: задача не вышла сама, её сняли. WARNING об этом уже
|
||||||
|
# напечатан в _await_scheduler — второй строкой её не «переобъявляем».
|
||||||
|
return
|
||||||
if shutdown_requested():
|
if shutdown_requested():
|
||||||
logger.info("scheduler_main: scheduler drained cleanly (SIGTERM)")
|
logger.info("scheduler_main: scheduler drained cleanly (SIGTERM)")
|
||||||
else:
|
else:
|
||||||
|
|
|
||||||
|
|
@ -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) и
|
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и
|
||||||
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
|
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
|
||||||
COALESCE: если ключа нет в counters — старое значение колонки сохраняется.
|
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)
|
total_seen, new_count = _column_counts(counters)
|
||||||
db.execute(
|
row = db.execute(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
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),
|
counters = CAST(:counters AS jsonb),
|
||||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
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,6 +633,11 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None
|
||||||
"total_seen": total_seen,
|
"total_seen": total_seen,
|
||||||
"new_count": new_count,
|
"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()
|
db.commit()
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
|
@ -698,9 +698,19 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, Any]) -> None
|
||||||
точка в БД жила лишь от сохранения бакета до ближайшего тика. Прод-след: 15
|
точка в БД жила лишь от сохранения бакета до ближайшего тика. Прод-след: 15
|
||||||
оборванных прогонов с доказанной работой (35 706 + 9 222 fetched) и БЕЗ ключа
|
оборванных прогонов с доказанной работой (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)
|
total_seen, new_count = _column_counts(counters)
|
||||||
db.execute(
|
row = db.execute(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
UPDATE scrape_runs
|
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),
|
counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb),
|
||||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
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,6 +728,11 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, Any]) -> None
|
||||||
"total_seen": total_seen,
|
"total_seen": total_seen,
|
||||||
"new_count": new_count,
|
"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()
|
db.commit()
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -32,7 +32,8 @@ from __future__ import annotations
|
||||||
import asyncio
|
import asyncio
|
||||||
import logging
|
import logging
|
||||||
import random
|
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 dataclasses import dataclass, field
|
||||||
from datetime import UTC, datetime, time, timedelta
|
from datetime import UTC, datetime, time, timedelta
|
||||||
from typing import TYPE_CHECKING, Any
|
from typing import TYPE_CHECKING, Any
|
||||||
|
|
@ -348,42 +349,141 @@ class SchedulerContext:
|
||||||
runs: Any = _kit_runs
|
runs: Any = _kit_runs
|
||||||
# #1182 P2: strong-ref'ы detached run-задач для graceful SIGTERM-drain'а.
|
# #1182 P2: strong-ref'ы detached run-задач для graceful SIGTERM-drain'а.
|
||||||
_inflight_tasks: set[asyncio.Task[None]] = field(default_factory=set)
|
_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).
|
"""create_task + регистрация в _inflight_tasks для graceful-drain'а (#1182 P2).
|
||||||
|
|
||||||
strong-ref в set'е держит задачу до завершения (RUF006) и даёт scheduler_loop'у
|
strong-ref в set'е держит задачу до завершения (RUF006) и даёт scheduler_loop'у
|
||||||
дождаться её на SIGTERM-drain'е; done-callback ретривит exception и убирает задачу.
|
дождаться её на SIGTERM-drain'е; done-callback ретривит exception и убирает задачу.
|
||||||
|
|
||||||
|
`run_id` (#3391) — строка scrape_runs этой задачи; None у задач без прогона.
|
||||||
"""
|
"""
|
||||||
task = asyncio.create_task(coro)
|
task = asyncio.create_task(coro)
|
||||||
self._inflight_tasks.add(task)
|
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:
|
def _on_done(t: asyncio.Task[None]) -> None:
|
||||||
self._inflight_tasks.discard(t)
|
self._inflight_tasks.discard(t)
|
||||||
|
self._inflight_run_ids.pop(t, None)
|
||||||
if not t.cancelled():
|
if not t.cancelled():
|
||||||
t.exception()
|
t.exception()
|
||||||
|
|
||||||
task.add_done_callback(_on_done)
|
task.add_done_callback(_on_done)
|
||||||
return task
|
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:
|
async def drain_inflight(self) -> None:
|
||||||
"""Дождаться завершения detached run-задач перед teardown'ом (#1182 P2).
|
"""Дождаться завершения detached run-задач перед teardown'ом (#1182 P2).
|
||||||
|
|
||||||
Кооперативные дети дочекивают текущий unit + mark_done и резолвятся;
|
Кооперативные дети дочекивают текущий unit + mark_done и резолвятся;
|
||||||
некооперативные упираются в _CHILD_DRAIN_TIMEOUT_S и остаются на внешний
|
некооперативные упираются в _CHILD_DRAIN_TIMEOUT_S либо во внешний hard-cancel
|
||||||
hard-cancel fallback. Не busy-spin — один asyncio.wait.
|
из scheduler_main — и в обоих случаях их строки помечаются `interrupted`
|
||||||
|
(#3391), иначе они доживают в 'running' до boot-reap'а. Не busy-spin — один
|
||||||
|
asyncio.wait.
|
||||||
"""
|
"""
|
||||||
pending = [t for t in self._inflight_tasks if not t.done()]
|
pending = [t for t in self._inflight_tasks if not t.done()]
|
||||||
if not pending:
|
if not pending:
|
||||||
return
|
return
|
||||||
logger.info("scheduler: draining %d in-flight run task(s) on shutdown", len(pending))
|
logger.info("scheduler: draining %d in-flight run task(s) on shutdown", len(pending))
|
||||||
|
try:
|
||||||
_done, still = await asyncio.wait(pending, timeout=_CHILD_DRAIN_TIMEOUT_S)
|
_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:
|
if still:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"scheduler: %d task(s) did not drain in %.0fs — leaving for hard-cancel",
|
"scheduler: %d task(s) did not drain in %.0fs — leaving for hard-cancel",
|
||||||
len(still),
|
len(still),
|
||||||
_CHILD_DRAIN_TIMEOUT_S,
|
_CHILD_DRAIN_TIMEOUT_S,
|
||||||
)
|
)
|
||||||
|
self.mark_inflight_interrupted(still)
|
||||||
else:
|
else:
|
||||||
logger.info("scheduler: all in-flight run task(s) drained cleanly")
|
logger.info("scheduler: all in-flight run task(s) drained cleanly")
|
||||||
|
|
||||||
|
|
@ -897,7 +997,7 @@ async def _dispatch(
|
||||||
finally:
|
finally:
|
||||||
run_db.close()
|
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)
|
logger.info("scheduler: triggered %s run_id=%d", handler.log_name, run_id)
|
||||||
return 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]
|
return [dict(r) for r in rows]
|
||||||
|
|
||||||
|
|
||||||
async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler]) -> None:
|
async def _tick_loop(ctx: SchedulerContext, registry: Mapping[str, Handler]) -> None:
|
||||||
"""Бесконечный async loop — tick каждые SCHEDULER_TICK_SEC секунд.
|
"""Тело tick-лупа: reap → get_due → dispatch, до выхода по SIGTERM-drain'у.
|
||||||
|
|
||||||
Структура (reap → get_due → dispatch → SIGTERM-drain) идентична боевому
|
Вынесено из `scheduler_loop` (#3391) ровно затем, чтобы отмену, прилетевшую в ТЕЛО
|
||||||
scheduler'у; 27-веточный if/elif заменён единственным `resolve_handler` + `_dispatch`.
|
тика, можно было поймать одним `except` — см. вызывающего.
|
||||||
"""
|
"""
|
||||||
logger.info("scheduler: started (tick=%ds)", SCHEDULER_TICK_SEC)
|
logger.info("scheduler: started (tick=%ds)", SCHEDULER_TICK_SEC)
|
||||||
boot_time = datetime.now(UTC) # граница «моих» прогонов для boot-reap (#3122)
|
boot_time = datetime.now(UTC) # граница «моих» прогонов для boot-reap (#3122)
|
||||||
|
|
@ -1258,6 +1358,24 @@ async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler])
|
||||||
break
|
break
|
||||||
await asyncio.sleep(SCHEDULER_TICK_SEC)
|
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 бесконечен): дожидаемся
|
# Tick-loop вышел только по SIGTERM-drain'у (иначе while True бесконечен): дожидаемся
|
||||||
# detached run-задач (spawn_tracked) — пусть докоммитят текущий unit и сделают
|
# detached run-задач (spawn_tracked) — пусть докоммитят текущий unit и сделают
|
||||||
# mark_done, а не дадим asyncio.run() teardown'у хард-кансельнуть их mid-await (#1182).
|
# mark_done, а не дадим asyncio.run() teardown'у хард-кансельнуть их mid-await (#1182).
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue