fix(tradein/avito): серия отказов обрывается и называет причину, а не выедает бюджет (#2674) #2739

Merged
bot-backend merged 1 commit from fix/avito-backfill-failure-breaker into main 2026-08-06 15:29:49 +00:00
3 changed files with 175 additions and 4 deletions

View file

@ -578,6 +578,7 @@ def mark_backfill_finished(
*,
source: str,
aborted_by_blocks: bool = False,
fail_hint: str | None = None,
) -> None:
"""Честный финал detail-backfill'а (#2674): нулевой прогон ≠ 'done'.
@ -601,11 +602,19 @@ def mark_backfill_finished(
`gone` (404 у avito) считается результатом наравне с `enriched`: прогон,
который подтвердил снятие объявлений, работу сделал.
`fail_hint` самая частая причина отказа этого прогона (задача считает её сама,
см. avito_detail_backfill._failure_signature). Дописывается в текст статуса,
потому что «blocked=5, обогащено 0» не отвечает на единственный вопрос, ради
которого статус и читают: отказала площадка или наш тракт (#2686, #2698). Логи
контейнера на этот вопрос отвечать не могут они исчезают при пересоздании
контейнера, то есть на первом же деплое после ночного прогона.
"""
attempted = int(counters.get("attempted") or 0)
enriched = int(counters.get("enriched") or 0)
blocked = int(counters.get("blocked") or 0)
produced = enriched + int(counters.get("gone") or 0)
hint = f"; причина: {fail_hint}" if fail_hint else ""
if attempted == 0:
mark_done(db, run_id, counters)
@ -614,7 +623,7 @@ def mark_backfill_finished(
if blocked and (aborted_by_blocks or produced == 0):
reason = (
f"backfill-honest-status: {source} остановлен блоками источника — "
f"blocked={blocked}, обогащено {enriched} из {attempted} попыток (#2674)"
f"blocked={blocked}, обогащено {enriched} из {attempted} попыток{hint} (#2674)"
)
logger.error("%s run_id=%d", reason, run_id)
mark_banned(db, run_id, reason, counters)
@ -624,7 +633,7 @@ def mark_backfill_finished(
reason = (
f"backfill-honest-status: {source} без результата — 0 обогащено из "
f"{attempted} попыток (failed={counters.get('failed', 0)}, "
f"blocked={blocked}) (#2674)"
f"blocked={blocked}){hint} (#2674)"
)
logger.error("%s run_id=%d", reason, run_id)
mark_failed(db, run_id, reason, counters)

View file

@ -13,6 +13,13 @@ session path as the detail-phase of `run_avito_city_sweep`
rotate IP on every block, abort after max_consecutive_blocks. Статус оборванного
блоками прогона 'banned' (#2674, runs.mark_backfill_finished): работу он не
доделал, остаток снапшота уедет в следующую ночь через NULL detail_enriched_at.
Отказы, не являющиеся блоками, до 2026-08-06 брейкера не имели вовсе: прогоны
3-5 августа делали ~1600 попыток, получали 1600 отказов, ноль обогащений и
выедали весь бюджет (9000 с) вместе с 1600 запросами через единственный прокси.
Теперь такая серия обрывается по max_consecutive_failures, а самая частая причина
отказа пишется в текст статуса прогона (_failure_signature) иначе она живёт
только в логах контейнера, а те исчезают на первом же деплое.
"""
from __future__ import annotations
@ -20,7 +27,9 @@ from __future__ import annotations
import asyncio
import logging
import random
import re
import time
from collections import Counter
from dataclasses import dataclass, field
from urllib.parse import urlparse
@ -99,6 +108,38 @@ _OBLAST_AVITO_URL_PATTERNS = tuple(
)
# Причина отказа карточки без её URL: 1576 отказов одного прогона должны схлопнуться
# в ОДНУ строку, иначе перепись бесполезна.
_URL_IN_MESSAGE_RE = re.compile(r"https?://\S+")
def _failure_signature(exc: BaseException) -> str:
"""Подпись причины отказа: тип исключения + текст без URL.
Зачем (замер 2026-08-06): у прогонов 3 и 4 августа counters говорили
`attempted=1576, failed=1576, blocked=0` и ничего больше. Кто отказал,
площадка или наш тракт, было видно ТОЛЬКО в логах контейнера, а тот
пересоздаётся на каждом деплое и уносит их с собой; в GlitchTip попадают
события уровня ERROR, а поштучные отказы WARNING. Разница между этими
двумя диагнозами разные владельцы задачи (#2686, #2698), поэтому она
обязана переживать перезапуск контейнера, то есть лежать в самом прогоне.
Тип исключения первый разряд диагноза (AvitoBlockedError = площадка
показала 403/firewall; сетевой класс curl_cffi = наш прокси-тракт;
ValueError = ответ пришёл, но не разобран), текст второй.
"""
message = _URL_IN_MESSAGE_RE.sub("<url>", str(exc)).strip()
return f"{type(exc).__name__}: {message}"[:160] if message else type(exc).__name__
def _top_failure(census: Counter[str]) -> str | None:
"""Самая частая причина отказа с её долей; None — отказов не было."""
if not census:
return None
reason, hits = census.most_common(1)[0]
return f"{reason} ({hits} из {sum(census.values())})"
@dataclass
class AvitoDetailBackfillResult:
"""Counters for one backfill run."""
@ -139,6 +180,8 @@ async def run_avito_detail_backfill(
budget_sec: float -- wall-clock budget per run, default 3600s.
request_delay_sec: float -- delay between listings, default 6.0s.
max_consecutive_blocks: int -- abort threshold, default 5.
max_consecutive_failures: int -- порог обрыва по отказам-не-блокам,
default 25 (см. комментарий у чтения параметра ниже).
Lifecycle: update_heartbeat -> snapshot -> loop with budget guard ->
mark_backfill_finished (done / banned при блоках / failed при нуле, #2674);
@ -149,6 +192,13 @@ async def run_avito_detail_backfill(
budget_sec = float(params.get("budget_sec", 3600))
request_delay_sec = float(params.get("request_delay_sec", 6.0))
max_consecutive_blocks = int(params.get("max_consecutive_blocks", 5))
# Брейкер на отказы-НЕ-блоки. Блоки свой брейкер имели с самого начала, отказы —
# нет, и это стоило трёх ночей подряд: 3-5 августа прогон делал ~1600 попыток,
# получал 1600 отказов, ноль обогащений и выедал весь бюджет 9000 с (плюс 1600
# запросов через единственный прокси, #2638). Порог заметно выше блочного: пачка
# мёртвых карточек (404 → ValueError в curl-режиме) не должна обрывать здоровый
# прогон, а 25 отказов подряд без единого успеха — уже не невезение.
max_consecutive_failures = int(params.get("max_consecutive_failures", 25))
warm_batch = int(params.get("warm_batch", 500))
research_every = int(params.get("research_every", 50))
block_cooldown_sec = float(params.get("block_cooldown_sec", 30.0))
@ -308,9 +358,13 @@ async def run_avito_detail_backfill(
)
consecutive_blocks = 0
consecutive_failures = 0
aborted_by_blocks = False
do_sleep = False
items_since_warm = 0
# Перепись причин (блоки + отказы) — переживает пересоздание контейнера,
# в отличие от логов; см. _failure_signature.
failure_census: Counter[str] = Counter()
for idx, row in enumerate(snapshot):
# Budget guard
@ -429,8 +483,9 @@ async def run_avito_detail_backfill(
if use_curl:
items_since_warm += 1
consecutive_blocks = 0
consecutive_failures = 0
except AvitoListingGoneError:
except AvitoListingGoneError as gone_exc:
# #2034: мёртвый листинг (404 / removed) — НЕ блок, НЕ failed.
# Координатные дыры в lat-null очереди в основном dead-листинги;
# browser-mode рендерит их «Ошибка 404» без item-view → раньше это
@ -440,6 +495,10 @@ async def run_avito_detail_backfill(
# и не сбрасываем). Метим is_active=FALSE → листинг уходит из scope
# (snapshot SELECT фильтрует is_active = TRUE) и не тратит фетчи впредь.
counters.gone += 1
# 404 — честный ответ площадки, значит тракт цел: серия отказов
# прерывается (блочный брейкер 404 не трогает — см. #2034).
consecutive_failures = 0
failure_census[_failure_signature(gone_exc)] += 1
try:
with db.begin_nested():
db.execute(
@ -481,6 +540,7 @@ async def run_avito_detail_backfill(
except (AvitoBlockedError, AvitoRateLimitedError) as e:
consecutive_blocks += 1
counters.blocked += 1
failure_census[_failure_signature(e)] += 1
do_sleep = False
logger.warning(
"avito_detail_backfill: run_id=%d BLOCKED #%d/%d (consecutive=%d): %s",
@ -538,12 +598,14 @@ async def run_avito_detail_backfill(
exc_info=True,
)
except TimeoutError:
except TimeoutError as e:
# asyncio.wait_for → TimeoutError (py3.12: asyncio.TimeoutError — alias).
# Ловим ДО общего Exception (TimeoutError ⊂ OSError ⊂ Exception). Зависший
# fetch отменён → листинг failed, переходим к следующему (loop не зависает,
# run не zombie #1950). Не считаем soft-блоком: rotate не дёргаем.
counters.failed += 1
consecutive_failures += 1
failure_census[_failure_signature(e)] += 1
logger.warning(
"avito_detail_backfill: run_id=%d listing %s TIMEOUT (>%.0fs) -- skip",
run_id,
@ -557,6 +619,8 @@ async def run_avito_detail_backfill(
except Exception as e:
counters.failed += 1
consecutive_failures += 1
failure_census[_failure_signature(e)] += 1
logger.warning(
"avito_detail_backfill: run_id=%d listing %s failed: %s",
run_id,
@ -568,6 +632,18 @@ async def run_avito_detail_backfill(
except Exception:
pass
if consecutive_failures >= max_consecutive_failures:
logger.error(
"avito_detail_backfill: run_id=%d ABORT -- %d отказов подряд без "
"единого успеха, частая причина: %s. enriched=%d attempted=%d",
run_id,
consecutive_failures,
_top_failure(failure_census) or "неизвестна",
counters.enriched,
counters.attempted,
)
break
if counters.attempted % 25 == 0:
current_counters = counters.to_dict()
runs_mod.update_heartbeat(db, run_id, current_counters)
@ -580,6 +656,7 @@ async def run_avito_detail_backfill(
current_counters,
source="avito_detail_backfill",
aborted_by_blocks=aborted_by_blocks,
fail_hint=_top_failure(failure_census),
)
logger.info(
"avito_detail_backfill: run_id=%d FINISHED -- attempted=%d enriched=%d "

View file

@ -790,3 +790,88 @@ async def test_backfill_use_curl_block_cooldown_research_no_rebuild() -> None:
mock_scraper.return_value._rotate_ip.assert_not_called()
runs.mark_backfill_finished.assert_called_once()
runs.mark_failed.assert_not_called()
# ---------------------------------------------------------------------------
# Отказы-не-блоки: брейкер + перепись причин (2026-08-06)
# ---------------------------------------------------------------------------
@pytest.mark.asyncio
async def test_backfill_aborts_on_consecutive_failures_and_names_the_reason() -> None:
"""Серия отказов без единого успеха обрывается, а причина попадает в статус прогона.
Прод 3-5 августа: attempted1600, failed1600, blocked=0, enriched=0, весь
бюджет 9000 с и 1600 запросов через единственный прокси и ни слова о том,
ЧТО именно отказало (поштучные отказы логируются WARNING, а логи контейнера
пропадают на первом деплое). Брейкера на отказы-не-блоки не было вовсе.
"""
snapshot = _make_snapshot(200)
db = _mock_db(snapshot)
runs = MagicMock()
# Тот же класс отказа, что видели у соседнего свипа в тот же день.
mock_fetch = AsyncMock(
side_effect=OSError(
"Failed to perform, curl: (56) CONNECT tunnel failed, response 502. "
"See https://curl.se/libcurl/c/libcurl-errors.html"
)
)
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True)
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION, return_value=AsyncMock()),
patch(_SCRAPER),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SLEEP, new_callable=AsyncMock),
):
result = await run_avito_detail_backfill(
db,
run_id=77,
params={
"batch_size": 200,
"budget_sec": 3600,
"max_consecutive_failures": 25,
},
)
assert result.attempted == 25, "серия отказов обязана обрываться, а не выедать бюджет"
assert result.failed == 25
assert result.blocked == 0
hint = runs.mark_backfill_finished.call_args.kwargs["fail_hint"]
assert hint is not None
assert "OSError" in hint # тип исключения = кому принадлежит отказ
assert "CONNECT tunnel failed" in hint
assert "25 из 25" in hint # доля, а не единичный пример
assert "https://curl.se" not in hint # URL вырезан, иначе 1600 «разных» причин
@pytest.mark.asyncio
async def test_backfill_success_resets_failure_streak() -> None:
"""Успех между отказами обнуляет серию — здоровый прогон брейкер не трогает."""
snapshot = _make_snapshot(5)
db = _mock_db(snapshot)
runs = MagicMock()
boom = ValueError("avito detail HTTP 500 for https://www.avito.ru/x")
# 2 отказа, успех, 2 отказа — при пороге 3 ни одна серия его не достигает.
mock_fetch = AsyncMock(side_effect=[boom, boom, MagicMock(), boom, boom])
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True)
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION, return_value=AsyncMock()),
patch(_SCRAPER),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
):
result = await run_avito_detail_backfill(
db,
run_id=78,
params={"batch_size": 5, "budget_sec": 3600, "max_consecutive_failures": 3},
)
assert result.attempted == 5
assert result.enriched == 1
assert result.failed == 4