From ae6ad9e440a8b16c85492437e168d556671e9932 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 6 Sep 2026 12:34:53 +0500 Subject: [PATCH] =?UTF-8?q?docs(#3390):=20=D1=83=D0=B1=D1=80=D0=B0=D1=82?= =?UTF-8?q?=D1=8C=20=D0=BB=D0=BE=D0=B6=D0=BD=D1=8B=D0=B5=20=C2=ABcounters?= =?UTF-8?q?=20=D0=97=D0=90=D0=9C=D0=95=D0=9D=D0=AF=D0=95=D0=A2=C2=BB=20?= =?UTF-8?q?=D0=B8=D0=B7=20=D0=BA=D0=BE=D0=BC=D0=BC=D0=B5=D0=BD=D1=82=D0=B0?= =?UTF-8?q?=D1=80=D0=B8=D0=B5=D0=B2/=D0=B4=D0=BE=D0=BA=D1=81=D1=82=D1=80?= =?UTF-8?q?=D0=B8=D0=BD=D0=B3=D0=BE=D0=B2;=20=D0=B3=D0=B5=D0=B9=D1=82=20te?= =?UTF-8?q?st=5F3168=20=E2=80=94=20=D1=87=D0=B5=D1=81=D1=82=D0=BD=D0=B0?= =?UTF-8?q?=D1=8F=20=D1=84=D0=BE=D1=80=D0=BC=D1=83=D0=BB=D0=B8=D1=80=D0=BE?= =?UTF-8?q?=D0=B2=D0=BA=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tradein-mvp/backend/app/services/scheduler.py | 14 ++++------ .../tests/test_3168_backfill_cursor_resume.py | 28 ++++++++----------- .../tests/test_3355_drain_mark_full_loads.py | 8 +++--- .../test_3384_no_proxy_at_batch_start.py | 13 +++++---- ...est_3391_drain_marks_inflight_app_tasks.py | 25 +++++++++-------- 5 files changed, 41 insertions(+), 47 deletions(-) diff --git a/tradein-mvp/backend/app/services/scheduler.py b/tradein-mvp/backend/app/services/scheduler.py index 79db37b5..0f67306d 100644 --- a/tradein-mvp/backend/app/services/scheduler.py +++ b/tradein-mvp/backend/app/services/scheduler.py @@ -32,8 +32,8 @@ from typing import Any # его алиас, реализация одна и counters везде МЕРЖАТСЯ (`counters || :counters`). До # #3390 копии было две, и app-копия counters ЗАМЕНЯЛА — тогда чекпоинт курсора # import_rosreestr_dkp (#3168) обязан был писаться именно kit-именем, иначе resume-вердикт -# со старта затирался первым же per-batch пульсом. Имя оставлено: на него смотрит -# текстовый гейт test_3168 (test_cursor_write_uses_merge_not_replace_heartbeat). +# со старта затирался первым же per-batch пульсом. Имя оставлено как есть: теперь это +# один объект, и переименование в runs_mod ничего не чинит и ничего не ломает. from scraper_kit.orchestration import runs as kit_runs # compute_next_run_at жил здесь ВТОРОЙ, побайтово одинаковой копией kit-версии (#2674). @@ -131,12 +131,10 @@ async def _execute_cian_backfill( поведению (пометка 'zombie' на 6-м часу).""" nonlocal counters # #3384: снимок измеренного едет не только в БД, но и в `counters` — этот словарь - # уезжает в mark_failed из общего except ниже, а тот counters ЗАМЕНЯЕТ - # (scrape_runs.py:738 `counters = CAST(:counters AS jsonb)`; мерж `||` — только у - # kit-копии, которую этот путь не зовёт). Пока снимок сюда не доезжал, любой отказ - # ПОСЛЕ пройденной стадии (пул опустел между стадиями, упал SELECT домов) писал - # поверх измеренного предынициализированные нули, и SQL-разбор простоя - # (#3288/#3367) читал «к площадке не ходили» про прогон, который ходил. + # уезжает в mark_failed из общего except ниже. Мерж (#3390) спасает лишь ключи, + # которых в payload нет; одноимённые он ПЕРЕЗАПИСЫВАЕТ, поэтому предынициализированные + # нули без этого присваивания легли бы поверх измеренного, и SQL-разбор простоя + # (#3288/#3367) прочитал бы «к площадке не ходили» про прогон, который ходил. # Присваивание ДО записи в БД: сбой heartbeat'а не должен стирать сам факт замера. counters = _counters(progress) try: diff --git a/tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py b/tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py index 25e7f499..eedd35ba 100644 --- a/tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py +++ b/tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py @@ -23,10 +23,10 @@ import_rosreestr_dkp (source='rosreestr_dkp_import'), шестой backfill, н штатным поведением). - Потолок возраста чекпоинта — 24ч (_DKP_CHECKPOINT_STALE_HOURS): старше — курсор считается негодным, прогон стартует с last_id=0, причина 'checkpoint_stale'. - - Вердикт пишется в scrape_runs.counters через kit_runs.update_heartbeat - (`counters || :counters` — merge), а не локальный runs_mod.update_heartbeat - (`CAST(:counters AS jsonb)` — полная замена): иначе первый же per-batch heartbeat - после старта стирает resume-вердикт. + - Вердикт пишется в scrape_runs.counters через update_heartbeat, а тот counters + МЕРЖИТ (`counters || :counters`) — иначе первый же per-batch heartbeat после старта + стёр бы resume-вердикт. На момент #3168 мерж был только у kit-копии, поэтому запись + шла именно kit-именем; с #3390 копия одна и семантика мержа — единственная. Обратимость (см. PR summary): временный откат last_id на литерал 0 красит test_resume_continues_from_saved_last_id (last_id == 123456 не совпадает с 0). @@ -34,7 +34,6 @@ test_resume_continues_from_saved_last_id (last_id == 123456 не совпада from __future__ import annotations -import inspect import json import os from typing import Any @@ -163,20 +162,15 @@ def test_checkpoint_just_under_ceiling_is_still_accepted() -> None: # ── 3. Запись курсора не затирает посторонние ключи в counters (мерж, не замена) ──── +# +# Текстового гейта на имя `kit_runs.update_heartbeat` здесь больше нет (#3390): пока копий +# было две, имя выбирало семантику, и проверять его по тексту имело смысл. Теперь +# `app.services.scrape_runs` — алиас kit'а, оба имени дают ОДИН объект, и гейт краснел бы +# на переименовании, ничего при этом не защищая. Мерж проверяется по значению — тестом +# ниже и test_3390_single_runs_module.py (оба пути импорта, heartbeat + финализатор). -def test_cursor_write_uses_merge_not_replace_heartbeat() -> None: - """import_rosreestr_dkp обязан писать чекпоинт через kit_runs.update_heartbeat - (merge: `counters || :counters`), а не локальный runs_mod.update_heartbeat (замена: - `CAST(:counters AS jsonb)`) — иначе resume-вердикт, записанный ДО цикла, стирается - первым же per-batch heartbeat'ом того же прогона. - """ - src = inspect.getsource(sched.import_rosreestr_dkp) - assert "kit_runs.update_heartbeat" in src - assert "runs_mod.update_heartbeat" not in src - - -def test_kit_runs_update_heartbeat_merges_into_existing_counters() -> None: +def test_update_heartbeat_merges_into_existing_counters() -> None: """Сама примитива слияния: второй write добавляет rows_fetched/last_id, НЕ стирая resume_from/resume_reason, записанные первым write'ом. diff --git a/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py b/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py index a5aaa286..5bc7bbde 100644 --- a/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py +++ b/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py @@ -247,10 +247,10 @@ async def test_domclick_sweep_drain_is_marked_interrupted() -> None: async def test_domclick_sweep_drain_keeps_inherited_checkpoint() -> None: """domclick: дрейн ОБЯЗАН переписать унаследованный чекпоинт в свой payload. - Мержа jsonb тут нет: app-level `scrape_runs.update_heartbeat`/`mark_done` пишут - `counters = CAST(:counters AS jsonb)` — полная перезапись, а scheduler при claim - done_buckets не наследует. Без явного ключа дрейн-прогон закрывался бы с пустым - чекпоинтом, и резюм пересобирал бы все шесть корзин заново. + Мерж jsonb (#3390) тут не спасает: он сливает payload со СВОЕЙ строкой прогона, а та + создана пустой (`create_run`) — унаследованный чекпоинт лежит в counters ПРЕДЫДУЩЕГО + прогона, и scheduler при claim его не переносит. Без явного ключа дрейн-прогон + закрывался бы с пустым чекпоинтом, и резюм пересобирал бы все шесть корзин заново. """ from scraper_kit.orchestration import pipeline as pl diff --git a/tradein-mvp/backend/tests/test_3384_no_proxy_at_batch_start.py b/tradein-mvp/backend/tests/test_3384_no_proxy_at_batch_start.py index c0e36231..5d2b9ea3 100644 --- a/tradein-mvp/backend/tests/test_3384_no_proxy_at_batch_start.py +++ b/tradein-mvp/backend/tests/test_3384_no_proxy_at_batch_start.py @@ -16,8 +16,9 @@ break`), которая и пишет `counters.no_proxy_stop`. Общий `exce `httpx.AsyncClient.post` — «к площадке не ходили» проверяется, а не предполагается. Третий тест — про соседний случай: пул опустел МЕЖДУ стадиями, то есть уже ПОСЛЕ -реальной работы. Там проверяется, что в запись прогона не уезжают нули поверх -измеренного (`mark_failed` у `app.services.scrape_runs` counters ЗАМЕНЯЕТ, а не мержит). +реальной работы. Там проверяется, что в запись прогона не уезжают нули поверх измеренного: +`mark_failed` counters МЕРЖИТ (#3390), но одноимённые ключи при мерже перезаписываются, +поэтому нули в payload'е финализатора всё так же стирают замер пульса. Сеть/БД замоканы; в сеть тест не ходит. """ @@ -199,10 +200,10 @@ async def test_cian_empty_pool_between_stages_keeps_measured_counters() -> None: У циана (в отличие от avito/домклика с их живым `counters.to_dict()`) в общий `except` приходит СТАРЫЙ словарь: реальные значения присваиваются уже после возврата из `backfill_cian_history`, а отказ бывает и посреди неё — пул опустел между - стадиями, упал SELECT домов. `runs_mod` здесь настоящий - (`app.services.scrape_runs`), и его `mark_failed` counters ЗАМЕНЯЕТ - (`counters = CAST(:counters AS jsonb)`, scrape_runs.py:738) — то есть в записи - прогона остаётся ровно то, что уехало последним аргументом. + стадиями, упал SELECT домов. `runs_mod` здесь настоящий (`app.services.scrape_runs`, + с #3390 — алиас kit'а), и его `mark_failed` counters МЕРЖИТ; мерж спасает только + ключи, которых в payload'е нет, а одноимённые ПЕРЕЗАПИСЫВАЕТ — нули поверх замера + всё равно недопустимы, и проверять их надо у вызывающего. Проверка по значению: смотрим jsonb-payload обоих UPDATE'ов (heartbeat и mark_failed), а не факт вызова. 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 index 631c3043..daf4d42f 100644 --- 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 @@ -15,15 +15,15 @@ cian_history_backfill 6173) остались в scrape_runs со статусо 3. hard-cancel из scheduler_main (CancelledError прямо в `asyncio.wait` дрейна) даёт тот же результат — это ровно прод-путь 07.09, ветка таймаута его не покрывает; 4. `scheduler_main` после hard-cancel'а больше не печатает «drained cleanly»; - 5. пульс НЕ пишет по финализированной строке (обе копии `update_heartbeat`) — иначе - помеченная, но ещё живая задача стирает метку следующим же ударом; + 5. пульс НЕ пишет по финализированной строке (оба пути импорта `update_heartbeat`) — + иначе помеченная, но ещё живая задача стирает метку следующим же ударом; 6. отмена, прилетевшая в тело тика (не в дрейн), помечает in-flight так же; 7. отказ SQL на одном run_id не уносит остальные, а WARNING перечисляет ровно помеченных. -Фейковый `mark_done` counters ЗАМЕНЯЕТ, а не мержит — так ведёт себя боевая app-копия -(`app.services.scrape_runs`, #3390), которая и инжектируется в прод-контекст. Голый -`{"interrupted": 1}` на этом фейке стёр бы чекпоинт, и (1) покраснел бы. +Фейковый `mark_done` counters ЗАМЕНЯЕТ, а не мержит — намеренно строже боевого (тот с +#3390 мержит везде): голый `{"interrupted": 1}` на этом фейке стёр бы чекпоинт, и (1) +покраснел бы. Дрейн обязан слать чекпоинт явно, а не полагаться на мерж в БД. """ from __future__ import annotations @@ -49,8 +49,8 @@ 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-модуля (с #3390 это ОДИН объект: app.services.scrape_runs — +# алиас kit'а): гейт по статусу проверяется через каждый из них одним телом теста. _RUNS_MODULES = {"kit": kit_runs, "app": app_runs} @@ -131,7 +131,7 @@ class _FakeRuns: 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) + row["counters"] = dict(counters) # строже боевого мержа (#3390) — см. докстринг модуля def _make_sched(source: str) -> dict[str, Any]: @@ -283,7 +283,7 @@ async def test_hard_cancel_does_not_log_drained_cleanly( assert "drained and exited cleanly" not in caplog.text -# ── 5. Пульс не пишет по финализированной строке (обе копии update_heartbeat) ──── +# ── 5. Пульс не пишет по финализированной строке (оба пути импорта) ───────────── class _RunRowDb: @@ -293,8 +293,8 @@ class _RunRowDb: Допустимые статусы вычитываются из SQL (`status = 'x'` / `status IN ('x', 'y')`), а не зашиты ожиданием теста: на коде без гейта (WHERE только по id) апдейт проходит по строке ЛЮБОГО статуса — и тест краснеет по значению, а не по отсутствию подстроки. - Мерж jsonb (`||`) против замены (`CAST(:counters AS jsonb)`) — тоже по тексту: у - kit- и app-копии он разный, а гейт нужен обеим. + Мерж jsonb (`||`) против замены — тоже по тексту statement'а, а не по ожиданию: + вернётся замена (было до #3390) — двойник это отразит, и тест покраснеет по значению. SELECT'ы alert-хуков (`_alert_on_run_id`) не моделируются — они best-effort и возвращают пусто. @@ -348,7 +348,8 @@ def test_heartbeat_does_not_erase_drain_mark_of_finalized_run(name: str) -> None Ровно эта последовательность достижима на проде: ветка таймаута `drain_inflight` оставляет задачу внешнему hard-cancel'у, и до него она успевает несколько итераций (`app/services/scheduler.py:141` — пульс на каждый батч). Без гейта по статусу - app-копия ЗАМЕНЯЛА counters и стирала метку: оборванный прогон снова читался как + пульс лёг бы поверх метки (одноимённые ключи мерж перезаписывает, а до #3390 + app-копия и вовсе ЗАМЕНЯЛА весь словарь): оборванный прогон снова читался бы как полный проход, а резюм его не подхватывал. """ mod = _RUNS_MODULES[name]