fix(domclick): честная мотивировка чекпоинта + один _payload() на все ветки
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
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 5m3s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
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 5m3s
Правка ревью к PR #3372: комментарии и текст коммита обещали восстановление потери, которой нет. Факт: pipeline.py импортирует scraper_kit.orchestration.runs, и все четыре его писателя МЕРЖАТ jsonb (`counters = COALESCE(counters,'{}') || CAST(:counters AS jsonb)`, runs.py:708/776/817/880), а _pick_resume наследует done_buckets при claim (scheduler.py:741-748, #3074) — чекпоинт domclick в БД не терялся. Перезаписывающий двойник app/services/scrape_runs.py этой функцией не вызывается. Поэтому done_buckets в каждом payload — единообразие и защита от смены писателя (если запись пойдёт через app-копию или kit-писатель станет перезаписывающим), а не спасение данных. Комментарии переписаны под этот факт. Дрейн-ветка и NoProxyAvailableError строили dict вручную — переведены на {**_payload(), "interrupted": 1} и _payload(), чтобы «один _payload() на функцию» было правдой. Тест: к cancel добавлены ассерты на финальный heartbeat и mark_done — оба payload'а несут унаследованное ∪ пройденное этим прогоном. Refs #3355, PR #3363. Closes #3369
This commit is contained in:
parent
b796702d9a
commit
140caa4ddf
2 changed files with 89 additions and 35 deletions
|
|
@ -1,15 +1,16 @@
|
|||
"""Cancel-ветка domclick-свипа обязана нести чекпоинт `done_buckets` (#3369).
|
||||
"""Каждая запись counters domclick-свипа несёт чекпоинт `done_buckets` (#3369).
|
||||
|
||||
Доводка #3355/PR #3363: дрейн- и бан-ветки `run_domclick_city_sweep` кладут
|
||||
`done_buckets` в свой payload явно, а cancel-ветка писала голый
|
||||
`counters.to_dict()`. Мержа jsonb тут нет: писатель — app-level
|
||||
`app/services/scrape_runs.py` (`counters = CAST(:counters AS jsonb)`, ПОЛНАЯ
|
||||
перезапись), а scheduler при claim `done_buckets` не наследует. `cancelled`
|
||||
входит в `_RESUME_STATUSES` → отменённый прогон закрывался с пустым чекпоинтом,
|
||||
и резюм пересобирал все шесть корзин заново.
|
||||
`done_buckets` в payload явно, а cancel-ветка, финальный heartbeat и финализаторы
|
||||
писали голый `counters.to_dict()`. Это НЕ было потерей данных: writer —
|
||||
`scraper_kit.orchestration.runs`, все его писатели мержат jsonb
|
||||
(`counters = COALESCE(counters,'{}') || CAST(:counters AS jsonb)`), а `_pick_resume`
|
||||
наследует ключ при claim (#3074). Инвариант ставится ради единообразия и
|
||||
независимости веток от того, какой писатель окажется на другом конце
|
||||
(перезаписывающий двойник `app/services/scrape_runs.py` существует).
|
||||
|
||||
Проверка по ЗНАЧЕНИЮ (не по факту записи): унаследованные две корзины обязаны
|
||||
дословно оказаться в payload'е отмены.
|
||||
Проверка по ЗНАЧЕНИЮ, а не по факту записи: чекпоинт обязан дословно лежать в
|
||||
payload'е отмены, финального heartbeat'а и `mark_done`.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -101,8 +102,8 @@ async def test_domclick_sweep_cancel_keeps_inherited_checkpoint() -> None:
|
|||
assert db.writes, "cancel не оставил ни одной записи counters"
|
||||
_sql, counters = db.writes[-1]
|
||||
assert counters.get("done_buckets") == inherited, (
|
||||
f"cancel: чекпоинт потерян ({counters.get('done_buckets')!r} вместо "
|
||||
f"{inherited!r}) — резюм пересоберёт уже собранные корзины"
|
||||
f"cancel: payload без чекпоинта ({counters.get('done_buckets')!r} вместо "
|
||||
f"{inherited!r}) — ветка зависит от того, мержит ли писатель jsonb"
|
||||
)
|
||||
|
||||
|
||||
|
|
@ -131,3 +132,59 @@ async def test_domclick_sweep_cancel_marker_is_not_resume_status_only() -> None:
|
|||
f"cancel без предшественника: ожидался пустой чекпоинт, получено "
|
||||
f"{counters.get('done_buckets')!r}"
|
||||
)
|
||||
|
||||
|
||||
class _CleanSweepScraper:
|
||||
"""Свип прошёл все корзины без блока: ветка честного статуса → mark_done."""
|
||||
|
||||
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
||||
self.blocked = False
|
||||
self.geo_filtered = 0
|
||||
self.fetch_errors = 0
|
||||
self.buckets_completed = 1
|
||||
self.buckets_total = 1
|
||||
self.completed_buckets = ["3:0:5000000"]
|
||||
self.bucket_start_index = 0
|
||||
self.request_delay_sec = 0.0
|
||||
|
||||
async def __aenter__(self) -> _CleanSweepScraper:
|
||||
return self
|
||||
|
||||
async def __aexit__(self, *_e: Any) -> None:
|
||||
return None
|
||||
|
||||
async def fetch_city(self, **_kw: Any) -> list[Any]:
|
||||
return []
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_domclick_sweep_final_heartbeat_and_mark_done_carry_checkpoint() -> None:
|
||||
"""Финальный heartbeat и mark_done несут унаследованное ∪ пройденное этим прогоном."""
|
||||
from scraper_kit.orchestration import pipeline as pl
|
||||
|
||||
inherited = ["1:0:5000000", "2:0:5000000"]
|
||||
expected = sorted([*inherited, "3:0:5000000"])
|
||||
db = _FakeDb(prev_counters={"done_buckets": inherited})
|
||||
with (
|
||||
patch.object(pl, "DomClickScraper", _CleanSweepScraper),
|
||||
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
||||
):
|
||||
await pl.run_domclick_city_sweep(
|
||||
db, # type: ignore[arg-type]
|
||||
run_id=3369,
|
||||
config=_config(),
|
||||
matcher=MagicMock(),
|
||||
request_delay_sec=0.0,
|
||||
shutdown_requested=lambda: False,
|
||||
resume_run_id=3368,
|
||||
)
|
||||
|
||||
assert len(db.writes) >= 2, f"ожидались heartbeat + финализатор, есть {len(db.writes)}"
|
||||
sql_done, counters_done = db.writes[-1]
|
||||
assert "finished_at" in sql_done, f"последняя запись — не финализатор: {sql_done!r}"
|
||||
_sql_hb, counters_hb = db.writes[-2]
|
||||
for label, payload in (("финальный heartbeat", counters_hb), ("mark_done", counters_done)):
|
||||
assert payload.get("done_buckets") == expected, (
|
||||
f"{label}: payload без чекпоинта ({payload.get('done_buckets')!r} вместо "
|
||||
f"{expected!r}) — ветка зависит от того, мержит ли писатель jsonb"
|
||||
)
|
||||
|
|
|
|||
|
|
@ -4426,11 +4426,16 @@ async def run_domclick_city_sweep(
|
|||
# Мутируемый контейнер для захвата scraper-ссылки из замыкания.
|
||||
_scraper_ref: list[DomClickScraper] = []
|
||||
|
||||
# #3369: чекпоинт (done_buckets) обязан ехать в КАЖДОЙ записи counters этой
|
||||
# функции. Писатель — app-level scrape_runs: `counters = CAST(:counters AS jsonb)`,
|
||||
# полная перезапись без jsonb-мержа, а scheduler при claim done_buckets не
|
||||
# наследует. Любая запись без ключа стирает чекпоинт, и резюм (cancelled/banned/
|
||||
# failed ∈ _RESUME_STATUSES) пересобирает уже собранные корзины заново.
|
||||
# #3369: чекпоинт (done_buckets) едет в КАЖДОЙ записи counters этой функции —
|
||||
# это единообразие и защита от смены писателя, а НЕ восстановление потери.
|
||||
# Как есть сегодня: пишет kit-овый scraper_kit.orchestration.runs (импорт выше),
|
||||
# и все четыре его писателя МЕРЖАТ jsonb — `counters = COALESCE(counters,'{}')
|
||||
# || CAST(:counters AS jsonb)`, так что ключ, записанный раньше, переживает
|
||||
# payload без него; плюс _pick_resume наследует done_buckets при claim (#3074).
|
||||
# То есть чекпоинт в БД не терялся. Перезаписывающий двойник существует —
|
||||
# app/services/scrape_runs.py (`counters = CAST(:counters AS jsonb)`), — но
|
||||
# этой функцией не вызывается. Полный payload делает ветки нечувствительными
|
||||
# к тому, какой из двух писателей окажется на другом конце.
|
||||
_checkpoint: list[str] = []
|
||||
|
||||
def _payload() -> dict[str, Any]:
|
||||
|
|
@ -4466,27 +4471,19 @@ async def run_domclick_city_sweep(
|
|||
# ── Cooperative cancel перед SERP-фазой ──────────────────────────────
|
||||
if runs.is_cancelled(db, run_id):
|
||||
logger.info("domclick-sweep run_id=%d: cancelled before SERP", run_id)
|
||||
# #3369: 'cancelled' ∈ _RESUME_STATUSES — payload без done_buckets стирал
|
||||
# унаследованный чекпоинт, и резюм пересобирал все корзины заново.
|
||||
# #3369: 'cancelled' ∈ _RESUME_STATUSES, прогон резюмируем — пишем
|
||||
# полный payload, как и все прочие ветки (см. _payload выше).
|
||||
runs.update_heartbeat(db, run_id, _payload())
|
||||
return counters
|
||||
elif shutdown_requested():
|
||||
# SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial),
|
||||
# минуя honest-status (это чистый drain, не QRATOR-блок).
|
||||
logger.info("domclick-sweep run_id=%d: SIGTERM-drain — stopping before SERP", run_id)
|
||||
# #3355: метка дрейна (см. city-свипы #3330/#3333) И явный перенос
|
||||
# чекпоинта. Никакого jsonb-мержа тут нет: domclick пишет counters
|
||||
# через app-level scrape_runs (update_heartbeat/mark_done делают
|
||||
# `counters = CAST(:counters AS jsonb)` — ПОЛНАЯ перезапись), а
|
||||
# scheduler при claim done_buckets не наследует. Без явного
|
||||
# done_buckets дрейн-прогон закрылся бы с пустым чекпоинтом и резюм
|
||||
# пересобирал бы все шесть корзин заново. Тот же приём, что в
|
||||
# ban-ветке ниже (except NoProxyAvailableError).
|
||||
_drain = {
|
||||
**counters.to_dict(),
|
||||
"done_buckets": sorted(skip_buckets),
|
||||
"interrupted": 1,
|
||||
}
|
||||
# #3355: метка дрейна (см. city-свипы #3330/#3333) поверх обычного
|
||||
# payload'а. Чекпоинт тут не «спасается»: kit-писатель мержит jsonb
|
||||
# (см. _payload выше) — done_buckets в payload'е нужен для
|
||||
# единообразия, а interrupted=1 отличает дрейн от честного done.
|
||||
_drain = {**_payload(), "interrupted": 1}
|
||||
runs.update_heartbeat(db, run_id, _drain)
|
||||
runs.mark_done(db, run_id, _drain)
|
||||
return counters
|
||||
|
|
@ -4571,16 +4568,16 @@ async def run_domclick_city_sweep(
|
|||
# хотя площадка вообще не была затронута. Тот же диагноз и тот же по духу
|
||||
# обработчик, что у run_avito_full_load/run_cian_full_load/run_yandex_full_load
|
||||
# (см. except NoProxyAvailableError там же) — ban_kind_of_exception() относит тип
|
||||
# исключения к BAN_KIND_INFRA («наша инфраструктура», не площадка). done_buckets
|
||||
# сохраняем как унаследованные skip_buckets — этот прогон новых не завершил, но
|
||||
# терять уже собранный чекпоинт (#2687-класс дефекта) resume не должен.
|
||||
# исключения к BAN_KIND_INFRA («наша инфраструктура», не площадка). В payload'е
|
||||
# done_buckets == унаследованные skip_buckets: этот прогон новых корзин не
|
||||
# завершил (_checkpoint обновляется только после SERP-фазы).
|
||||
logger.error("domclick-sweep run_id=%d: no proxy available — %s", run_id, exc)
|
||||
counters.errors_count += 1
|
||||
runs.mark_banned(
|
||||
db,
|
||||
run_id,
|
||||
f"domclick sweep aborted: {exc}",
|
||||
{**counters.to_dict(), "done_buckets": sorted(skip_buckets)},
|
||||
_payload(),
|
||||
ban_kind=ban_kind_of_exception(exc),
|
||||
)
|
||||
return counters
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue