fix(tradein/scheduler): SIGTERM-drain снимает с 'running' in-flight app-task'и (#3391)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 10s
CI / changes (pull_request) Successful in 12s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m55s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 10s
CI / changes (pull_request) Successful in 12s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m55s
Прод 07.09 02:36 UTC, первый настоящий drain после #3363: hard-cancel из scheduler_main оборвал дрейн, и два бэкфилла (cian_detail_backfill 6167, cian_history_backfill 6173) остались в scrape_runs со статусом 'running' — boot-reap следующего контейнера сделал их 'zombie' (boot_reaped=true), метки interrupted не было. interrupted=1 при дрейне писали только kit-пайплайны и DKP-импорт: у задач, чьё тело живёт в app, ставить её было некому. Метка ставится в единственной точке, через которую проходит любая detached run-задача — SchedulerContext.drain_inflight: и по истечении _CHILD_DRAIN_TIMEOUT_S, и в обработчике CancelledError (тот самый прод-путь). run_id берётся из нового реестра {task: run_id}, который заполняет _dispatch сразу после claim'а; claim-логика не тронута. Статус строки перечитывается перед записью, поэтому успевший финализироваться сам прогон не перезаписывается, а counters читаются из строки и дописываются — app-копия mark_done их ЗАМЕНЯЕТ (#3390), голая {"interrupted": 1} стёрла бы чекпоинт. scheduler_main: _await_scheduler возвращает признак hard-cancel'а, и строка «scheduler drained cleanly (SIGTERM)» больше не печатается сразу за WARNING'ом о превышении grace — на проде эти две строки стояли подряд и противоречили друг другу. Запас времени на запись: docker stop_grace_period 120s − _DRAIN_TIMEOUT_S 100s = 20 с после hard-cancel'а, запись синхронная (несколько statement'ов).
This commit is contained in:
parent
3fe2d310b2
commit
30e3bacc5e
3 changed files with 364 additions and 9 deletions
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -0,0 +1,266 @@
|
|||
"""SIGTERM-drain снимает с 'running' прогоны app-task'ов, не доживших до финализации (#3391).
|
||||
|
||||
Прод 07.09 02:36 UTC (первый настоящий SIGTERM-drain после #3363): деплой пересоздал
|
||||
`tradein-scraper`, kit-scheduler напечатал «draining 2 in-flight run task(s)», через 100 с
|
||||
`scheduler_main` хард-кансельнул задачу — и оба бэкфилла (cian_detail_backfill 6167,
|
||||
cian_history_backfill 6173) остались в scrape_runs со статусом 'running' до boot-reap'а
|
||||
следующего контейнера, где стали 'zombie' с `boot_reaped=true`. `interrupted=1` при дрейне
|
||||
пишут только kit-пайплайны (#3363) и DKP-импорт: у задач, чьё тело живёт в `app`, метку
|
||||
ставить некому.
|
||||
|
||||
Проверяется НЕ «функция вызвалась», а значение в строке прогона:
|
||||
1. grace истёк → строка застрявшей задачи 'done' + counters.interrupted=1, прежние
|
||||
счётчики (в т.ч. чекпоинт) на месте;
|
||||
2. соседняя задача, успевшая финализироваться сама, НЕ перезаписана (без interrupted);
|
||||
3. hard-cancel из scheduler_main (CancelledError прямо в `asyncio.wait` дрейна) даёт тот
|
||||
же результат — это ровно прод-путь 07.09, ветка таймаута его не покрывает;
|
||||
4. `scheduler_main` после hard-cancel'а больше не печатает «drained cleanly».
|
||||
|
||||
Фейковый `mark_done` counters ЗАМЕНЯЕТ, а не мержит — так ведёт себя боевая app-копия
|
||||
(`app.services.scrape_runs`, #3390), которая и инжектируется в прод-контекст. Голый
|
||||
`{"interrupted": 1}` на этом фейке стёр бы чекпоинт, и (1) покраснел бы.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
import types
|
||||
from contextlib import suppress
|
||||
from typing import Any
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
||||
|
||||
import pytest
|
||||
from scraper_kit.orchestration import scheduler as kit_sched
|
||||
from scraper_kit.orchestration.scheduler import Handler, SchedulerContext, _dispatch
|
||||
|
||||
import app.scheduler_main as sm
|
||||
from app.core import shutdown as sd
|
||||
|
||||
|
||||
class _Table:
|
||||
"""Мини-`scrape_runs`: id → {status, counters}."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.rows: dict[int, dict[str, Any]] = {}
|
||||
self._next_id = 6167
|
||||
|
||||
def claim(self, source: str) -> int:
|
||||
run_id = self._next_id
|
||||
self._next_id += 6
|
||||
self.rows[run_id] = {"status": "running", "counters": {}, "source": source}
|
||||
return run_id
|
||||
|
||||
|
||||
class _FakeResult:
|
||||
def __init__(self, *, fetchone: Any = None, scalar: Any = None) -> None:
|
||||
self._fetchone = fetchone
|
||||
self._scalar = scalar
|
||||
|
||||
def fetchone(self) -> Any:
|
||||
return self._fetchone
|
||||
|
||||
def scalar(self) -> Any:
|
||||
return self._scalar
|
||||
|
||||
|
||||
class _FakeDb:
|
||||
"""Session поверх `_Table`: claim-путь `_claim_run` + SELECT статуса дрейна."""
|
||||
|
||||
def __init__(self, table: _Table, *, running_states: list[bool] | None = None) -> None:
|
||||
self.table = table
|
||||
self._running = list(running_states or [])
|
||||
self.closed = False
|
||||
|
||||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
|
||||
sql = str(stmt)
|
||||
if "SELECT status, counters FROM scrape_runs" in sql:
|
||||
row = self.table.rows.get(int((params or {})["id"]))
|
||||
if row is None:
|
||||
return _FakeResult(fetchone=None)
|
||||
return _FakeResult(
|
||||
fetchone=types.SimpleNamespace(status=row["status"], counters=row["counters"])
|
||||
)
|
||||
if "SELECT 1 FROM scrape_runs" in sql:
|
||||
return _FakeResult(fetchone=(1,) if self._running.pop(0) else None)
|
||||
if "pg_try_advisory_xact_lock" in sql:
|
||||
return _FakeResult(scalar=True)
|
||||
return _FakeResult()
|
||||
|
||||
def commit(self) -> None:
|
||||
pass
|
||||
|
||||
def rollback(self) -> None:
|
||||
pass
|
||||
|
||||
def close(self) -> None:
|
||||
self.closed = True
|
||||
|
||||
|
||||
class _FakeRuns:
|
||||
"""ctx.runs: create_run + mark_done с боевым гейтом `WHERE status = 'running'`."""
|
||||
|
||||
def __init__(self, table: _Table) -> None:
|
||||
self.table = table
|
||||
|
||||
def create_run(self, db: Any, *, source: str, params: dict[str, Any]) -> int:
|
||||
return self.table.claim(source)
|
||||
|
||||
def mark_done(self, db: Any, run_id: int, counters: dict[str, Any]) -> None:
|
||||
row = self.table.rows.get(run_id)
|
||||
if row is None or row["status"] != "running":
|
||||
return # боевой UPDATE ... WHERE status = 'running' — no-op
|
||||
row["status"] = "done"
|
||||
row["counters"] = dict(counters) # app-копия ЗАМЕНЯЕТ counters (#3390)
|
||||
|
||||
|
||||
def _make_sched(source: str) -> dict[str, Any]:
|
||||
return {
|
||||
"id": 1,
|
||||
"source": source,
|
||||
"enabled": True,
|
||||
"window_start_hour": 2,
|
||||
"window_end_hour": 5,
|
||||
"default_params": {},
|
||||
}
|
||||
|
||||
|
||||
def _ctx(table: _Table) -> SchedulerContext:
|
||||
return SchedulerContext(
|
||||
config=MagicMock(),
|
||||
matcher=MagicMock(),
|
||||
enrichment=MagicMock(),
|
||||
session_factory=lambda: _FakeDb(table),
|
||||
runs=_FakeRuns(table),
|
||||
)
|
||||
|
||||
|
||||
async def _dispatch_stuck(ctx: SchedulerContext, table: _Table, started: asyncio.Event) -> int:
|
||||
"""app-task, который на дрейн не реагирует (долгий await вместо checkpoint'а)."""
|
||||
|
||||
async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
|
||||
started.set()
|
||||
await asyncio.sleep(10)
|
||||
|
||||
run_id = await _dispatch(
|
||||
Handler(_job, "cian_detail_backfill"),
|
||||
_FakeDb(table, running_states=[False, False]),
|
||||
_make_sched("cian_detail_backfill"),
|
||||
ctx,
|
||||
)
|
||||
assert run_id is not None
|
||||
await asyncio.wait_for(started.wait(), timeout=2.0)
|
||||
# Пульс до дрейна: собранное этим прогоном и его чекпоинт.
|
||||
table.rows[run_id]["counters"] = {"lots_fetched": 120, "done_buckets": ["1:0:5"]}
|
||||
return run_id
|
||||
|
||||
|
||||
async def _cancel_inflight(ctx: SchedulerContext) -> None:
|
||||
for task in list(ctx._inflight_tasks):
|
||||
task.cancel()
|
||||
with suppress(asyncio.CancelledError):
|
||||
await task
|
||||
|
||||
|
||||
# ── 1-2. grace истёк: застрявший помечен, самофинализировавшийся не тронут ───────
|
||||
|
||||
|
||||
async def test_drain_grace_expiry_marks_stuck_run_interrupted(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
monkeypatch.setattr(kit_sched, "_CHILD_DRAIN_TIMEOUT_S", 0.05)
|
||||
table = _Table()
|
||||
ctx = _ctx(table)
|
||||
|
||||
stuck_id = await _dispatch_stuck(ctx, table, asyncio.Event())
|
||||
|
||||
finished = asyncio.Event()
|
||||
|
||||
async def _quick(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
|
||||
c.runs.mark_done(db, run_id, {"lots_fetched": 5})
|
||||
finished.set()
|
||||
|
||||
quick_id = await _dispatch(
|
||||
Handler(_quick, "cian_history_backfill"),
|
||||
_FakeDb(table, running_states=[False, False]),
|
||||
_make_sched("cian_history_backfill"),
|
||||
ctx,
|
||||
)
|
||||
await asyncio.wait_for(finished.wait(), timeout=2.0)
|
||||
|
||||
await ctx.drain_inflight()
|
||||
|
||||
assert table.rows[stuck_id]["status"] == "done"
|
||||
assert table.rows[stuck_id]["counters"]["interrupted"] == 1
|
||||
# Чекпоинт и собранное пережили пометку — иначе резюм подхватывать нечего.
|
||||
assert table.rows[stuck_id]["counters"]["done_buckets"] == ["1:0:5"]
|
||||
assert table.rows[stuck_id]["counters"]["lots_fetched"] == 120
|
||||
# Финализировавшийся сам — не перезаписан: ни статуса, ни лишней метки.
|
||||
assert table.rows[quick_id]["status"] == "done"
|
||||
assert table.rows[quick_id]["counters"] == {"lots_fetched": 5}
|
||||
|
||||
await _cancel_inflight(ctx)
|
||||
|
||||
|
||||
# ── 3. hard-cancel из scheduler_main — прод-путь 07.09 ──────────────────────────
|
||||
|
||||
|
||||
async def test_drain_hard_cancel_marks_stuck_run_interrupted(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
"""CancelledError прилетает прямо в `asyncio.wait` дрейна (grace scheduler_main истёк)."""
|
||||
monkeypatch.setattr(kit_sched, "_CHILD_DRAIN_TIMEOUT_S", 30.0)
|
||||
table = _Table()
|
||||
ctx = _ctx(table)
|
||||
stuck_id = await _dispatch_stuck(ctx, table, asyncio.Event())
|
||||
|
||||
drain = asyncio.create_task(ctx.drain_inflight())
|
||||
await asyncio.sleep(0) # дать дрейну дойти до asyncio.wait
|
||||
await asyncio.sleep(0)
|
||||
drain.cancel()
|
||||
with suppress(asyncio.CancelledError):
|
||||
await drain
|
||||
|
||||
assert table.rows[stuck_id]["status"] == "done"
|
||||
assert table.rows[stuck_id]["counters"]["interrupted"] == 1
|
||||
|
||||
await _cancel_inflight(ctx)
|
||||
|
||||
|
||||
# ── 4. scheduler_main: после hard-cancel'а «drained cleanly» не печатается ───────
|
||||
|
||||
|
||||
@pytest.fixture()
|
||||
def _reset_shutdown() -> Any:
|
||||
sd.reset_shutdown()
|
||||
yield
|
||||
sd.reset_shutdown()
|
||||
|
||||
|
||||
async def test_hard_cancel_does_not_log_drained_cleanly(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
_reset_shutdown: Any,
|
||||
) -> None:
|
||||
"""Прод 07.09: WARNING о hard-cancel'е и INFO «drained cleanly» стояли подряд."""
|
||||
caplog.set_level(logging.INFO)
|
||||
monkeypatch.setattr(sm, "_DRAIN_TIMEOUT_S", 0.05)
|
||||
started = asyncio.Event()
|
||||
|
||||
async def _stub_loop() -> None:
|
||||
started.set()
|
||||
await asyncio.sleep(30) # некооперативная задача: shutdown не смотрит
|
||||
|
||||
monkeypatch.setattr(sm, "_run_kit_scheduler", _stub_loop)
|
||||
|
||||
task = asyncio.create_task(sm._run())
|
||||
await asyncio.wait_for(started.wait(), timeout=2.0)
|
||||
sd.request_shutdown()
|
||||
await asyncio.wait_for(task, timeout=5.0)
|
||||
|
||||
assert "hard-cancelling scheduler task" in caplog.text
|
||||
assert "drained cleanly" not in caplog.text
|
||||
assert "drained and exited cleanly" not in caplog.text
|
||||
|
|
@ -32,7 +32,7 @@ from __future__ import annotations
|
|||
import asyncio
|
||||
import logging
|
||||
import random
|
||||
from collections.abc import Callable, Coroutine, Mapping
|
||||
from collections.abc import Callable, Coroutine, Iterable, Mapping
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import UTC, datetime, time, timedelta
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
|
@ -348,42 +348,117 @@ class SchedulerContext:
|
|||
runs: Any = _kit_runs
|
||||
# #1182 P2: strong-ref'ы detached run-задач для graceful SIGTERM-drain'а.
|
||||
_inflight_tasks: set[asyncio.Task[None]] = field(default_factory=set)
|
||||
# #3391: task → run_id заклеймленного прогона. Без него дрейн знает, что задача
|
||||
# не дожила до финализации, но не знает, КАКУЮ строку scrape_runs снимать с
|
||||
# 'running' (claim-логика run_id никуда не публиковала).
|
||||
_inflight_run_ids: dict[asyncio.Task[None], int] = field(default_factory=dict)
|
||||
|
||||
def spawn_tracked(self, coro: Coroutine[Any, Any, None]) -> asyncio.Task[None]:
|
||||
def spawn_tracked(
|
||||
self, coro: Coroutine[Any, Any, None], *, run_id: int | None = None
|
||||
) -> asyncio.Task[None]:
|
||||
"""create_task + регистрация в _inflight_tasks для graceful-drain'а (#1182 P2).
|
||||
|
||||
strong-ref в set'е держит задачу до завершения (RUF006) и даёт scheduler_loop'у
|
||||
дождаться её на SIGTERM-drain'е; done-callback ретривит exception и убирает задачу.
|
||||
|
||||
`run_id` (#3391) — строка scrape_runs этой задачи; None у задач без прогона.
|
||||
"""
|
||||
task = asyncio.create_task(coro)
|
||||
self._inflight_tasks.add(task)
|
||||
if run_id is not None:
|
||||
self._inflight_run_ids[task] = run_id
|
||||
|
||||
def _on_done(t: asyncio.Task[None]) -> None:
|
||||
self._inflight_tasks.discard(t)
|
||||
self._inflight_run_ids.pop(t, None)
|
||||
if not t.cancelled():
|
||||
t.exception()
|
||||
|
||||
task.add_done_callback(_on_done)
|
||||
return task
|
||||
|
||||
def mark_inflight_interrupted(self, tasks: Iterable[asyncio.Task[None]]) -> int:
|
||||
"""Снять с 'running' прогоны задач, не доживших до собственной финализации (#3391).
|
||||
|
||||
Прод 07.09 (первый настоящий SIGTERM-drain после #3363): два app-task бэкфилла
|
||||
(cian_detail_backfill 6167, cian_history_backfill 6173) не уложились в grace,
|
||||
scheduler_main их хард-кансельнул — и строки остались 'running' до boot-reap'а
|
||||
следующего контейнера, где стали 'zombie' с `boot_reaped=true`. `interrupted=1`
|
||||
при дрейне пишут только kit-пайплайны (#3363) и DKP-импорт: у задач, чьё тело
|
||||
живёт в app и о дрейне не знает, метку ставить некому. Ставим её ЗДЕСЬ — в
|
||||
единственной точке, через которую проходит любая detached run-задача.
|
||||
|
||||
Синхронно (sync-сессия, ни одного await): вызывается в том числе из except
|
||||
CancelledError, где лишний await мог бы не дожить до конца.
|
||||
|
||||
Статус строки перечитывается перед записью: задача, успевшая финализироваться
|
||||
сама (done/failed/banned), не перезаписывается. Гонку добивает сам `mark_done`
|
||||
(`WHERE status = 'running'`) — но тогда счётчики уже прочитаны, и их merge был бы
|
||||
холостым. counters берём из строки и дописываем `interrupted`, а не отдаём
|
||||
`{"interrupted": 1}` голым: app-копия `mark_done` counters ЗАМЕНЯЕТ, а не мержит
|
||||
(#3390) — голый словарь стёр бы всю бухгалтерию прогона, включая чекпоинт.
|
||||
"""
|
||||
run_ids = [rid for t in tasks if (rid := self._inflight_run_ids.get(t)) is not None]
|
||||
if not run_ids:
|
||||
return 0
|
||||
marked = 0
|
||||
db = self.session_factory()
|
||||
try:
|
||||
for run_id in run_ids:
|
||||
try:
|
||||
row = db.execute(
|
||||
text("SELECT status, counters FROM scrape_runs WHERE id = :id"),
|
||||
{"id": run_id},
|
||||
).fetchone()
|
||||
if row is None or row.status != "running":
|
||||
continue
|
||||
counters = dict(row.counters or {})
|
||||
counters["interrupted"] = 1
|
||||
self.runs.mark_done(db, run_id, counters)
|
||||
marked += 1
|
||||
except Exception:
|
||||
# Один непроходимый прогон не должен утащить остальные: дрейну
|
||||
# осталось секунды до SIGKILL, а каждая незакрытая строка — это
|
||||
# ещё один 'zombie' у следующего контейнера.
|
||||
logger.exception("scheduler: drain — не удалось пометить run_id=%d", run_id)
|
||||
finally:
|
||||
db.close()
|
||||
if marked:
|
||||
logger.warning(
|
||||
"scheduler: drain — %d прогон(ов) сняты с 'running' как interrupted: %s",
|
||||
marked,
|
||||
run_ids,
|
||||
)
|
||||
return marked
|
||||
|
||||
async def drain_inflight(self) -> None:
|
||||
"""Дождаться завершения detached run-задач перед teardown'ом (#1182 P2).
|
||||
|
||||
Кооперативные дети дочекивают текущий unit + mark_done и резолвятся;
|
||||
некооперативные упираются в _CHILD_DRAIN_TIMEOUT_S и остаются на внешний
|
||||
hard-cancel fallback. Не busy-spin — один asyncio.wait.
|
||||
некооперативные упираются в _CHILD_DRAIN_TIMEOUT_S либо во внешний hard-cancel
|
||||
из scheduler_main — и в обоих случаях их строки помечаются `interrupted`
|
||||
(#3391), иначе они доживают в 'running' до boot-reap'а. Не busy-spin — один
|
||||
asyncio.wait.
|
||||
"""
|
||||
pending = [t for t in self._inflight_tasks if not t.done()]
|
||||
if not pending:
|
||||
return
|
||||
logger.info("scheduler: draining %d in-flight run task(s) on shutdown", len(pending))
|
||||
_done, still = await asyncio.wait(pending, timeout=_CHILD_DRAIN_TIMEOUT_S)
|
||||
try:
|
||||
_done, still = await asyncio.wait(pending, timeout=_CHILD_DRAIN_TIMEOUT_S)
|
||||
except asyncio.CancelledError:
|
||||
# Hard-cancel из scheduler_main (grace 100s истёк раньше нашего дрейна —
|
||||
# такт tick-сна + _CHILD_DRAIN_TIMEOUT_S могут превысить его). Пишем метку
|
||||
# синхронно и пробрасываем отмену дальше.
|
||||
self.mark_inflight_interrupted([t for t in pending if not t.done()])
|
||||
raise
|
||||
if still:
|
||||
logger.warning(
|
||||
"scheduler: %d task(s) did not drain in %.0fs — leaving for hard-cancel",
|
||||
len(still),
|
||||
_CHILD_DRAIN_TIMEOUT_S,
|
||||
)
|
||||
self.mark_inflight_interrupted(still)
|
||||
else:
|
||||
logger.info("scheduler: all in-flight run task(s) drained cleanly")
|
||||
|
||||
|
|
@ -897,7 +972,7 @@ async def _dispatch(
|
|||
finally:
|
||||
run_db.close()
|
||||
|
||||
ctx.spawn_tracked(_run())
|
||||
ctx.spawn_tracked(_run(), run_id=run_id)
|
||||
logger.info("scheduler: triggered %s run_id=%d", handler.log_name, run_id)
|
||||
return run_id
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue