feat(tradein): checkpoint/resume for cian full-load — survive restart (#930) #932
4 changed files with 181 additions and 17 deletions
|
|
@ -982,6 +982,14 @@ class CianFullLoadRequest(BaseModel):
|
||||||
le=50,
|
le=50,
|
||||||
description="Сколько свежих листингов обогатить detail (актуально при enrich_detail=True).",
|
description="Сколько свежих листингов обогатить detail (актуально при enrich_detail=True).",
|
||||||
)
|
)
|
||||||
|
resume_run_id: int | None = Field(
|
||||||
|
default=None,
|
||||||
|
description=(
|
||||||
|
"Если задан — читает done_buckets из counters прошлого run и пропускает "
|
||||||
|
"уже завершённые бакеты. Передать run_id предыдущего (прерванного) "
|
||||||
|
"cian_full_load прогона для resume. Без этого поля — full walk с нуля."
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
@router.post("/scrape/cian-full-load", response_model=CitySweepStartResponse)
|
@router.post("/scrape/cian-full-load", response_model=CitySweepStartResponse)
|
||||||
|
|
@ -1016,6 +1024,7 @@ async def start_cian_full_load(
|
||||||
concurrency=payload.concurrency,
|
concurrency=payload.concurrency,
|
||||||
enrich_detail=payload.enrich_detail,
|
enrich_detail=payload.enrich_detail,
|
||||||
detail_top_n=payload.detail_top_n,
|
detail_top_n=payload.detail_top_n,
|
||||||
|
resume_run_id=payload.resume_run_id,
|
||||||
)
|
)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("cian-full-load background task run_id=%d crashed", run_id)
|
logger.exception("cian-full-load background task run_id=%d crashed", run_id)
|
||||||
|
|
|
||||||
|
|
@ -1558,6 +1558,7 @@ async def run_cian_full_load(
|
||||||
request_delay_sec: float = 4.0,
|
request_delay_sec: float = 4.0,
|
||||||
enrich_detail: bool = False,
|
enrich_detail: bool = False,
|
||||||
concurrency: int = 5,
|
concurrency: int = 5,
|
||||||
|
resume_run_id: int | None = None,
|
||||||
) -> CianFullLoadCounters:
|
) -> CianFullLoadCounters:
|
||||||
"""Exhaustive региональный сбор Cian ЕКБ вторички (БЕЗ anchor'ов).
|
"""Exhaustive региональный сбор Cian ЕКБ вторички (БЕЗ anchor'ов).
|
||||||
|
|
||||||
|
|
@ -1581,22 +1582,52 @@ async def run_cian_full_load(
|
||||||
|
|
||||||
counters = CianFullLoadCounters()
|
counters = CianFullLoadCounters()
|
||||||
|
|
||||||
def _on_bucket(lots: list) -> None: # type: ignore[type-arg]
|
# ── Checkpoint/resume: читаем done_buckets из прошлого run ───────────────
|
||||||
"""Инкрементальный save после каждого leaf-бакета."""
|
skip_set: set[str] = set()
|
||||||
|
if resume_run_id is not None:
|
||||||
|
_prev_row = db.execute(
|
||||||
|
text("SELECT counters FROM scrape_runs WHERE id = CAST(:rid AS bigint)"),
|
||||||
|
{"rid": resume_run_id},
|
||||||
|
).fetchone()
|
||||||
|
if _prev_row is not None and _prev_row.counters:
|
||||||
|
_prev_counters: dict = (
|
||||||
|
_prev_row.counters if isinstance(_prev_row.counters, dict) else {}
|
||||||
|
)
|
||||||
|
skip_set = set(_prev_counters.get("done_buckets", []))
|
||||||
|
logger.info(
|
||||||
|
"cian-full-load run_id=%d: resuming from run %d — %d buckets already done",
|
||||||
|
run_id,
|
||||||
|
resume_run_id,
|
||||||
|
len(skip_set),
|
||||||
|
)
|
||||||
|
|
||||||
|
done: set[str] = set(skip_set) # накапливаем завершённые бакеты этого прогона
|
||||||
|
|
||||||
|
def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg]
|
||||||
|
"""Инкрементальный save после каждого leaf-бакета. Дописывает bucket_key в done."""
|
||||||
|
nonlocal done
|
||||||
if scrape_runs.is_cancelled(db, run_id):
|
if scrape_runs.is_cancelled(db, run_id):
|
||||||
logger.info("cian-full-load run_id=%d: cancel detected in on_bucket", run_id)
|
logger.info("cian-full-load run_id=%d: cancel detected in on_bucket", run_id)
|
||||||
raise RuntimeError("cancelled")
|
raise RuntimeError("cancelled")
|
||||||
if not lots:
|
if not lots:
|
||||||
|
done.add(bucket_key)
|
||||||
|
scrape_runs.update_heartbeat(
|
||||||
|
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
|
||||||
|
)
|
||||||
return
|
return
|
||||||
inserted, updated = save_listings(db, lots, run_id=run_id)
|
inserted, updated = save_listings(db, lots, run_id=run_id)
|
||||||
# save_listings вызывает db.commit() внутри — данные в БД сразу
|
# save_listings вызывает db.commit() внутри — данные в БД сразу
|
||||||
counters.saved_inserted += inserted
|
counters.saved_inserted += inserted
|
||||||
counters.saved_updated += updated
|
counters.saved_updated += updated
|
||||||
counters.unique_fetched += len(lots)
|
counters.unique_fetched += len(lots)
|
||||||
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
|
done.add(bucket_key)
|
||||||
|
scrape_runs.update_heartbeat(
|
||||||
|
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
|
||||||
|
)
|
||||||
logger.info(
|
logger.info(
|
||||||
"cian-full-load run_id=%d: bucket saved ins=%d upd=%d total_unique=%d",
|
"cian-full-load run_id=%d: bucket %s saved ins=%d upd=%d total_unique=%d",
|
||||||
run_id,
|
run_id,
|
||||||
|
bucket_key,
|
||||||
inserted,
|
inserted,
|
||||||
updated,
|
updated,
|
||||||
counters.unique_fetched,
|
counters.unique_fetched,
|
||||||
|
|
@ -1615,6 +1646,7 @@ async def run_cian_full_load(
|
||||||
concurrency=concurrency,
|
concurrency=concurrency,
|
||||||
on_bucket=_on_bucket,
|
on_bucket=_on_bucket,
|
||||||
on_progress=_on_progress,
|
on_progress=_on_progress,
|
||||||
|
skip_buckets=skip_set if skip_set else None,
|
||||||
)
|
)
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
|
|
@ -1687,7 +1719,7 @@ async def run_cian_full_load(
|
||||||
)
|
)
|
||||||
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
|
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
|
||||||
|
|
||||||
scrape_runs.mark_done(db, run_id, counters.to_dict())
|
scrape_runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
||||||
logger.info(
|
logger.info(
|
||||||
"cian-full-load run_id=%d done: unique=%d ins=%d upd=%d detail=%d/%d errors=%d",
|
"cian-full-load run_id=%d done: unique=%d ins=%d upd=%d detail=%d/%d errors=%d",
|
||||||
run_id,
|
run_id,
|
||||||
|
|
@ -1710,7 +1742,7 @@ async def run_cian_full_load(
|
||||||
counters.saved_inserted,
|
counters.saved_inserted,
|
||||||
counters.saved_updated,
|
counters.saved_updated,
|
||||||
)
|
)
|
||||||
scrape_runs.mark_done(db, run_id, counters.to_dict())
|
scrape_runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
||||||
return counters
|
return counters
|
||||||
logger.exception("cian-full-load run_id=%d: fatal error", run_id)
|
logger.exception("cian-full-load run_id=%d: fatal error", run_id)
|
||||||
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
||||||
|
|
|
||||||
|
|
@ -320,6 +320,7 @@ class CianScraper(BaseScraper):
|
||||||
concurrency: int = 5,
|
concurrency: int = 5,
|
||||||
on_bucket: Callable[..., Any] | None = None,
|
on_bucket: Callable[..., Any] | None = None,
|
||||||
on_progress: Callable[[int], None] | None = None,
|
on_progress: Callable[[int], None] | None = None,
|
||||||
|
skip_buckets: set[str] | None = None,
|
||||||
) -> list[ScrapedLot]:
|
) -> list[ScrapedLot]:
|
||||||
"""Exhaustive-загрузка Cian ЕКБ вторички через партиционирование КОМНАТНОСТЬ × ЦЕНА.
|
"""Exhaustive-загрузка Cian ЕКБ вторички через партиционирование КОМНАТНОСТЬ × ЦЕНА.
|
||||||
|
|
||||||
|
|
@ -334,9 +335,12 @@ class CianScraper(BaseScraper):
|
||||||
price_cap_per_bucket: максимум офферов в бакете перед делением (< 1500).
|
price_cap_per_bucket: максимум офферов в бакете перед делением (< 1500).
|
||||||
max_pages_per_bucket: Cian hard cap ~54; не превышать.
|
max_pages_per_bucket: Cian hard cap ~54; не превышать.
|
||||||
concurrency: максимум параллельных page-фетчей в leaf-бакете (default=5).
|
concurrency: максимум параллельных page-фетчей в leaf-бакете (default=5).
|
||||||
on_bucket: опциональный callback(list[ScrapedLot]) после каждого leaf-бакета.
|
on_bucket: опциональный callback(bucket_key, list[ScrapedLot]) после каждого
|
||||||
Может быть async или sync. Если кидает исключение — прерывает прогон.
|
leaf-бакета. Может быть async или sync. Если кидает исключение — прерывает прогон.
|
||||||
on_progress: опциональный callback(unique_count) для heartbeat (per room-bucket).
|
on_progress: опциональный callback(unique_count) для heartbeat (per room-bucket).
|
||||||
|
skip_buckets: множество ключей «room_label:lo:hi» уже завершённых бакетов —
|
||||||
|
пагинация и on_bucket для них пропускаются. Probe-запросы (для split-решения)
|
||||||
|
всё равно выполняются (их мало по сравнению с пагинацией).
|
||||||
|
|
||||||
Возвращает list[ScrapedLot] уникальных лотов (дедуп по source_id/source_url).
|
Возвращает list[ScrapedLot] уникальных лотов (дедуп по source_id/source_url).
|
||||||
"""
|
"""
|
||||||
|
|
@ -361,6 +365,7 @@ class CianScraper(BaseScraper):
|
||||||
max_pages_per_bucket=max_pages_per_bucket,
|
max_pages_per_bucket=max_pages_per_bucket,
|
||||||
concurrency=concurrency,
|
concurrency=concurrency,
|
||||||
on_bucket=on_bucket,
|
on_bucket=on_bucket,
|
||||||
|
skip_buckets=skip_buckets,
|
||||||
)
|
)
|
||||||
room_collected = len(seen) - before
|
room_collected = len(seen) - before
|
||||||
logger.info(
|
logger.info(
|
||||||
|
|
@ -386,6 +391,7 @@ class CianScraper(BaseScraper):
|
||||||
max_pages_per_bucket: int,
|
max_pages_per_bucket: int,
|
||||||
concurrency: int = 5,
|
concurrency: int = 5,
|
||||||
on_bucket: Callable[..., Any] | None = None,
|
on_bucket: Callable[..., Any] | None = None,
|
||||||
|
skip_buckets: set[str] | None = None,
|
||||||
_depth: int = 0,
|
_depth: int = 0,
|
||||||
) -> None:
|
) -> None:
|
||||||
"""Рекурсивное адаптивное бинарное партиционирование ценового диапазона [lo, hi].
|
"""Рекурсивное адаптивное бинарное партиционирование ценового диапазона [lo, hi].
|
||||||
|
|
@ -396,8 +402,10 @@ class CianScraper(BaseScraper):
|
||||||
3. Если totalOffers > cap → разбить бакет пополам (рекурсия).
|
3. Если totalOffers > cap → разбить бакет пополам (рекурсия).
|
||||||
Guard: hi - lo < _MIN_BRACKET → пагинировать как есть (логируем WARNING).
|
Guard: hi - lo < _MIN_BRACKET → пагинировать как есть (логируем WARNING).
|
||||||
|
|
||||||
После пагинации leaf-бакета вызывает on_bucket(bucket_lots) если задан.
|
После пагинации leaf-бакета вызывает on_bucket(bucket_key, bucket_lots) если задан.
|
||||||
|
bucket_key = "room_label:lo:hi" (checkpoint-ключ для resume).
|
||||||
on_bucket может быть async или sync. Исключение в on_bucket прерывает прогон.
|
on_bucket может быть async или sync. Исключение в on_bucket прерывает прогон.
|
||||||
|
skip_buckets: если bucket_key в skip_buckets — пагинацию и on_bucket пропускаем.
|
||||||
"""
|
"""
|
||||||
effective_hi = hi
|
effective_hi = hi
|
||||||
|
|
||||||
|
|
@ -481,6 +489,7 @@ class CianScraper(BaseScraper):
|
||||||
max_pages_per_bucket=max_pages_per_bucket,
|
max_pages_per_bucket=max_pages_per_bucket,
|
||||||
concurrency=concurrency,
|
concurrency=concurrency,
|
||||||
on_bucket=on_bucket,
|
on_bucket=on_bucket,
|
||||||
|
skip_buckets=skip_buckets,
|
||||||
_depth=_depth + 1,
|
_depth=_depth + 1,
|
||||||
)
|
)
|
||||||
# [mid+1, hi]
|
# [mid+1, hi]
|
||||||
|
|
@ -493,10 +502,22 @@ class CianScraper(BaseScraper):
|
||||||
max_pages_per_bucket=max_pages_per_bucket,
|
max_pages_per_bucket=max_pages_per_bucket,
|
||||||
concurrency=concurrency,
|
concurrency=concurrency,
|
||||||
on_bucket=on_bucket,
|
on_bucket=on_bucket,
|
||||||
|
skip_buckets=skip_buckets,
|
||||||
_depth=_depth + 1,
|
_depth=_depth + 1,
|
||||||
)
|
)
|
||||||
return
|
return
|
||||||
|
|
||||||
|
# ── Checkpoint-ключ leaf-бакета ──────────────────────────────────────
|
||||||
|
bucket_key = f"{room_label}:{lo}:{hi}"
|
||||||
|
|
||||||
|
# Если бакет уже завершён в предыдущем запуске — пропускаем пагинацию и on_bucket.
|
||||||
|
if skip_buckets and bucket_key in skip_buckets:
|
||||||
|
logger.info(
|
||||||
|
"cian: skip bucket %s — already done (resume)",
|
||||||
|
bucket_key,
|
||||||
|
)
|
||||||
|
return
|
||||||
|
|
||||||
# ── Параллельная пагинация leaf-бакета ────────────────────────────────
|
# ── Параллельная пагинация leaf-бакета ────────────────────────────────
|
||||||
max_pages = min(
|
max_pages = min(
|
||||||
math.ceil(total / _CIAN_OFFERS_PER_PAGE),
|
math.ceil(total / _CIAN_OFFERS_PER_PAGE),
|
||||||
|
|
@ -563,7 +584,7 @@ class CianScraper(BaseScraper):
|
||||||
|
|
||||||
# ── on_bucket callback: инкрементальный save ──────────────────────────
|
# ── on_bucket callback: инкрементальный save ──────────────────────────
|
||||||
if on_bucket is not None and bucket_lots:
|
if on_bucket is not None and bucket_lots:
|
||||||
res_cb = on_bucket(bucket_lots)
|
res_cb = on_bucket(bucket_key, bucket_lots)
|
||||||
if inspect.isawaitable(res_cb):
|
if inspect.isawaitable(res_cb):
|
||||||
await res_cb
|
await res_cb
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -6,8 +6,12 @@
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
from unittest.mock import AsyncMock, MagicMock, patch
|
from unittest.mock import AsyncMock, MagicMock, patch
|
||||||
|
|
||||||
|
# Settings requires DATABASE_URL at import time — set dummy DSN before any app import.
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from app.services.scrapers.base import ScrapedLot
|
from app.services.scrapers.base import ScrapedLot
|
||||||
|
|
@ -213,10 +217,10 @@ async def test_fetch_all_secondary_min_bracket_guard(scraper: CianScraper) -> No
|
||||||
async def test_fetch_all_secondary_on_bucket_called_per_leaf(scraper: CianScraper) -> None:
|
async def test_fetch_all_secondary_on_bucket_called_per_leaf(scraper: CianScraper) -> None:
|
||||||
"""on_bucket вызывается после каждого leaf-бакета с лотами бакета."""
|
"""on_bucket вызывается после каждого leaf-бакета с лотами бакета."""
|
||||||
# totalOffers=56 ≤ cap → leaf-бакет, пагинируется параллельно 2 страницы
|
# totalOffers=56 ≤ cap → leaf-бакет, пагинируется параллельно 2 страницы
|
||||||
bucket_calls: list[list[ScrapedLot]] = []
|
bucket_calls: list[tuple[str, list[ScrapedLot]]] = []
|
||||||
|
|
||||||
def fake_on_bucket(lots: list[ScrapedLot]) -> None:
|
def fake_on_bucket(bucket_key: str, lots: list[ScrapedLot]) -> None:
|
||||||
bucket_calls.append(list(lots))
|
bucket_calls.append((bucket_key, list(lots)))
|
||||||
|
|
||||||
async def fake_fetch_page_html(
|
async def fake_fetch_page_html(
|
||||||
rooms: tuple, page: int, min_price: int | None, max_price: int | None
|
rooms: tuple, page: int, min_price: int | None, max_price: int | None
|
||||||
|
|
@ -249,17 +253,18 @@ async def test_fetch_all_secondary_on_bucket_called_per_leaf(scraper: CianScrape
|
||||||
|
|
||||||
# on_bucket вызван ровно 1 раз (один leaf-бакет для одной комнатности)
|
# on_bucket вызван ровно 1 раз (один leaf-бакет для одной комнатности)
|
||||||
assert len(bucket_calls) == 1, f"Ожидался 1 вызов on_bucket, получено {len(bucket_calls)}"
|
assert len(bucket_calls) == 1, f"Ожидался 1 вызов on_bucket, получено {len(bucket_calls)}"
|
||||||
|
key, lots = bucket_calls[0]
|
||||||
|
# Ключ бакета содержит room_label:lo:hi
|
||||||
|
assert "room1" in key, f"bucket_key должен содержать room1, got {key!r}"
|
||||||
# Лоты из обеих страниц переданы в on_bucket
|
# Лоты из обеих страниц переданы в on_bucket
|
||||||
assert (
|
assert len(lots) == 20, f"Ожидалось 20 лотов в on_bucket (2 стр × 10), получено {len(lots)}"
|
||||||
len(bucket_calls[0]) == 20
|
|
||||||
), f"Ожидалось 20 лотов в on_bucket (2 стр × 10), получено {len(bucket_calls[0])}"
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_fetch_all_secondary_on_bucket_cancel_stops_run(scraper: CianScraper) -> None:
|
async def test_fetch_all_secondary_on_bucket_cancel_stops_run(scraper: CianScraper) -> None:
|
||||||
"""Если on_bucket кидает RuntimeError('cancelled') — прогон прерывается."""
|
"""Если on_bucket кидает RuntimeError('cancelled') — прогон прерывается."""
|
||||||
|
|
||||||
def cancel_on_bucket(lots: list[ScrapedLot]) -> None:
|
def cancel_on_bucket(bucket_key: str, lots: list[ScrapedLot]) -> None:
|
||||||
raise RuntimeError("cancelled")
|
raise RuntimeError("cancelled")
|
||||||
|
|
||||||
async def fake_fetch_page_html(
|
async def fake_fetch_page_html(
|
||||||
|
|
@ -331,3 +336,100 @@ async def test_fetch_all_secondary_concurrent_pages_deduped(scraper: CianScraper
|
||||||
assert 3 in fetch_pages or 3 in [
|
assert 3 in fetch_pages or 3 in [
|
||||||
p for p in fetch_pages
|
p for p in fetch_pages
|
||||||
], f"Страница 3 не была запрошена, fetch_pages={fetch_pages}"
|
], f"Страница 3 не была запрошена, fetch_pages={fetch_pages}"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_skip_buckets_skips_pagination_and_on_bucket(scraper: CianScraper) -> None:
|
||||||
|
"""skip_buckets: для скипнутого бакета пагинация не вызывается, on_bucket не вызывается."""
|
||||||
|
fetch_calls: list[tuple] = []
|
||||||
|
on_bucket_calls: list[str] = []
|
||||||
|
|
||||||
|
async def fake_fetch_page_html(
|
||||||
|
rooms: tuple, page: int, min_price: int | None, max_price: int | None
|
||||||
|
) -> str:
|
||||||
|
fetch_calls.append((rooms, page, min_price, max_price))
|
||||||
|
return f"<html>page={page}</html>"
|
||||||
|
|
||||||
|
def fake_extract_total_offers(html: str) -> int | None:
|
||||||
|
return 28 # один leaf-бакет, 1 страница
|
||||||
|
|
||||||
|
def fake_parse_serp_html(html: str) -> list[ScrapedLot]:
|
||||||
|
return [_make_lot("lot_1")]
|
||||||
|
|
||||||
|
def fake_on_bucket(bucket_key: str, lots: list[ScrapedLot]) -> None:
|
||||||
|
on_bucket_calls.append(bucket_key)
|
||||||
|
|
||||||
|
# Сформируем ключ того же бакета, который получится при walk (room1:0:200000000)
|
||||||
|
bucket_key_to_skip = f"room1:0:{_MAX_PRICE}"
|
||||||
|
skip_set = {bucket_key_to_skip}
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(scraper, "_fetch_page_html", side_effect=fake_fetch_page_html),
|
||||||
|
patch.object(scraper, "_extract_total_offers", side_effect=fake_extract_total_offers),
|
||||||
|
patch.object(scraper, "_parse_serp_html", side_effect=fake_parse_serp_html),
|
||||||
|
patch.object(scraper, "sleep_between_requests", new_callable=AsyncMock),
|
||||||
|
patch.object(scraper, "_rotate_ip", return_value=False),
|
||||||
|
):
|
||||||
|
lots = await scraper.fetch_all_secondary(
|
||||||
|
rooms_buckets=[(1,)],
|
||||||
|
price_cap_per_bucket=1400,
|
||||||
|
on_bucket=fake_on_bucket,
|
||||||
|
skip_buckets=skip_set,
|
||||||
|
)
|
||||||
|
|
||||||
|
# Probe (page=1) ВСЁ РАВНО вызывается (нужен для split-решения),
|
||||||
|
# но НЕ должно быть страниц пагинации (т.к. бакет в skip_buckets).
|
||||||
|
# Поскольку probe = page=1 и мы пропускаем ДО пагинации — fetch_page_html не должен вызываться.
|
||||||
|
assert (
|
||||||
|
len(fetch_calls) == 1
|
||||||
|
), f"Ожидался ровно 1 вызов (probe), получено {len(fetch_calls)}: {fetch_calls}"
|
||||||
|
assert fetch_calls[0][1] == 1, "Probe должен быть page=1"
|
||||||
|
# on_bucket не вызывается для skip-бакета
|
||||||
|
assert (
|
||||||
|
len(on_bucket_calls) == 0
|
||||||
|
), f"on_bucket не должен вызываться для skip-бакета, got: {on_bucket_calls}"
|
||||||
|
# lots пустой — ничего не собрали (пагинация пропущена)
|
||||||
|
assert len(lots) == 0, f"Ожидалось 0 лотов для skip-бакета, got {len(lots)}"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_on_bucket_receives_key_and_lots(scraper: CianScraper) -> None:
|
||||||
|
"""on_bucket получает (bucket_key, lots) — bucket_key в формате room_label:lo:hi."""
|
||||||
|
received: list[tuple[str, int]] = [] # (key, len(lots))
|
||||||
|
|
||||||
|
def fake_on_bucket(bucket_key: str, lots: list[ScrapedLot]) -> None:
|
||||||
|
received.append((bucket_key, len(lots)))
|
||||||
|
|
||||||
|
async def fake_fetch_page_html(
|
||||||
|
rooms: tuple, page: int, min_price: int | None, max_price: int | None
|
||||||
|
) -> str:
|
||||||
|
return f"<html>page={page}</html>"
|
||||||
|
|
||||||
|
def fake_extract_total_offers(html: str) -> int | None:
|
||||||
|
return 28 # 1 страница
|
||||||
|
|
||||||
|
def fake_parse_serp_html(html: str) -> list[ScrapedLot]:
|
||||||
|
import re
|
||||||
|
|
||||||
|
page_m = re.search(r"page=(\d+)", html)
|
||||||
|
page = int(page_m.group(1)) if page_m else 1
|
||||||
|
return [_make_lot(f"lot_p{page}_{i}") for i in range(5)]
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(scraper, "_fetch_page_html", side_effect=fake_fetch_page_html),
|
||||||
|
patch.object(scraper, "_extract_total_offers", side_effect=fake_extract_total_offers),
|
||||||
|
patch.object(scraper, "_parse_serp_html", side_effect=fake_parse_serp_html),
|
||||||
|
patch.object(scraper, "sleep_between_requests", new_callable=AsyncMock),
|
||||||
|
patch.object(scraper, "_rotate_ip", return_value=False),
|
||||||
|
):
|
||||||
|
await scraper.fetch_all_secondary(
|
||||||
|
rooms_buckets=[(2,)],
|
||||||
|
price_cap_per_bucket=1400,
|
||||||
|
on_bucket=fake_on_bucket,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert len(received) == 1, f"Ожидался 1 вызов on_bucket, получено {len(received)}"
|
||||||
|
key, count = received[0]
|
||||||
|
# Ключ должен быть «room2:0:<MAX_PRICE>»
|
||||||
|
assert key == f"room2:0:{_MAX_PRICE}", f"Неверный bucket_key: {key!r}"
|
||||||
|
assert count == 5, f"Ожидалось 5 лотов, получено {count}"
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue