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

Merged
bot-backend merged 2 commits from fix/drain-marks-inflight-app-tasks into main 2026-09-06 04:14:07 +00:00
Collaborator

Closes #3391.

Дефект (прод 07.09 02:36 UTC, первый реальный SIGTERM-drain после #3363). Деплой пересоздал scraper при идущих cian_detail_backfill 6167 (75 мин, 162 объявления) и cian_history_backfill 6173; через 100 с grace — hard-cancel; строки остались running, boot-reap нового контейнера сделал их zombie без interrupted. interrupted=1 на SIGTERM писали только kit-пайплайны и DKP; app-task бэкфиллы — нет. Плюс scheduler_main печатал «scheduler drained cleanly (SIGTERM)» и после hard-cancel.

Фикс в одном месте (kit-scheduler). orchestration/scheduler.py:354-380 — реестр _inflight_run_ids: dict[Task,int] через spawn_tracked(coro, run_id=…) из _dispatch (:975; claim не тронут), done-callback чистит. mark_inflight_interrupted(tasks) (:380, синхронная — без await, живёт внутри except CancelledError): для каждого run_id SELECT status, counters; если runningcounters["interrupted"]=1 поверх прочитанных и mark_done. Вызывается и по истечении grace (:461), и в except asyncio.CancelledError (:449, затем raise) — на проде hard-cancel прилетает внутрь asyncio.wait (хвост tick-сна до 60 с + 80 с дрейна > 100 с grace). Counters читаются из строки, потому что боевой ctx.runsapp.services.scrape_runs, где mark_done counters заменяет (#3390): голый {"interrupted":1} стёр бы done_buckets. «Не трогать успевших» — проверка статуса + собственный гейт mark_done WHERE status='running'.

app/scheduler_main.py:156,237-244_await_scheduler возвращает признак hard-cancel; «drained cleanly» печатается только когда задача вышла сама.

stop_grace_period: 120 с у scraper против _DRAIN_TIMEOUT_S=100 → 20 с на запись (SELECT + UPDATE RETURNING на прогон, единицы мс) — по построению успевает; компоуз не менял.

Тесты tests/test_3391_drain_marks_inflight_app_tasks.py: настоящие SchedulerContext/_dispatch/drain_inflight; застрявшая корутина (sleep(10)), соседняя финализируется сама; фейковый mark_done моделирует app-копию (замена + гейт running) — голая метка стёрла бы чекпоинт. Плюс caplog-тест на отсутствие «drained cleanly» при hard-cancel.

Фальсификация (git apply -R исходников, тесты на месте): assert 'running' == 'done' ×2 и 'drained cleanly' not in … — 3 failed, rc=1; после возврата — зелёное.

Прогоны: полный backend 5521 passed, 35 skipped (rc=0); ruff OK.

Приёмка на проде: следующий деплой при идущем detail-бэкфилле — в логах scraper scheduler: drain — N прогон(ов) сняты с 'running' как interrupted: […] и нет «drained cleanly» после «hard-cancelling»; scrape_runs: status='done' (или failed по honest-status-гейту при доле отказов ≥15 %), interrupted='1', boot_reaped пуст, прежние счётчики на месте; новый контейнер: boot-reap — … нет (0).

Closes #3391. **Дефект (прод 07.09 02:36 UTC, первый реальный SIGTERM-drain после #3363).** Деплой пересоздал scraper при идущих `cian_detail_backfill` 6167 (75 мин, 162 объявления) и `cian_history_backfill` 6173; через 100 с grace — hard-cancel; строки остались `running`, boot-reap нового контейнера сделал их `zombie` без `interrupted`. `interrupted=1` на SIGTERM писали только kit-пайплайны и DKP; app-task бэкфиллы — нет. Плюс `scheduler_main` печатал «scheduler drained cleanly (SIGTERM)» и после hard-cancel. **Фикс в одном месте (kit-scheduler).** `orchestration/scheduler.py:354-380` — реестр `_inflight_run_ids: dict[Task,int]` через `spawn_tracked(coro, run_id=…)` из `_dispatch` (`:975`; claim не тронут), done-callback чистит. `mark_inflight_interrupted(tasks)` (`:380`, синхронная — без `await`, живёт внутри `except CancelledError`): для каждого run_id `SELECT status, counters`; если `running` → `counters["interrupted"]=1` **поверх прочитанных** и `mark_done`. Вызывается и по истечении grace (`:461`), и в `except asyncio.CancelledError` (`:449`, затем `raise`) — на проде hard-cancel прилетает внутрь `asyncio.wait` (хвост tick-сна до 60 с + 80 с дрейна > 100 с grace). Counters читаются из строки, потому что боевой `ctx.runs` — `app.services.scrape_runs`, где `mark_done` counters **заменяет** (#3390): голый `{"interrupted":1}` стёр бы `done_buckets`. «Не трогать успевших» — проверка статуса + собственный гейт `mark_done` `WHERE status='running'`. `app/scheduler_main.py:156,237-244` — `_await_scheduler` возвращает признак hard-cancel; «drained cleanly» печатается только когда задача вышла сама. **`stop_grace_period`:** 120 с у scraper против `_DRAIN_TIMEOUT_S=100` → 20 с на запись (SELECT + UPDATE RETURNING на прогон, единицы мс) — по построению успевает; компоуз не менял. **Тесты** `tests/test_3391_drain_marks_inflight_app_tasks.py`: настоящие `SchedulerContext`/`_dispatch`/`drain_inflight`; застрявшая корутина (`sleep(10)`), соседняя финализируется сама; фейковый `mark_done` моделирует app-копию (замена + гейт `running`) — голая метка стёрла бы чекпоинт. Плюс caplog-тест на отсутствие «drained cleanly» при hard-cancel. **Фальсификация** (`git apply -R` исходников, тесты на месте): `assert 'running' == 'done'` ×2 и `'drained cleanly' not in …` — 3 failed, rc=1; после возврата — зелёное. **Прогоны:** полный backend `5521 passed, 35 skipped` (rc=0); ruff OK. **Приёмка на проде:** следующий деплой при идущем detail-бэкфилле — в логах scraper `scheduler: drain — N прогон(ов) сняты с 'running' как interrupted: […]` и нет «drained cleanly» после «hard-cancelling»; `scrape_runs`: `status='done'` (или `failed` по honest-status-гейту при доле отказов ≥15 %), `interrupted='1'`, `boot_reaped` пуст, прежние счётчики на месте; новый контейнер: `boot-reap — … нет (0)`.
bot-backend added 1 commit 2026-09-06 02:59:57 +00:00
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
30e3bacc5e
Прод 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'ов).
Author
Collaborator

Deep-ревью (07.09): approve с предупреждениями. Подтверждено по коду: пометка синхронная, своя сессия (RealSessionFactory), между SELECT и mark_done нет точек передачи управления; гейт WHERE status='running' есть в обеих копиях mark_done/mark_failed; резюм для failed+interrupted работает ('failed' в _RESUME_STATUSES); после hard-cancel новых await нет, 20 с до SIGKILL хватает; обе мутации краснеют ('running' == 'done', KeyError: 'done_buckets').

Докатываю до мержа: (1) update_heartbeat в обеих копиях без гейта по статусу — живая задача следующим пульсом стирает interrupted (app-копия заменяет counters) → AND status='running'; (2) hard-cancel в теле тика (не внутри drain_inflight) метку не ставит — except CancelledError вокруг тик-лупа; (3) db.rollback() в except — иначе первый отказ SQL утащит остальные run_id; (4) session_factory()/close() внутри try, чтобы не подменить CancelledError; (5) честный список помеченных в WARNING + потолок «БД недоступна → SIGKILL через 20 с» в докстринге.

Отдельно — #3393: оборванные деплоем прогоны теперь попадают в лестницы стриков (failed-ratio по частичным counters; done обнуляет стрик банов источника) — до #3392 такие строки были zombie и в выборки не попадали.

Deep-ревью (07.09): approve с предупреждениями. Подтверждено по коду: пометка синхронная, своя сессия (`RealSessionFactory`), между SELECT и `mark_done` нет точек передачи управления; гейт `WHERE status='running'` есть в обеих копиях `mark_done`/`mark_failed`; резюм для `failed+interrupted` работает (`'failed'` в `_RESUME_STATUSES`); после hard-cancel новых `await` нет, 20 с до SIGKILL хватает; обе мутации краснеют (`'running' == 'done'`, `KeyError: 'done_buckets'`). Докатываю до мержа: (1) `update_heartbeat` в обеих копиях без гейта по статусу — живая задача следующим пульсом стирает `interrupted` (app-копия заменяет counters) → `AND status='running'`; (2) hard-cancel в теле тика (не внутри `drain_inflight`) метку не ставит — `except CancelledError` вокруг тик-лупа; (3) `db.rollback()` в `except` — иначе первый отказ SQL утащит остальные run_id; (4) `session_factory()/close()` внутри `try`, чтобы не подменить `CancelledError`; (5) честный список помеченных в WARNING + потолок «БД недоступна → SIGKILL через 20 с» в докстринге. Отдельно — #3393: оборванные деплоем прогоны теперь попадают в лестницы стриков (failed-ratio по частичным counters; `done` обнуляет стрик банов источника) — до #3392 такие строки были `zombie` и в выборки не попадали.
Light1YT added 1 commit 2026-09-06 04:05:17 +00:00
fix(#3391): пульс не пишет по финализированной строке; отмена в теле тика тоже помечает; rollback/try/честный лог
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 9s
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 4m56s
a7362bc5fa
Пять замечаний deep-ревью к PR #3392, ровно они.

1. Пульс стирал метку дрейна. `update_heartbeat` обеих копий бил `WHERE id = :run_id`
   без гейта по статусу, а app-копия counters ЗАМЕНЯЕТ (#3390): задача, помеченная
   `interrupted`, но ещё живая (ветка таймаута drain_inflight отдаёт её внешнему
   hard-cancel'у — несколько итераций спустя, пульс на каждый батч —
   app/services/scheduler.py:141), следующим же ударом стирала метку, и оборванный
   прогон снова читался как полный проход. Гейт — `IN ('running', 'cancelled')`, а не
   `= 'running'`: 'cancelled' финализирует строку, но задача встаёт лишь на ближайшей
   границе якоря, и её последний пульс — ЕДИНСТВЕННЫЙ писатель чекпоинта в этот момент
   (pipeline.py:1308/2368/2986/4488, mark_done там уже no-op по своему гейту), а
   'cancelled' входит в _RESUME_STATUSES — сужение до 'running' молча съело бы точку
   возобновления у каждой отмены. Возвращаемое значение update_heartbeat не читает
   никто (обе копии -> None, ни одного присваивания на 130 сайтах вызова), так что
   «0 строк обновлено» ломать нечего; no-op логируется WARNING'ом, как у mark_done.

2. Отмена вне drain_inflight. Hard-cancel приходит по расписанию grace'а
   scheduler_main, а не по нашему, и может застать ТЕЛО тика (reap / stale-digest /
   `_dispatch` с сетевым pre_claim). `except Exception` тика CancelledError не ловит,
   до `await ctx.drain_inflight()` дело не доходит — строки оставались 'running'.
   Тело вынесено в `_tick_loop`, `scheduler_loop` ловит CancelledError, помечает
   in-flight и пробрасывает отмену.

3. rollback в except пометки: отказавший statement оставляет сессию в aborted-tx, и
   первый же непроходимый run_id утаскивал все следующие (образец — defensive rollback
   в mark_failed/mark_banned).

4. session_factory()/db.close() втянуты в try: исключение оттуда ЗАМЕНИЛО бы собой
   CancelledError, а suppress(CancelledError) в scheduler_main его не глушит — процесс
   уходил бы с трейсбеком вместо чистого drain-выхода.

5. WARNING перечисляет marked_ids, а не весь run_ids (там были и пропущенные по
   статусу). В докстринге назван потолок: SELECT синхронный, у движка нет ни connect-,
   ни statement-таймаута (app/core/db.py:8-19) — недоступная БД блокирует луп до
   SIGKILL'а через 20 с docker-grace; данные при этом не хуже прежних (строки остаются
   'running' → boot-reap).

Тесты — по значению, не по факту вызова; на исходниках 30e3bacc краснеют все пять:
'listings_processed' дописан в финализированную строку (kit), KeyError: 'interrupted'
(app), assert 'running' == 'done' (отмена в теле тика), assert 0 == 1 (второй run_id
не помечен после отказа первого), RuntimeError наружу (недоступная БД).
bot-backend merged commit aeb1b2b3b9 into main 2026-09-06 04:14:07 +00:00
Sign in to join this conversation.
No reviewers
No milestone
No project
No assignees
2 participants
Notifications
Due date
The due date is invalid or out of range. Please use the format "yyyy-mm-dd".

No due date set.

Dependencies

No dependencies set.

Reference: lekss361/gendesign#3392
No description provided.