Merge pull request 'Яндекс: пять вечно недогружаемых карточек больше не обрывают каждый прогон добора' (#3574) from fix/yandex-perpetual-underloaded into main
Some checks failed
Deploy Trade-In / test (push) Blocked by required conditions
Deploy Trade-In / build-backend (push) Blocked by required conditions
Deploy Trade-In / build-frontend (push) Blocked by required conditions
Deploy Trade-In / build-browser (push) Blocked by required conditions
Deploy Trade-In / deploy (push) Blocked by required conditions
Deploy Trade-In / perimeter-smoke (push) Blocked by required conditions
Deploy Trade-In / deploy-status (push) Blocked by required conditions
Deploy Trade-In / changes (push) Has been cancelled
Some checks failed
Deploy Trade-In / test (push) Blocked by required conditions
Deploy Trade-In / build-backend (push) Blocked by required conditions
Deploy Trade-In / build-frontend (push) Blocked by required conditions
Deploy Trade-In / build-browser (push) Blocked by required conditions
Deploy Trade-In / deploy (push) Blocked by required conditions
Deploy Trade-In / perimeter-smoke (push) Blocked by required conditions
Deploy Trade-In / deploy-status (push) Blocked by required conditions
Deploy Trade-In / changes (push) Has been cancelled
This commit is contained in:
commit
4c11a7dab5
4 changed files with 369 additions and 7 deletions
|
|
@ -52,6 +52,13 @@ non-200 ответ считается блоком, а его диагноз б
|
|||
размерный порог settings.yandex_detail_min_html_bytes) и считается исходом
|
||||
incomplete ⊆ failed — обогащения нет, значит следующий снапшот
|
||||
(detail_enriched_at IS NULL) возьмёт её снова.
|
||||
|
||||
Вечный недогруз (#3191, 13.09–16.09): карточка, которая не догружается НИКОГДА, без
|
||||
счётчика попыток возвращалась в голову каждого снапшота. Пять таких объявлений подряд
|
||||
выбивали брейкер, и 16 прогонов закончились ровно attempted=5 incomplete=5. Теперь
|
||||
недогруз пишется на объявление (миграция 321): сутки его не берём, после
|
||||
INCOMPLETE_MAX_ATTEMPTS не берём совсем, а повторный недогруз уже известной карточки
|
||||
не двигает брейкер — это сведения о карточке, а не о площадке.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -137,6 +144,13 @@ OFFER_URL_PATTERN = "/offer/[0-9]+"
|
|||
OFFER_ID_PATTERN = "^[0-9]+$"
|
||||
CANONICAL_URL_SQL = "'https://realty.yandex.ru/offer/' || source_id || '/'"
|
||||
|
||||
# Недогруз на объявлении (#3191). Разовый недогруз (1,8 МБ вместо 3,9 в замере 28.08)
|
||||
# с другого захода отдаётся целиком, поэтому пауза, а не приговор. Вечный — пустая
|
||||
# SPA-оболочка ~381 КБ без INITIAL_STATE, одинаковая с разных IP (проба 17.09) — после
|
||||
# трёх суточных попыток из очереди выходит и считается в incomplete_given_up.
|
||||
INCOMPLETE_RETRY_AFTER_HOURS = 24
|
||||
INCOMPLETE_MAX_ATTEMPTS = 3
|
||||
|
||||
|
||||
@dataclass
|
||||
class YandexDetailBackfillResult:
|
||||
|
|
@ -161,6 +175,8 @@ class YandexDetailBackfillResult:
|
|||
url_from_offer_id: int = 0
|
||||
# Ждут обогащения и адресовать их НЕЧЕМ: ни offer-URL, ни числового source_id.
|
||||
unenrichable_pending: int = 0
|
||||
# Ждут обогащения, но вышли из очереди после INCOMPLETE_MAX_ATTEMPTS недогрузов.
|
||||
incomplete_given_up: int = 0
|
||||
duration_sec: float = field(default=0.0)
|
||||
|
||||
def to_dict(self) -> dict[str, int]:
|
||||
|
|
@ -172,6 +188,7 @@ class YandexDetailBackfillResult:
|
|||
"failed": self.failed,
|
||||
"url_from_offer_id": self.url_from_offer_id,
|
||||
"unenrichable_pending": self.unenrichable_pending,
|
||||
"incomplete_given_up": self.incomplete_given_up,
|
||||
"duration_sec": int(self.duration_sec),
|
||||
}
|
||||
|
||||
|
|
@ -228,7 +245,8 @@ async def run_yandex_detail_backfill(
|
|||
WHEN source_url ~ CAST(:offer_url_pattern AS text)
|
||||
THEN source_url
|
||||
ELSE {CANONICAL_URL_SQL}
|
||||
END AS source_url
|
||||
END AS source_url,
|
||||
detail_incomplete_count
|
||||
FROM listings
|
||||
WHERE source = 'yandex'
|
||||
AND detail_enriched_at IS NULL
|
||||
|
|
@ -239,6 +257,15 @@ async def run_yandex_detail_backfill(
|
|||
)
|
||||
OR source_id ~ CAST(:offer_id_pattern AS text)
|
||||
)
|
||||
-- #3191: недогруженную карточку не берём сутки и не берём
|
||||
-- совсем после INCOMPLETE_MAX_ATTEMPTS.
|
||||
AND detail_incomplete_count < CAST(:incomplete_max_attempts AS int)
|
||||
AND (
|
||||
detail_incomplete_at IS NULL
|
||||
OR detail_incomplete_at < now() - make_interval(
|
||||
hours => CAST(:incomplete_retry_hours AS int)
|
||||
)
|
||||
)
|
||||
ORDER BY is_active DESC NULLS LAST, scraped_at DESC NULLS LAST
|
||||
LIMIT CAST(:batch_size AS int)
|
||||
"""
|
||||
|
|
@ -249,6 +276,8 @@ async def run_yandex_detail_backfill(
|
|||
"batch_size": batch_size,
|
||||
"offer_url_pattern": OFFER_URL_PATTERN,
|
||||
"offer_id_pattern": OFFER_ID_PATTERN,
|
||||
"incomplete_max_attempts": INCOMPLETE_MAX_ATTEMPTS,
|
||||
"incomplete_retry_hours": INCOMPLETE_RETRY_AFTER_HOURS,
|
||||
},
|
||||
)
|
||||
.mappings()
|
||||
|
|
@ -283,6 +312,19 @@ async def run_yandex_detail_backfill(
|
|||
).one()
|
||||
counters.url_from_offer_id = int(pending.url_from_offer_id)
|
||||
counters.unenrichable_pending = int(pending.unenrichable_pending)
|
||||
counters.incomplete_given_up = int(
|
||||
db.execute(
|
||||
text(
|
||||
"""
|
||||
SELECT count(*) FROM listings
|
||||
WHERE source = 'yandex'
|
||||
AND detail_enriched_at IS NULL
|
||||
AND detail_incomplete_count >= CAST(:incomplete_max_attempts AS int)
|
||||
"""
|
||||
),
|
||||
{"incomplete_max_attempts": INCOMPLETE_MAX_ATTEMPTS},
|
||||
).scalar_one()
|
||||
)
|
||||
if counters.url_from_offer_id:
|
||||
logger.info(
|
||||
"yandex_detail_backfill: run_id=%d — у %d объявлений сохранённый "
|
||||
|
|
@ -443,26 +485,59 @@ async def run_yandex_detail_backfill(
|
|||
# объявление уедет в БД с detail_enriched_at, выбыв из очереди
|
||||
# навсегда. Здесь оно исхода 'enriched' не получает, значит в
|
||||
# следующем прогоне снова попадёт в снапшот (detail_enriched_at
|
||||
# IS NULL). Серия таких страниц двигает consecutive_none — тот же
|
||||
# брейкер, что у parse→None: вечно недогружаемая карточка упрётся
|
||||
# в max_consecutive_blocks и оборвёт прогон, а не будет молотиться
|
||||
# (per-listing счётчика попыток в схеме нет, см. отчёт #3191).
|
||||
# IS NULL) — но не раньше чем через сутки: недогруз пишется на
|
||||
# объявление, и после INCOMPLETE_MAX_ATTEMPTS снапшот его не берёт.
|
||||
incomplete_reason = detail_incomplete_reason(
|
||||
resp.text, min_html_bytes=settings.yandex_detail_min_html_bytes
|
||||
)
|
||||
if incomplete_reason is not None:
|
||||
counters.incomplete += 1
|
||||
counters.failed += 1
|
||||
consecutive_none += 1
|
||||
# Брейкер ловит СИСТЕМНЫЙ недогруз (площадка отдаёт оболочки
|
||||
# всем). Повторный недогруз карточки, уже недогружавшейся
|
||||
# раньше, — сведения о ней, а не о площадке: он брейкер не
|
||||
# двигает, иначе вечные карточки, стоящие в очереди подряд,
|
||||
# рвали бы прогон при каждом повторе (#3191, 16 прогонов
|
||||
# attempted=5 incomplete=5). Первый недогруз двигает, как раньше.
|
||||
prior_incomplete = int(row.get("detail_incomplete_count") or 0)
|
||||
if prior_incomplete == 0:
|
||||
consecutive_none += 1
|
||||
# Площадка ОТВЕТИЛА (HTTP 200) — серии блоков нет (#3196).
|
||||
consecutive_blocks = 0
|
||||
try:
|
||||
db.execute(
|
||||
text(
|
||||
"""
|
||||
UPDATE listings
|
||||
SET detail_incomplete_count = detail_incomplete_count + 1,
|
||||
detail_incomplete_at = now()
|
||||
WHERE id = CAST(:listing_id AS bigint)
|
||||
"""
|
||||
),
|
||||
{"listing_id": listing_id},
|
||||
)
|
||||
db.commit()
|
||||
except Exception as mark_exc:
|
||||
# Не записали — карточка просто вернётся в следующий
|
||||
# снапшот без паузы (поведение до #3191-счётчика).
|
||||
db.rollback()
|
||||
logger.warning(
|
||||
"yandex_detail_backfill: run_id=%d listing_id=%d "
|
||||
"недогруз не записан: %s",
|
||||
run_id,
|
||||
listing_id,
|
||||
mark_exc,
|
||||
)
|
||||
logger.warning(
|
||||
"yandex_detail_backfill: run_id=%d listing_id=%d source_url=%s "
|
||||
"-> недогруженная карточка, отказ: %s (consecutive=%d)",
|
||||
"-> недогруженная карточка, отказ: %s (попытка %d/%d, "
|
||||
"consecutive=%d)",
|
||||
run_id,
|
||||
listing_id,
|
||||
source_url,
|
||||
incomplete_reason,
|
||||
prior_incomplete + 1,
|
||||
INCOMPLETE_MAX_ATTEMPTS,
|
||||
consecutive_none,
|
||||
)
|
||||
if consecutive_none >= max_consecutive_blocks:
|
||||
|
|
|
|||
|
|
@ -0,0 +1,52 @@
|
|||
-- 321_listings_detail_incomplete_attempts.sql
|
||||
-- Счётчик недогруженных карточек Яндекса на объявлении (#3191, вечный недогруз).
|
||||
--
|
||||
-- Apply after: 310_yandex_seed_decimal_slips.sql
|
||||
-- Deploy order: схема первой (деплой применяет миграции до пересоздания контейнеров);
|
||||
-- старый код новые колонки не читает, поэтому окно между шагами безопасно.
|
||||
--
|
||||
-- ── ЧТО НЕ ТАК ────────────────────────────────────────────────────────────────
|
||||
-- PR #3364 сделал недогруженную карточку отказом: detail_enriched_at не ставится,
|
||||
-- объявление остаётся в очереди yandex_detail_backfill. Попыток на объявлении при
|
||||
-- этом никто не считал, и карточка, которая не догружается НИКОГДА, возвращалась в
|
||||
-- голову каждого снапшота. Прод 13.09 21:14 – 16.09 18:52: пять объявлений
|
||||
-- (10776456, 10775964, 10775945, 10775747, 10775594; scraped_at 12.09 15:02, подряд
|
||||
-- в порядке очереди) давали 381 КБ без encryptedPhones/redirectPhones в каждом
|
||||
-- прогоне, пять недогрузов подряд выбивали брейкер consecutive_none=5, и 16 прогонов
|
||||
-- закончились ровно attempted=5 incomplete=5, а остальные обрывались, дойдя до них.
|
||||
-- Страница — пустая SPA-оболочка без INITIAL_STATE (id оффера в ней не встречается),
|
||||
-- та же с другого IP (проба 17.09): это свойство объявления, а не сети.
|
||||
--
|
||||
-- ── ЧТО ДЕЛАЕТ ФАЙЛ ───────────────────────────────────────────────────────────
|
||||
-- detail_incomplete_count — сколько раз карточка пришла недогруженной;
|
||||
-- detail_incomplete_at — когда последний раз. Читает и пишет только
|
||||
-- app/tasks/yandex_detail_backfill.py: снапшот не берёт карточку в течение суток
|
||||
-- после недогруза и не берёт совсем после трёх.
|
||||
--
|
||||
-- ── СТОИМОСТЬ ─────────────────────────────────────────────────────────────────
|
||||
-- ADD COLUMN с константным DEFAULT в PostgreSQL 11+ не переписывает heap (значение
|
||||
-- лежит в каталоге), поэтому 19 ГБ listings не копируются. ACCESS EXCLUSIVE берётся
|
||||
-- на миллисекунды; ожидание лока ограничено lock_timeout (.claude/rules/sql.md).
|
||||
--
|
||||
-- IDEMPOTENCY: ADD COLUMN IF NOT EXISTS. Данные существующих строк не трогаются.
|
||||
--
|
||||
-- Критерий приёмки (записан ДО применения):
|
||||
-- 1. Запись в _schema_migrations по имени этого файла.
|
||||
-- 2. information_schema.columns: обе колонки у listings.
|
||||
-- 3. После первого прогона yandex_detail_backfill, дошедшего до пяти объявлений:
|
||||
-- у них detail_incomplete_count = 1, и следующий прогон их не запрашивает.
|
||||
|
||||
BEGIN;
|
||||
|
||||
SET LOCAL lock_timeout = '5s';
|
||||
|
||||
ALTER TABLE listings
|
||||
ADD COLUMN IF NOT EXISTS detail_incomplete_count smallint NOT NULL DEFAULT 0,
|
||||
ADD COLUMN IF NOT EXISTS detail_incomplete_at timestamptz;
|
||||
|
||||
COMMENT ON COLUMN listings.detail_incomplete_count IS
|
||||
'Сколько раз detail-страница приходила недогруженной (#3191). Пишет yandex_detail_backfill; после 3 объявление выходит из очереди добора.';
|
||||
COMMENT ON COLUMN listings.detail_incomplete_at IS
|
||||
'Время последнего недогруза detail-страницы (#3191). Суточная пауза перед повтором в yandex_detail_backfill.';
|
||||
|
||||
COMMIT;
|
||||
|
|
@ -182,3 +182,8 @@ tests/test_3404_egress_run_attribution.py::test_attribution_does_not_commit_call
|
|||
# снятия подписи батча, сужения окна до ×10 и снятия порога — проверено вручную 17.09.
|
||||
tests/test_3385_migration_310_yandex_seed_slips.py::test_deletes_only_seed_off_by_order_and_is_idempotent
|
||||
tests/test_3385_migration_310_yandex_seed_slips.py::test_stops_and_rolls_back_when_too_many
|
||||
|
||||
# Вечный недогруз Яндекса (#3191, миграция 321): снапшот-SELECT очереди добора с
|
||||
# суточной паузой и выбыванием после трёх недогрузов судит только Postgres — на
|
||||
# мок-лэйне БД нет. В ci-tradein.yml бежит по-настоящему (postgres-сервис, #2745).
|
||||
tests/test_3191_yandex_perpetual_underloaded.py::test_perpetual_underloaded_cards_leave_queue_head
|
||||
|
|
|
|||
|
|
@ -0,0 +1,230 @@
|
|||
"""Вечно недогружаемые карточки Яндекса не держат очередь добора (#3191).
|
||||
|
||||
Прод 13.09 21:14 – 16.09 18:52: пять объявлений (10776456, 10775964, 10775945,
|
||||
10775747, 10775594), стоящих в очереди подряд, в каждом прогоне приходили пустой
|
||||
SPA-оболочкой ~381 КБ без блока контактов. Пять недогрузов подряд выбивали брейкер
|
||||
consecutive_none=5: 16 прогонов закончились ровно attempted=5 incomplete=5, а
|
||||
«успешные» обрывались, едва дойдя до них. За ними — 20 тыс. необогащённых объявлений.
|
||||
|
||||
Проверки ПО ЗНАЧЕНИЮ:
|
||||
* мок-часть — какие карточки запрошены и какие счётчики вышли;
|
||||
* живой Postgres (в CI есть, локально skip) — что НАСТОЯЩИЙ снапшот-SELECT после
|
||||
недогруза не берёт карточку сутки, берёт её после паузы и выбрасывает после трёх.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import sys
|
||||
import uuid
|
||||
from typing import Any
|
||||
from unittest.mock import AsyncMock, MagicMock, patch
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
_wp_mock = MagicMock()
|
||||
sys.modules.setdefault("weasyprint", _wp_mock)
|
||||
|
||||
import pytest # noqa: E402
|
||||
from sqlalchemy import text # noqa: E402
|
||||
|
||||
from app.tasks.yandex_detail_backfill import run_yandex_detail_backfill # noqa: E402
|
||||
|
||||
_ASYNC_SESSION = "app.tasks.yandex_detail_backfill.AsyncSession"
|
||||
_PARSE = "app.tasks.yandex_detail_backfill.YandexDetailScraper.parse"
|
||||
_SAVE = "app.tasks.yandex_detail_backfill.save_detail_enrichment"
|
||||
_RUNS = "app.tasks.yandex_detail_backfill.runs_mod"
|
||||
_SLEEP = "app.tasks.yandex_detail_backfill.asyncio.sleep"
|
||||
_RESOLVE_PROXY_URL = "app.tasks.yandex_detail_backfill.resolve_proxy_url"
|
||||
|
||||
# Оболочка как на проде: 381 КБ, контактов нет. Полная — 3,9 МБ с маркером.
|
||||
SHELL_HTML = "<html>" + "x" * 381_157 + "</html>"
|
||||
FULL_HTML = "<html>" + "x" * 3_920_119 + '"encryptedPhones":["a"]</html>'
|
||||
|
||||
PERPETUAL = 5
|
||||
NORMAL = 3
|
||||
|
||||
|
||||
def _resp(html: str) -> MagicMock:
|
||||
resp = MagicMock()
|
||||
resp.status_code = 200
|
||||
resp.text = html
|
||||
return resp
|
||||
|
||||
|
||||
async def _run(db: Any, perpetual_urls: set[str]) -> tuple[Any, list[str]]:
|
||||
"""Один прогон; оболочку отдаём по URL, всё остальное — полная карточка."""
|
||||
requested: list[str] = []
|
||||
|
||||
async def get(url: str, **_kw: Any) -> MagicMock:
|
||||
requested.append(url)
|
||||
return _resp(SHELL_HTML if url in perpetual_urls else FULL_HTML)
|
||||
|
||||
session = AsyncMock()
|
||||
session.get = AsyncMock(side_effect=get)
|
||||
ctx = MagicMock()
|
||||
ctx.__aenter__ = AsyncMock(return_value=session)
|
||||
ctx.__aexit__ = AsyncMock(return_value=None)
|
||||
|
||||
with (
|
||||
patch(_ASYNC_SESSION, MagicMock(return_value=ctx)),
|
||||
patch(_PARSE, return_value=MagicMock()),
|
||||
patch(_SAVE, MagicMock(return_value=True)),
|
||||
patch(_RUNS, MagicMock()),
|
||||
patch(_SLEEP, new_callable=AsyncMock),
|
||||
patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://proxy:3128")),
|
||||
):
|
||||
counters = await run_yandex_detail_backfill(
|
||||
db, run_id=3191, params={"batch_size": PERPETUAL + NORMAL, "budget_sec": 3600}
|
||||
)
|
||||
return counters, requested
|
||||
|
||||
|
||||
# ── 1. Брейкер: повтор известной карточки не рвёт прогон, первый недогруз — рвёт ──
|
||||
|
||||
|
||||
def _mock_db(prior_incomplete: int) -> tuple[MagicMock, set[str], list[str]]:
|
||||
rows = [
|
||||
{
|
||||
"id": i + 1,
|
||||
"source_url": f"https://realty.yandex.ru/offer/{i + 1}/",
|
||||
"detail_incomplete_count": prior_incomplete if i < PERPETUAL else 0,
|
||||
}
|
||||
for i in range(PERPETUAL + NORMAL)
|
||||
]
|
||||
db = MagicMock()
|
||||
sel = MagicMock()
|
||||
sel.mappings.return_value.all.return_value = rows
|
||||
sel.one.return_value = MagicMock(url_from_offer_id=0, unenrichable_pending=0)
|
||||
sel.scalar_one.return_value = 0
|
||||
db.execute.return_value = sel
|
||||
urls = [r["source_url"] for r in rows]
|
||||
return db, set(urls[:PERPETUAL]), urls[PERPETUAL:]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_retry_of_known_underloaded_cards_does_not_abort_run() -> None:
|
||||
"""Пять карточек, уже недогружавшихся раньше, в голове снапшота — прогон идёт дальше."""
|
||||
db, perpetual, normal = _mock_db(prior_incomplete=1)
|
||||
|
||||
counters, requested = await _run(db, perpetual)
|
||||
|
||||
assert set(normal) <= set(requested), (
|
||||
f"запрошено {len(requested)} из {PERPETUAL + NORMAL}: повторный недогруз известных "
|
||||
"карточек оборвал прогон раньше, чем очередь дошла до следующих"
|
||||
)
|
||||
assert (counters.attempted, counters.enriched, counters.incomplete) == (8, 3, 5)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_first_underload_series_still_aborts_run() -> None:
|
||||
"""Защита от системного недогруза жива: пять ПЕРВЫХ недогрузов подряд — обрыв."""
|
||||
db, perpetual, normal = _mock_db(prior_incomplete=0)
|
||||
|
||||
counters, requested = await _run(db, perpetual)
|
||||
|
||||
assert not set(normal) & set(requested), (
|
||||
"пять первых недогрузов подряд не оборвали прогон — брейкер системного "
|
||||
"недогруза больше не срабатывает"
|
||||
)
|
||||
assert (counters.attempted, counters.enriched, counters.incomplete) == (5, 0, 5)
|
||||
|
||||
|
||||
# ── 2. Очередь на живом Postgres ─────────────────────────────────────────────
|
||||
|
||||
|
||||
def _live_session() -> Any | None:
|
||||
try:
|
||||
from sqlalchemy import create_engine
|
||||
from sqlalchemy.orm import sessionmaker
|
||||
|
||||
dsn = os.environ.get("TEST_DATABASE_URL") or os.environ.get("DATABASE_URL", "")
|
||||
if not dsn or "localhost:5432/test" in dsn:
|
||||
return None
|
||||
engine = create_engine(dsn, future=True)
|
||||
conn = engine.connect()
|
||||
conn.execute(text("SELECT 1"))
|
||||
conn.close()
|
||||
return sessionmaker(bind=engine, future=True)()
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
@pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB")
|
||||
@pytest.mark.asyncio
|
||||
async def test_perpetual_underloaded_cards_leave_queue_head() -> None:
|
||||
db = _live_session()
|
||||
tag = f"t3191-{uuid.uuid4().hex[:8]}"
|
||||
try:
|
||||
ids: list[int] = []
|
||||
urls: list[str] = []
|
||||
for i in range(PERPETUAL + NORMAL):
|
||||
url = f"https://realty.yandex.ru/offer/3191{uuid.uuid4().int % 10**12}/"
|
||||
# Вечные — самые свежие (голова очереди), как 12.09 15:02 на проде.
|
||||
# Дата в будущем: чужие необогащённые строки тестовой БД не встанут впереди.
|
||||
row_id = db.execute(
|
||||
text(
|
||||
"""
|
||||
INSERT INTO listings (source, source_url, source_id, dedup_hash,
|
||||
price_rub, is_active, scraped_at)
|
||||
VALUES ('yandex', :url, :sid, :sid, 5000000, true,
|
||||
TIMESTAMPTZ '2100-01-01' - make_interval(mins => :i))
|
||||
RETURNING id
|
||||
"""
|
||||
),
|
||||
{"url": url, "sid": f"{tag}-{i}", "i": i},
|
||||
).scalar_one()
|
||||
ids.append(row_id)
|
||||
urls.append(url)
|
||||
db.commit()
|
||||
perpetual = set(urls[:PERPETUAL])
|
||||
normal = urls[PERPETUAL:]
|
||||
|
||||
# Прогон 1: карточки ещё не недогружались — первый недогруз, обрыв на пятой.
|
||||
_c1, req1 = await _run(db, perpetual)
|
||||
assert set(req1) == perpetual, req1
|
||||
|
||||
# Прогон 2 сразу следом: голова очереди больше не держит — берутся следующие.
|
||||
c2, req2 = await _run(db, perpetual)
|
||||
assert not set(req2) & perpetual, (
|
||||
"недогруженные час назад карточки снова в снапшоте — прогон опять "
|
||||
"упрётся в них, как 16 прогонов 13.09–16.09"
|
||||
)
|
||||
assert set(normal) <= set(req2)
|
||||
assert c2.enriched == NORMAL
|
||||
|
||||
# Сутки спустя: повтор (разовый недогруз мог пройти) и прогон не обрывается.
|
||||
def age(hours: int) -> None:
|
||||
db.execute(
|
||||
text(
|
||||
"UPDATE listings SET detail_incomplete_at = now() - make_interval(hours => :h) "
|
||||
"WHERE id = ANY(:ids)"
|
||||
),
|
||||
{"h": hours, "ids": ids[:PERPETUAL]},
|
||||
)
|
||||
db.commit()
|
||||
|
||||
age(25)
|
||||
c3, req3 = await _run(db, perpetual)
|
||||
assert perpetual <= set(req3), "после суток паузы карточки обязаны вернуться в очередь"
|
||||
assert set(normal) <= set(req3), "повтор известных карточек оборвал прогон"
|
||||
assert c3.incomplete == PERPETUAL
|
||||
|
||||
age(25)
|
||||
await _run(db, perpetual) # третья попытка
|
||||
age(25)
|
||||
c5, req5 = await _run(db, perpetual)
|
||||
assert not set(req5) & perpetual, "после трёх недогрузов карточка всё ещё в очереди"
|
||||
assert c5.incomplete_given_up >= PERPETUAL, (
|
||||
f"incomplete_given_up={c5.incomplete_given_up}: выбывшие из очереди не посчитаны"
|
||||
)
|
||||
counts = db.execute(
|
||||
text("SELECT detail_incomplete_count FROM listings WHERE id = ANY(:ids) ORDER BY id"),
|
||||
{"ids": ids},
|
||||
).scalars()
|
||||
assert list(counts) == [3] * PERPETUAL + [0] * NORMAL
|
||||
finally:
|
||||
db.rollback()
|
||||
db.execute(text("DELETE FROM listings WHERE source_id LIKE :p"), {"p": f"{tag}-%"})
|
||||
db.commit()
|
||||
db.close()
|
||||
Loading…
Add table
Reference in a new issue