3 Commits

Author SHA1 Message Date
qananasikq
3c3aea8eb3 tune vps config and update tests 2026-04-15 21:58:10 +03:00
qananasikq
37d01c4ba1 clean mapper comments 2026-04-15 21:58:03 +03:00
qananasikq
6165c65159 add segmented listing sync 2026-04-15 21:57:40 +03:00
14 changed files with 718 additions and 466 deletions

57
.env.example Normal file
View File

@@ -0,0 +1,57 @@
IAAI_HEADLESS=true
IAAI_LOG_LEVEL=INFO
IAAI_CAPTURE_SAME_ORIGIN_ONLY=true
IAAI_MAX_CAPTURED_REQUESTS=40
IAAI_MAX_CAPTURED_JSON_RESPONSES=20
IAAI_CARS_LISTING_URL=https://www.iaai.com/Vehiclelisting/Cars
IAAI_LISTING_SEGMENTS=auto
IAAI_MAX_PAGES_PER_RUN=999999
IAAI_MAX_VEHICLES_PER_RUN=999999
IAAI_PAGE_LINK_LIMIT=999999
IAAI_INCLUDE_PAGINATION=true
IAAI_COLLECT_CURRENT_PAGE_ONLY=false
IAAI_HUMAN_PACE_ENABLED=true
IAAI_PARALLEL_TABS=10
CELERY_BATCH_SIZE=100
IAAI_AFTER_LISTING_OPEN_MIN_S=0.2
IAAI_AFTER_LISTING_OPEN_MAX_S=0.5
IAAI_AFTER_FILTER_ACTION_MIN_S=0.2
IAAI_AFTER_FILTER_ACTION_MAX_S=0.5
IAAI_BEFORE_VEHICLE_OPEN_MIN_S=0.02
IAAI_BEFORE_VEHICLE_OPEN_MAX_S=0.08
IAAI_AFTER_VEHICLE_OPEN_MIN_S=0.02
IAAI_AFTER_VEHICLE_OPEN_MAX_S=0.08
IAAI_BETWEEN_VEHICLES_MIN_S=0.02
IAAI_BETWEEN_VEHICLES_MAX_S=0.08
IAAI_AFTER_PAGE_CHANGE_MIN_S=0.2
IAAI_AFTER_PAGE_CHANGE_MAX_S=0.5
IAAI_SYNC_ONLY_NEW=false
IAAI_TOKENS_FILE=/data/tokens.json
IAAI_RUNTIME_CONFIG_FILE=/app/runtime_config.json
IAAI_MAX_RETRIES=5
IAAI_RETRY_DELAY_SECONDS=4
IAAI_RETRY_BACKOFF_MULTIPLIER=2.0
IAAI_RETRY_JITTER_SECONDS=0.5
IAAI_TIMEOUT_MS=90000
IAAI_FALLBACK_NAV_TIMEOUT_MS=30000
IAAI_DATABASE_URL=postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper
IAAI_DATABASE_ECHO=false
IAAI_DATABASE_POOL_SIZE=5
IAAI_DATABASE_MAX_OVERFLOW=5
IAAI_REDIS_URL=redis://redis:6379/0
CELERY_BROKER_URL=redis://redis:6379/0
CELERY_RESULT_BACKEND=redis://redis:6379/0
CELERY_TASK_SOFT_TIME_LIMIT=86400
CELERY_TASK_TIME_LIMIT=86520
CELERY_BROKER_VISIBILITY_TIMEOUT=90000
CELERY_WORKER_CONCURRENCY=1
CELERY_BEAT_SYNC_INTERVAL_MINUTES=60
CELERY_BEAT_SYNC_LIMIT=0

3
.gitignore vendored
View File

@@ -3,8 +3,7 @@ __pycache__/
.coverage
coverage.xml
.env*
!.env
.env.example
!.env.example
.venv/
venv/
*.pyc

View File

@@ -15,6 +15,7 @@ x-app-env: &app-env
CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-200}
IAAI_PARALLEL_TABS: ${IAAI_PARALLEL_TABS:-40}
IAAI_BLOCK_RESOURCES: ${IAAI_BLOCK_RESOURCES:-true}
IAAI_LISTING_SEGMENTS: ${IAAI_LISTING_SEGMENTS:-auto}
IAAI_MAX_PAGES_PER_RUN: ${IAAI_MAX_PAGES_PER_RUN:-9999}
IAAI_MAX_VEHICLES_PER_RUN: ${IAAI_MAX_VEHICLES_PER_RUN:-50000}
IAAI_HUMAN_PACE_ENABLED: ${IAAI_HUMAN_PACE_ENABLED:-true}

View File

@@ -133,9 +133,9 @@ class ListingCollector:
return not old_first_href
def open_cars_listing(self, page: Page) -> None:
logger.info("Opening cars listing page: %s", self.settings.listing.cars_url)
url = self.settings.listing.cars_url
def open_cars_listing(self, page: Page, *, url_override: str | None = None) -> None:
url = url_override or self.settings.listing.cars_url
logger.info("Opening cars listing page: %s", url)
last_err = None
for attempt in range(3):
try:
@@ -202,16 +202,82 @@ class ListingCollector:
# Короткая пауза вместо длинного sleep.
time.sleep(1.0)
def apply_filters(self, page: Page, make: str | None = None, model: str | None = None) -> dict[str, str | None]:
applied = {"make": None, "model": None}
def apply_filters(
self,
page: Page,
make: str | None = None,
model: str | None = None,
year_min: int | None = None,
year_max: int | None = None,
) -> dict[str, str | int | None]:
applied: dict[str, str | int | None] = {"make": None, "model": None, "year_min": None, "year_max": None}
if make and self._try_fill_filter_input(page, ["input[placeholder*='Make']", "input[aria-label*='Make']"], make):
applied["make"] = make
self.pacer.after_filter_action()
if model and self._try_fill_filter_input(page, ["input[placeholder*='Model']", "input[aria-label*='Model']"], model):
applied["model"] = model
self.pacer.after_filter_action()
if year_min is not None or year_max is not None:
if self._apply_year_range(page, year_min, year_max):
applied["year_min"] = year_min
applied["year_max"] = year_max
self.pacer.after_filter_action()
return applied
def _apply_year_range(self, page: Page, year_min: int | None, year_max: int | None) -> bool:
"""Заполняет поля фильтра Year и нажимает Apply Year."""
if year_min is None and year_max is None:
return False
try:
success = page.evaluate(
"""([yearMin, yearMax]) => {
const inputs = Array.from(document.querySelectorAll('input'));
const yearInputs = inputs.filter(inp => {
const v = parseInt(inp.value, 10);
return !isNaN(v) && v >= 1900 && v <= 2100;
});
if (yearInputs.length < 2) return false;
yearInputs.sort((a, b) => parseInt(a.value) - parseInt(b.value));
const setVal = (el, val) => {
const setter = Object.getOwnPropertyDescriptor(
HTMLInputElement.prototype, 'value'
).set;
setter.call(el, String(val));
el.dispatchEvent(new Event('input', {bubbles: true}));
el.dispatchEvent(new Event('change', {bubbles: true}));
};
if (yearMin !== null) setVal(yearInputs[0], yearMin);
if (yearMax !== null) setVal(yearInputs[yearInputs.length - 1], yearMax);
const container = yearInputs[0].closest(
'[class*="filter"], [class*="year"], section, fieldset'
) || yearInputs[0].parentElement.parentElement;
if (container) {
const btn = Array.from(container.querySelectorAll(
'button, a, [role="button"], span[class*="apply"]'
)).find(el => /apply|\u043f\u0440\u0438\u043c\u0435\u043d/i.test(el.textContent));
if (btn) { btn.click(); return true; }
}
yearInputs[yearInputs.length - 1].dispatchEvent(
new KeyboardEvent('keydown', {
key: 'Enter', code: 'Enter', keyCode: 13, bubbles: true
})
);
return true;
}""",
[year_min, year_max],
)
if success:
try:
page.wait_for_load_state("domcontentloaded", timeout=15_000)
except Exception:
pass
self._wait_for_listing_content(page)
logger.info("Applied year range filter: %s%s", year_min, year_max)
return True
except Exception as exc:
logger.warning("Failed to apply year range filter: %s", exc)
return False
def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult:
# Считываем ссылки одним проходом по DOM.
self._accept_cookie_banner(page)

View File

@@ -1,3 +1,4 @@
import json
import os
from dataclasses import dataclass, field
from pathlib import Path
@@ -111,6 +112,79 @@ class ListingConfig:
# Порог раннего останова: если доля уже известных машин на странице >= этого значения,
# прекращаем листать — все новые машины уже найдены. 0 = отключено.
early_stop_threshold: float = _env_float("IAAI_EARLY_STOP_THRESHOLD", 0.8)
# Сегментация листинга по брендам для обхода лимита пагинации IAAI (~22 600 машин).
# JSON-массив объектов: [{"make":"TOYOTA"},{"make":"FORD"},...] или "auto" для авто-списка.
# Пустая строка = без сегментации (backward compatible).
listing_segments_json: str = _env_str("IAAI_LISTING_SEGMENTS", "")
# Список брендов IAAI для автоматической сегментации.
# Покрывает >99% автомобилей на сайте. Порядок: от крупных к мелким.
IAAI_DEFAULT_MAKES: tuple[str, ...] = (
"TOYOTA", "FORD", "CHEVROLET", "HONDA", "NISSAN", "HYUNDAI",
"KIA", "DODGE", "JEEP", "BMW", "MERCEDES-BENZ", "SUBARU",
"VOLKSWAGEN", "GMC", "MAZDA", "LEXUS", "CHRYSLER", "AUDI",
"RAM", "BUICK", "CADILLAC", "ACURA", "INFINITI", "LINCOLN",
"MITSUBISHI", "VOLVO", "JAGUAR", "LAND ROVER", "PORSCHE",
"MINI", "TESLA", "GENESIS", "FIAT", "ALFA ROMEO", "MASERATI",
"SCION", "PONTIAC", "SATURN", "MERCURY", "SAAB", "SUZUKI",
"OLDSMOBILE", "ISUZU", "HUMMER", "PLYMOUTH", "SMART",
"RIVIAN", "LUCID", "POLESTAR", "FERRARI", "LAMBORGHINI",
"BENTLEY", "ROLLS-ROYCE", "ASTON MARTIN", "MCLAREN", "LOTUS",
"MAYBACH", "FISKER", "GEO", "DAEWOO", "EAGLE",
)
# Пагинационный потолок IAAI: ~226 страниц × 100 = 22 600 результатов.
IAAI_PAGINATION_CEILING = 22_600
# Бренды, потенциально превышающие потолок пагинации — разбиваем по годам.
_LARGE_MAKES: frozenset[str] = frozenset({
"TOYOTA", "FORD", "CHEVROLET", "HONDA", "NISSAN", "HYUNDAI",
"KIA", "DODGE", "JEEP",
})
_YEAR_SPLITS: tuple[tuple[int, int], ...] = (
(1900, 2012),
(2013, 2019),
(2020, 2027),
)
def parse_listing_segments(raw: str) -> list[dict[str, str | int | None]]:
"""Парсит IAAI_LISTING_SEGMENTS в список сегментов.
Каждый сегмент — dict с ключами: make (str), year_min/year_max (int|None).
Специальное значение ``"auto"`` генерирует сегменты из IAAI_DEFAULT_MAKES.
Крупные бренды автоматически разбиваются по диапазонам годов.
"""
raw = raw.strip()
if not raw:
return []
if raw.lower() == "auto":
segments: list[dict[str, str | int | None]] = []
for m in IAAI_DEFAULT_MAKES:
if m in _LARGE_MAKES:
for yr_min, yr_max in _YEAR_SPLITS:
segments.append({"make": m, "year_min": yr_min, "year_max": yr_max})
else:
segments.append({"make": m, "year_min": None, "year_max": None})
return segments
try:
data = json.loads(raw)
except (json.JSONDecodeError, ValueError):
return []
if not isinstance(data, list):
return []
segments = []
for item in data:
if isinstance(item, str):
segments.append({"make": item.upper(), "year_min": None, "year_max": None})
elif isinstance(item, dict):
segments.append({
"make": str(item.get("make") or "").upper() or None,
"year_min": int(item["year_min"]) if item.get("year_min") is not None else None,
"year_max": int(item["year_max"]) if item.get("year_max") is not None else None,
})
return segments
# - Конфиг PostgreSQL (URL, пул соединений, pool_recycle)

View File

@@ -379,7 +379,7 @@ class CarMapper:
digest = hashlib.sha256(origin_id.encode()).digest()
alphabet = cls._PARSER_ID_ALPHABET
base = len(alphabet)
num = int.from_bytes(digest[:17], "big") # 17 bytes = 136 bits, enough for 22 chars
num = int.from_bytes(digest[:17], "big") # 17 байт = 136 бит, хватает на 22 символа
chars: list[str] = []
for _ in range(22):
num, idx = divmod(num, base)

View File

@@ -20,7 +20,7 @@ from playwright.sync_api import BrowserContext, Page, sync_playwright
from playwright.sync_api import TimeoutError as PlaywrightTimeoutError
from .browser import BrowserFactory, HumanPacer, NetworkCapture
from .core.config import Settings, settings
from .core.config import Settings, settings, parse_listing_segments
from .core.exceptions import AntiBotDetectedError, SiteStructureChangedError
from .core.logs import set_trace_id, setup_logging
from .core.retry import retryable
@@ -120,8 +120,6 @@ class IAAIScraper:
self.car_mapper = CarMapper()
self.persistence = PersistenceService(self.settings)
self._shutdown_requested = False
# HTTP-клиент для fast-path.
# При наличии прокси используем ProxyManager.
proxy_url = self.settings.proxy.server
if proxy_url:
_proxy_kwargs: dict = {
@@ -182,6 +180,13 @@ class IAAIScraper:
pass
finally:
self.playwright = None
if getattr(self, "_http_pool", None) is not None:
try:
self._http_pool.clear()
except Exception:
pass
finally:
self._http_pool = None
def _new_context(self) -> BrowserContext:
if self.browser is None:
@@ -199,21 +204,17 @@ class IAAIScraper:
return {"status": "ok", "database_url": self.settings.database.url}
def _warmup_visit(self, page: Page) -> None:
# Прогрев главной страницы.
try:
logger.info("Warmup: visiting homepage to pass anti-bot challenge...")
page.goto(self.settings.home_url, wait_until="commit", timeout=30_000)
# Ждём domcontentloaded вместо networkidle.
try:
page.wait_for_load_state("domcontentloaded", timeout=6_000)
except PlaywrightTimeoutError:
pass
# Короткая проверка, что страница стабилизировалась.
try:
page.wait_for_function("() => document.title && document.title.length > 3", timeout=4_000)
except PlaywrightTimeoutError:
pass
# Короткая пауза для установки cookies.
time.sleep(0.3)
logger.info("Warmup done: %s (title=%s)", page.url, page.title()[:50])
except Exception as e:
@@ -223,12 +224,10 @@ class IAAIScraper:
def _get_page_with_warmup(self) -> Page:
context = self._new_context()
page = context.new_page()
# Прогрев с первой антибот-проверкой.
self._warmup_visit(page)
return page
def _dedupe_urls(self, raw_urls: list[str]) -> list[str]:
# Нормализация и дедупликация URL.
vehicle_urls: list[str] = []
seen: set[str] = set()
for raw in raw_urls:
@@ -239,7 +238,6 @@ class IAAIScraper:
return vehicle_urls
def _filter_known_urls(self, vehicle_urls: list[str]) -> tuple[list[str], int]:
# Отсев уже известных URL.
url_to_origin_id = {
url: self._extract_db_origin_id_from_url(url)
for url in vehicle_urls
@@ -276,6 +274,14 @@ class IAAIScraper:
page_urls.append(normalized)
return page_urls
@staticmethod
def _build_segment_listing_url(base_url: str, make: str | None) -> str:
"""Построить URL листинга для сегмента. Make передаётся через query-параметр."""
if not make:
return base_url
sep = "&" if "?" in base_url else "?"
return f"{base_url}{sep}Make={make.replace(' ', '%20')}"
def _recover_empty_listing_page(
self,
page: Page,
@@ -284,11 +290,12 @@ class IAAIScraper:
all_raw_urls: list[str],
seen_urls: set[str],
all_listing_origin_urls: set[str] | None = None,
listing_url: str | None = None,
) -> tuple[ListingPageResult, list[str]]:
if page_number == 1:
logger.warning("Page 1 returned 0 links — retrying listing open...")
time.sleep(3)
self.listing_collector.open_cars_listing(page)
self.listing_collector.open_cars_listing(page, url_override=listing_url)
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
return page_result, self._extract_page_urls(page_result, all_raw_urls, seen_urls, all_listing_origin_urls)
@@ -309,10 +316,24 @@ class IAAIScraper:
target_page_number: int,
make: str | None,
model: str | None,
listing_url: str | None = None,
year_min: int | None = None,
year_max: int | None = None,
max_nav_pages: int = 10,
) -> Page:
# Ограничиваем глубину навигации: если до цели > max_nav_pages кликов — не пытаемся.
if target_page_number > max_nav_pages + 1:
raise RuntimeError(
f"Cannot resume at page {target_page_number}: "
f"exceeds max navigation depth ({max_nav_pages} pages)"
)
page = self._get_page_with_warmup()
self.listing_collector.open_cars_listing(page)
self.listing_collector.apply_filters(page, make=make, model=model)
try:
self.listing_collector.open_cars_listing(page, url_override=listing_url)
self.listing_collector.apply_filters(page, make=make, model=model, year_min=year_min, year_max=year_max)
except Exception:
page.close()
raise
for expected_page in range(2, target_page_number + 1):
if not self.listing_collector.go_to_next_page(page, expected_page_number=expected_page):
@@ -322,120 +343,6 @@ class IAAIScraper:
logger.warning("Listing resumed at page %d after recovery", target_page_number)
return page
def _collect_listing_iterative(
self,
*,
make: str | None = None,
model: str | None = None,
limit: int,
max_duration_seconds: float | None = None,
) -> tuple[list[str], list[str], dict, int]:
# Поэтапный сбор листинга до нужного лимита новых URL.
new_urls: list[str] = []
all_raw_urls: list[str] = []
seen: set[str] = set()
skipped_existing = 0
pages_info: list[dict] = []
max_pages = self.settings.listing.max_pages_per_run
started_at = time.perf_counter()
truncated_by_time_budget = False
page = self._get_page_with_warmup()
try:
self.listing_collector.open_cars_listing(page)
self.listing_collector.apply_filters(page, make=make, model=model)
for page_number in range(1, max_pages + 1):
if max_duration_seconds is not None and max_duration_seconds > 0:
elapsed = time.perf_counter() - started_at
if elapsed >= max_duration_seconds:
truncated_by_time_budget = True
logger.warning(
"Iterative listing stopped by time budget: page=%d elapsed=%.1fs budget=%.1fs",
page_number,
elapsed,
max_duration_seconds,
)
break
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
pages_info.append({
"page_number": page_result.page_number,
"links_found": len(page_result.vehicle_links),
})
# Собираем URL со страницы.
page_urls = self._extract_page_urls(page_result, all_raw_urls, seen)
if not page_urls:
page_result, page_urls = self._recover_empty_listing_page(
page,
page_number=page_number,
all_raw_urls=all_raw_urls,
seen_urls=seen,
)
if not page_urls and page_number > 1:
logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number)
page.close()
page = self._reopen_listing_and_resume(
target_page_number=page_number,
make=make,
model=model,
)
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
page_urls = self._extract_page_urls(page_result, all_raw_urls, seen)
if not page_urls:
logger.info("Page %d: 0 new links, stopping pagination", page_number)
break
# Фильтруем known.
fresh, page_skipped = self._filter_known_urls(page_urls)
skipped_existing += page_skipped
new_urls.extend(fresh)
logger.info(
"Page %d: %d links, %d new, %d known (total new: %d/%d)",
page_number, len(page_urls), len(fresh), page_skipped,
len(new_urls), limit,
)
# Стопаем только если набрали нужное количество И на этой странице уже нет новых.
# Если последняя страница дала новые — проверяем следующую (там могут быть ещё).
if len(new_urls) >= limit and len(fresh) == 0:
break
if len(new_urls) >= limit:
# Набрали достаточно, дальше не листаем
break
if not page_result.next_page_detected:
logger.info("No next page detected, stopping")
break
if not self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1):
logger.warning("Failed to navigate to page %d — reopening listing and resuming", page_number + 1)
page.close()
page = self._reopen_listing_and_resume(
target_page_number=page_number + 1,
make=make,
model=model,
)
finally:
page.close()
# Обрезаем до limit.
new_urls = new_urls[:limit]
listing = {
"status": "ok",
"listing_url": self.settings.listing.cars_url,
"pages_collected": len(pages_info),
"vehicles_collected": len(all_raw_urls),
"vehicle_urls": all_raw_urls,
"early_stopped": False,
"truncated_by_time_budget": truncated_by_time_budget,
"pages": pages_info,
}
return new_urls, all_raw_urls, listing, skipped_existing
def collect_listing(
self,
make: str | None = None,
@@ -472,6 +379,9 @@ class IAAIScraper:
started_at: float,
start_page: int = 1,
progress_callback: Callable[[int], None] | None = None,
listing_url: str | None = None,
year_min: int | None = None,
year_max: int | None = None,
) -> dict[str, Any]:
batch_size = self.settings.celery.batch_size
pending_urls: list[str] = []
@@ -500,9 +410,9 @@ class IAAIScraper:
except Exception as exc:
logger.warning("Could not load known origin_ids for early-stop: %s", exc)
# Совместимость: если only_new без явного limit, используем старую схему
# collect_listing -> filter_known -> batch upsert (это ожидают тесты и API-поток only_new).
if effective_only_new and (limit is None or limit <= 0):
# Совместимость: если only_new без limit и без сегментного URL — старая схема.
# При сегментированном скрапинге (listing_url/year фильтры) всегда streaming.
if effective_only_new and (limit is None or limit <= 0) and not listing_url and year_min is None and year_max is None:
listing_payload = self.collect_listing(
make=make,
model=model,
@@ -612,15 +522,20 @@ class IAAIScraper:
try:
if start_page <= 1:
assert page is not None
self.listing_collector.open_cars_listing(page)
applied_filters = self.listing_collector.apply_filters(page, make=make, model=model)
self.listing_collector.open_cars_listing(page, url_override=listing_url)
applied_filters = self.listing_collector.apply_filters(
page, make=make, model=model, year_min=year_min, year_max=year_max,
)
else:
page = self._reopen_listing_and_resume(
target_page_number=start_page,
make=make,
model=model,
listing_url=listing_url,
year_min=year_min,
year_max=year_max,
)
applied_filters = {"make": make, "model": model}
applied_filters = {"make": make, "model": model, "year_min": year_min, "year_max": year_max}
for page_number in range(start_page, max(start_page, self.settings.listing.max_pages_per_run) + 1):
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
@@ -643,14 +558,19 @@ class IAAIScraper:
all_raw_urls=all_raw_urls,
seen_urls=seen_urls,
all_listing_origin_urls=all_listing_origin_urls,
listing_url=listing_url,
)
if not page_urls and page_number > 1:
logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number)
page.close()
try:
page = self._reopen_listing_and_resume(
target_page_number=page_number,
make=make,
model=model,
listing_url=listing_url,
year_min=year_min,
year_max=year_max,
)
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
page_urls = self._extract_page_urls(
@@ -659,6 +579,13 @@ class IAAIScraper:
seen_urls,
all_listing_origin_urls,
)
except RuntimeError as resume_exc:
logger.warning(
"Cannot resume at page %d (%s) — stopping pagination for this segment",
page_number, resume_exc,
)
page = self._get_page_with_warmup()
page_urls = []
if not page_urls:
logger.info("Page %d: 0 new links, stopping pagination", page_number)
break
@@ -708,7 +635,7 @@ class IAAIScraper:
limit if limit is not None else "",
)
# Явный прогресс в логах даже при более высоком уровне логирования.
# Прогресс каждые 10 страниц на уровне WARNING.
if page_number % 10 == 0:
logger.warning(
"sync_listing progress: pages=%d, discovered=%d, upserted=%d, failed=%d, pending=%d",
@@ -743,11 +670,22 @@ class IAAIScraper:
if not self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1):
logger.warning("Failed to navigate to page %d — reopening listing and resuming", page_number + 1)
page.close()
try:
page = self._reopen_listing_and_resume(
target_page_number=page_number + 1,
make=make,
model=model,
listing_url=listing_url,
year_min=year_min,
year_max=year_max,
)
except RuntimeError as resume_exc:
logger.warning(
"Cannot resume at page %d (%s) — stopping pagination for this segment",
page_number + 1, resume_exc,
)
page = self._get_page_with_warmup()
break
if pending_urls:
_process_pending_batch(pending_urls, cars_upserted + cars_failed)
@@ -760,6 +698,7 @@ class IAAIScraper:
cars_failed,
)
finally:
if page is not None:
page.close()
return {
@@ -790,14 +729,12 @@ class IAAIScraper:
@retryable(max_attempts=3, jitter_seconds=0.25)
def scrape_vehicle_detail(self, vehicle_url: str):
# Открыть страницу, перехватить JSON, вернуть данные.
page = self._get_page()
try:
return self._scrape_on_page(page, vehicle_url)
finally:
page.close()
# Быстрый извлекатель данных из inline JSON.
_JS_EXTRACT = """
() => {
try {
@@ -867,7 +804,6 @@ class IAAIScraper:
@staticmethod
def _build_car_from_js(js_data: dict, vehicle_url: str) -> dict:
# Сбор vehicle_summary из JS-данных.
if not js_data or not js_data.get("ok"):
return {}
@@ -909,7 +845,6 @@ class IAAIScraper:
@staticmethod
def _build_payload_insights(vehicle_summary: dict) -> dict:
# Сбор payload_insights из vehicle_summary.
return {
"vehicle_core": vehicle_summary,
"pricing": {
@@ -929,7 +864,6 @@ class IAAIScraper:
trace_id = self._new_trace_id("scrape")
started_at = time.perf_counter()
# Блокируем тяжёлые ресурсы страницы.
if self.settings.block_resources:
BrowserFactory.enable_resource_blocking(page)
@@ -938,13 +872,11 @@ class IAAIScraper:
page.goto(vehicle_url, wait_until="commit", timeout=60_000)
# На VPS достаточно domcontentloaded.
try:
page.wait_for_load_state("domcontentloaded", timeout=5_000)
except PlaywrightTimeoutError:
pass
# Быстрый путь: читаем данные из inline JSON.
js_data: dict = {}
try:
js_data = page.evaluate(self._JS_EXTRACT) or {}
@@ -952,9 +884,7 @@ class IAAIScraper:
pass
if js_data.get("ok"):
# Успешное извлечение структуры авто.
vehicle_summary = self._build_car_from_js(js_data, vehicle_url)
# Быстрая проверка валидности данных.
has_identity = bool(vehicle_summary.get("make") or vehicle_summary.get("lot_number"))
if not has_identity:
raise SiteStructureChangedError(f"JS extraction returned no vehicle identity for {vehicle_url}")
@@ -982,7 +912,7 @@ class IAAIScraper:
save_to_json(network_dump, self.settings.raw_output_json)
return result
# Резервный путь: полный парсинг страницы.
# Резервный путь: полный HTML-парсинг.
try:
page.wait_for_selector("#VehicleDetailViewModel, .veh-details, .vehicle-details, [data-uname='vehicleDetailPage']", timeout=400)
except PlaywrightTimeoutError:
@@ -1029,7 +959,6 @@ class IAAIScraper:
if pos < 0:
return None
# Ищем начало внешнего JSON-объекта.
script_start = html.rfind("<script", 0, pos)
if script_start < 0:
return None
@@ -1040,7 +969,6 @@ class IAAIScraper:
if content_start < 0:
return None
# Быстрый поиск конца JSON через raw_decode.
decoder = json.JSONDecoder()
try:
obj, _ = decoder.raw_decode(html, content_start)
@@ -1104,7 +1032,6 @@ class IAAIScraper:
}
def _scrape_via_context_request(self, vehicle_url: str) -> CarRecord:
# Быстрый запрос через context.request.
if not self.context:
self._new_context()
assert self.context is not None
@@ -1119,7 +1046,6 @@ class IAAIScraper:
html = response.text()
# Пытаемся извлечь JSON напрямую из HTML.
js_data = self._extract_iaai_json_from_html(html, vehicle_url)
if not js_data:
raise RuntimeError(f"HTTP fast-path missing embedded JSON for {vehicle_url}")
@@ -1200,7 +1126,6 @@ class IAAIScraper:
return CarRecord.model_validate(db_record.model_dump(mode="json"))
def sync_vehicle(self, vehicle_url: str, lane: str = "iaai"):
# Скрапинг и upsert одного авто.
trace_id = self._new_trace_id("sync-vehicle")
started_at = time.perf_counter()
self.persistence.create_tables()
@@ -1247,12 +1172,12 @@ class IAAIScraper:
error_summary=error_summary,
)
# ── Tunables for sync_batch ──
_HTTP_MICRO_BATCH = 25 # URLs per micro-batch (avoid mass rate-limit)
_HTTP_MAX_RETRIES = 2 # Retries per URL before giving up to browser
_HTTP_RETRY_DELAYS = (0.4, 1.0) # Backoff between retries
_FALLBACK_PARALLEL_PAGES = 4 # Concurrent browser tabs for fallback
_FALLBACK_NAV_TIMEOUT_MS = 8000 # Reduced from 15 000
# ── Настройки sync_batch ──
_HTTP_MICRO_BATCH = 25 # URL за микро-батч
_HTTP_MAX_RETRIES = 2 # Повторов на URL до перехода в браузер
_HTTP_RETRY_DELAYS = (0.4, 1.0) # Задержки между повторами
_FALLBACK_PARALLEL_PAGES = 4 # Параллельных вкладок для fallback
_FALLBACK_NAV_TIMEOUT_MS = 8000 # Таймаут навигации fallback
def _browser_fallback_parallel(
self,
@@ -1260,13 +1185,11 @@ class IAAIScraper:
total: int,
n_pages: int,
) -> dict:
# Обработка fallback URL через браузер.
records: list[CarRecord] = []
failures: list[dict[str, str]] = []
cars_failed = 0
protection_events = 0
# Деление URL по страницам.
slices: list[list[tuple[str, int]]] = [[] for _ in range(n_pages)]
for i, item in enumerate(fallback_urls):
slices[i % n_pages].append(item)
@@ -1384,7 +1307,7 @@ class IAAIScraper:
"protection_events": local_protection,
}
# Один page обрабатываем в главном потоке.
# Одна страница — в главном потоке.
if n_pages == 1:
res = _process_slice(slices[0] if slices else [])
records.extend(res["records"])
@@ -1392,7 +1315,7 @@ class IAAIScraper:
cars_failed += res["cars_failed"]
protection_events += res["protection_events"]
else:
# Несколько страниц — параллельно через потоки (каждый поток со своей Page).
# Несколько страниц — параллельно, каждый поток со своей Page.
with ThreadPoolExecutor(max_workers=n_pages) as executor:
futures = [executor.submit(_process_slice, s) for s in slices if s]
for fut in as_completed(futures):
@@ -1415,7 +1338,6 @@ class IAAIScraper:
lane: str = "iaai_cars",
parallel_tabs: int | None = None,
) -> dict:
# Пакетный скрапинг списка URL.
trace_id = self._new_trace_id("sync-batch")
started_at = time.perf_counter()
num_workers = parallel_tabs or self.settings.parallel_tabs
@@ -1433,14 +1355,13 @@ class IAAIScraper:
total = len(vehicle_urls)
logger.info("sync_batch: %d vehicles, %d parallel workers", total, num_workers)
# Подготовка для HTTP fast-path: cookies + user-agent из Playwright сессии.
cookie_header = self._cookie_header_for_context()
user_agent = (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:128.0) "
"Gecko/20100101 Firefox/128.0"
)
# ── Phase 1: HTTP fast-path — все URLs одним ThreadPoolExecutor ──
# ── Фаза 1: HTTP fast-path ──
fallback_urls: list[tuple[str, int]] = []
def _fetch_one_with_retry(idx_url: tuple[int, str]) -> tuple[int, CarRecord | Exception]:
@@ -1459,7 +1380,6 @@ class IAAIScraper:
time.sleep(self._HTTP_RETRY_DELAYS[min(attempt, len(self._HTTP_RETRY_DELAYS) - 1)])
return idx, last_exc # type: ignore[return-value]
# Запускаем все URLs сразу — один пул потоков для минимального времени ожидания.
indexed_urls = list(enumerate(vehicle_urls))
with ThreadPoolExecutor(max_workers=min(num_workers, total)) as executor:
for idx, result in executor.map(_fetch_one_with_retry, indexed_urls):
@@ -1483,7 +1403,7 @@ class IAAIScraper:
http_successes, http_fallbacks, time.perf_counter() - started_at,
)
# ── Phase 2: browser fallback (single-threaded — Playwright sync API is not thread-safe) ──
# ── Фаза 2: браузерный fallback ──
if fallback_urls:
logger.info(
"Fallback browser mode: %d/%d URLs, single page",
@@ -1495,9 +1415,8 @@ class IAAIScraper:
protection_events += fb_results["protection_events"]
failures.extend(fb_results["failures"])
# ── Phase 3: single DB flush ──
# ── Фаза 3: запись в БД ──
if records:
# Дедупликация перед записью.
seen_origins: set[str] = set()
unique_records: list[CarRecord] = []
for rec in records:
@@ -1540,9 +1459,11 @@ class IAAIScraper:
only_new: bool | None = None,
start_page: int = 1,
progress_callback: Callable[[int], None] | None = None,
listing_url: str | None = None,
year_min: int | None = None,
year_max: int | None = None,
):
# Листинг + sync всех найденных машин.
# Применяем runtime_config как дефолты (CLI/API аргументы имеют приоритет).
# runtime_config — дефолты; CLI/API аргументы приоритетнее.
rc = self.runtime_config.sync
if limit is None and rc.limit is not None:
limit = rc.limit
@@ -1576,6 +1497,9 @@ class IAAIScraper:
started_at=started_at,
start_page=max(1, int(start_page)),
progress_callback=progress_callback,
listing_url=listing_url,
year_min=year_min,
year_max=year_max,
)
listing = stream_result["listing"]
total = stream_result["total"]
@@ -1588,8 +1512,7 @@ class IAAIScraper:
logger.info("Streaming sync processed %d vehicles", total)
# Помечаем авто как проданные, если они исчезли из листинга.
# Только если сканирование было полным (не ограниченным limit/only_new/early_stop).
# Помечаем проданные авто, исчезнувшие из листинга (только при полном скане).
is_partial_scan = (
effective_only_new
or (limit is not None and limit > 0)
@@ -1610,7 +1533,6 @@ class IAAIScraper:
if not failures:
failures.append({"vehicle_url": "collect_listing", "error": str(exc)})
logger.error("sync_listing failed: %s (partial progress: %d upserted)", exc, cars_upserted)
# Не пробрасываем — partial progress уже записан в БД через finish_sync_run ниже.
finally:
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
error_summary = "; ".join(item["error"] for item in failures[:10]) if failures else None
@@ -1632,7 +1554,9 @@ class IAAIScraper:
"trace_id": trace_id,
"status": status,
"run_id": run_id,
"full_scan_completed": (not is_partial_scan) and status == "success",
# partial_success допустим — отдельные машины могли не спарситься,
# это не повод повторять весь bootstrap.
"full_scan_completed": (not is_partial_scan) and status in ("success", "partial_success"),
"only_new_effective": effective_only_new,
"listing": listing,
"cars_upserted": cars_upserted,
@@ -1644,9 +1568,132 @@ class IAAIScraper:
"failures": failures,
}
def sync_listing_segmented(
self,
segments: list[dict[str, Any]],
lane: str = "iaai_cars",
only_new: bool | None = None,
start_segment: int = 0,
start_page: int = 1,
progress_callback: Callable[[int, int], None] | None = None,
) -> dict[str, Any]:
"""Итеративный sync_listing по списку сегментов (бренд / бренд+годы).
Args:
segments: список dict с ключами make, year_min, year_max.
start_segment: индекс сегмента для resume (0-based).
start_page: страница внутри start_segment для resume.
progress_callback: вызывается (segment_index, page_number) после каждой страницы.
"""
trace_id = self._new_trace_id("sync-segmented")
started_at = time.perf_counter()
base_url = self.settings.listing.cars_url
total_cars_upserted = 0
total_cars_failed = 0
total_images_upserted = 0
total_skipped = 0
total_discovered = 0
all_failures: list[dict[str, str]] = []
segment_results: list[dict[str, Any]] = []
completed_all = True
logger.warning(
"Starting segmented sync: %d segments, resume from segment=%d page=%d",
len(segments), start_segment, start_page,
)
for seg_idx in range(start_segment, len(segments)):
seg = segments[seg_idx]
seg_make = seg.get("make")
seg_year_min = seg.get("year_min")
seg_year_max = seg.get("year_max")
seg_url = self._build_segment_listing_url(base_url, seg_make) if seg_make else None
seg_start_page = start_page if seg_idx == start_segment else 1
seg_label = f"{seg_make or 'ALL'}"
if seg_year_min is not None or seg_year_max is not None:
seg_label += f" ({seg_year_min}-{seg_year_max})"
logger.warning(
"Segment %d/%d: %s (start_page=%d)",
seg_idx + 1, len(segments), seg_label, seg_start_page,
)
def _seg_progress(page_number: int, _si=seg_idx) -> None:
if progress_callback:
progress_callback(_si, page_number)
try:
result = self.sync_listing(
make=None if seg_url else seg_make,
model=None,
lane=lane,
only_new=only_new,
start_page=seg_start_page,
progress_callback=_seg_progress,
listing_url=seg_url,
year_min=seg_year_min,
year_max=seg_year_max,
)
total_cars_upserted += result.get("cars_upserted", 0)
total_cars_failed += result.get("cars_failed", 0)
total_images_upserted += result.get("images_upserted", 0)
total_skipped += result.get("skipped_existing", 0)
total_discovered += result.get("listing", {}).get("vehicles_collected", 0)
all_failures.extend(result.get("failures", []))
segment_results.append({
"segment": seg,
"segment_index": seg_idx,
"status": result.get("status"),
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"vehicles_collected": result.get("listing", {}).get("vehicles_collected", 0),
})
logger.warning(
"Segment %d/%d done: %s → upserted=%d, failed=%d, collected=%d",
seg_idx + 1, len(segments), seg_label,
result.get("cars_upserted", 0),
result.get("cars_failed", 0),
result.get("listing", {}).get("vehicles_collected", 0),
)
except Exception as exc:
logger.error("Segment %d/%d failed: %s%s", seg_idx + 1, len(segments), seg_label, exc)
all_failures.append({"vehicle_url": f"segment_{seg_idx}_{seg_label}", "error": str(exc)})
completed_all = False
# Продолжаем оставшиеся сегменты — одна ошибка не должна убивать весь прогон
continue
elapsed = round(time.perf_counter() - started_at, 3)
status = "success" if not all_failures else "partial_success" if total_cars_upserted else "failed"
logger.warning(
"Segmented sync done: %d/%d segments, upserted=%d, failed=%d, discovered=%d, elapsed=%.1fs",
len(segment_results), len(segments),
total_cars_upserted, total_cars_failed, total_discovered, elapsed,
)
return {
"trace_id": trace_id,
"status": status,
# Bootstrap считается завершённым если все сегменты пройдены,
# даже если часть машин failed (они будут обновлены в следующих циклах).
"full_scan_completed": completed_all,
"segments_total": len(segments),
"segments_completed": len(segment_results),
"cars_upserted": total_cars_upserted,
"cars_failed": total_cars_failed,
"images_upserted": total_images_upserted,
"skipped_existing": total_skipped,
"total_discovered": total_discovered,
"elapsed_seconds": elapsed,
"failures": all_failures,
"segment_results": segment_results,
}
def run_scheduled(self) -> None:
# Резерв для standalone-режима.
# В текущей архитектуре планирование выполняется через Celery beat + worker/tasks.py.
"""Standalone-планировщик (в продакшене используется Celery beat)."""
def _handle_shutdown(signum, frame):
logger.info("Received signal %s, shutting down gracefully...", signum)
self._shutdown_requested = True

View File

@@ -11,7 +11,7 @@ from billiard.exceptions import SoftTimeLimitExceeded
from celery import shared_task
from redis import Redis
from ..core.config import Settings
from ..core.config import Settings, parse_listing_segments
from ..scraper import IAAIScraper
from ..storage.db import PersistenceService
@@ -166,7 +166,7 @@ def _release_lock_if_owner(redis_client: Redis, key: str, owner_token: str) -> N
logger.warning("Failed to release lock %s", key, exc_info=True)
def _has_running_sync_listing_tasks(celery_app) -> bool:
def _has_running_sync_listing_tasks(celery_app, *, exclude_task_id: str | None = None) -> bool:
try:
inspector = celery_app.control.inspect(timeout=1.0)
snapshots = [
@@ -182,12 +182,17 @@ def _has_running_sync_listing_tasks(celery_app) -> bool:
for entries in snapshot.values():
for entry in entries or []:
task_name = str(entry.get("name") or entry.get("request", {}).get("name") or "")
if task_name == SYNC_LISTING_TASK_NAME:
if task_name != SYNC_LISTING_TASK_NAME:
continue
# Исключаем текущую задачу — она не считается "другой запущенной"
entry_id = str(entry.get("id") or entry.get("request", {}).get("id") or "")
if exclude_task_id and entry_id == exclude_task_id:
continue
return True
return False
def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app) -> bool:
def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app, *, current_task_id: str | None = None) -> bool:
try:
owner_token = redis_client.get(SYNC_LISTING_LOCK_KEY)
if not owner_token:
@@ -196,7 +201,7 @@ def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app) -> bool:
logger.warning("Failed to read sync listing lock before cleanup", exc_info=True)
return False
if _has_running_sync_listing_tasks(celery_app):
if _has_running_sync_listing_tasks(celery_app, exclude_task_id=current_task_id):
logger.info("sync_listing lock preserved: active task still detected")
return False
@@ -256,6 +261,7 @@ def _save_sync_checkpoint(
make: str | None,
model: str | None,
lane: str,
segment_index: int | None = None,
) -> None:
payload = {
"status": "in_progress",
@@ -264,6 +270,7 @@ def _save_sync_checkpoint(
"make": make,
"model": model,
"lane": lane,
"segment_index": segment_index,
"updated_at": int(time.time()),
}
try:
@@ -384,7 +391,7 @@ def sync_listing_task(
lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl)
if not lock_acquired:
orphan_cleared = _clear_orphan_sync_listing_lock(redis_client, self.app)
orphan_cleared = _clear_orphan_sync_listing_lock(redis_client, self.app, current_task_id=task_id)
if orphan_cleared:
lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl)
@@ -409,16 +416,20 @@ def sync_listing_task(
checkpoint_make = checkpoint.get("make")
checkpoint_model = checkpoint.get("model")
checkpoint_lane = checkpoint.get("lane")
checkpoint_segment = checkpoint.get("segment_index")
same_scope = (
checkpoint_make == make
and checkpoint_model == model
and checkpoint_lane == lane
)
if checkpoint_page > 0 and same_scope:
# Для сегментированного режима: совпадение по lane + наличие segment_index
is_segmented_checkpoint = checkpoint_segment is not None and checkpoint_lane == lane
if checkpoint_page > 0 and (same_scope or is_segmented_checkpoint):
resume_from_page = checkpoint_page + 1
logger.warning(
"Resuming sync_listing from page %d using checkpoint",
"Resuming sync_listing from page %d (segment=%s) using checkpoint",
resume_from_page,
checkpoint_segment,
)
elif checkpoint_page > 0:
logger.info("Ignoring stale checkpoint due to different sync parameters")
@@ -442,8 +453,38 @@ def sync_listing_task(
lock_ttl,
)
self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id})
# Определяем сегменты из конфига.
settings = Settings()
segments = parse_listing_segments(settings.listing.listing_segments_json)
use_segmented = bool(segments) and make is None and model is None
resume_from_segment = 0
if use_segmented and checkpoint and str(checkpoint.get("status") or "") == "in_progress":
cp_segment = checkpoint.get("segment_index")
if cp_segment is not None and int(cp_segment) >= 0:
resume_from_segment = int(cp_segment)
# resume_from_page уже вычислен выше
def _job():
with IAAIScraper() as scraper:
if use_segmented:
return scraper.sync_listing_segmented(
segments=segments,
lane=lane,
only_new=effective_only_new,
start_segment=resume_from_segment,
start_page=resume_from_page,
progress_callback=lambda seg_idx, page_number: _save_sync_checkpoint(
redis_client,
task_id=task_id,
page_number=page_number,
make=None,
model=None,
lane=lane,
segment_index=seg_idx,
),
)
return scraper.sync_listing(
make=make,
model=model,

View File

@@ -30,23 +30,19 @@ class _FakePage:
class TestListingUnit(unittest.TestCase):
def test_has_next_page_true_for_known_selector(self) -> None:
page = _FakePage({"a[aria-label*='Next']": 1})
self.assertTrue(ListingCollector._has_next_page(page))
def test_has_next_page_false_when_no_selectors(self) -> None:
page = _FakePage({})
self.assertFalse(ListingCollector._has_next_page(page))
def test_has_next_page_true_for_numeric_pagination_fallback(self) -> None:
def test_pagination_detection_and_page_number(self) -> None:
# Есть кнопка Next → True.
self.assertTrue(ListingCollector._has_next_page(_FakePage({"a[aria-label*='Next']": 1})))
# Нет селекторов → False.
self.assertFalse(ListingCollector._has_next_page(_FakePage({})))
# Числовая пагинация через JS → True.
page = _FakePage({})
page._evaluate_result = True
self.assertTrue(ListingCollector._has_next_page(page))
def test_get_current_page_number_from_text_counter(self) -> None:
page = _FakePage({})
page._evaluate_values = [5]
self.assertEqual(ListingCollector._get_current_page_number(page), 5)
# Номер текущей страницы.
page2 = _FakePage({})
page2._evaluate_values = [5]
self.assertEqual(ListingCollector._get_current_page_number(page2), 5)
def test_extract_vehicle_links_from_html_finds_detail_urls(self) -> None:
collector = ListingCollector(Settings(), HumanPacer(Settings()))

View File

@@ -41,30 +41,27 @@ class TestCarMapper(unittest.TestCase):
self.assertFalse(record.is_damaged)
def test_unknown_and_empty_values_fallbacks(self) -> None:
def test_unknown_empty_values_and_normalization(self) -> None:
record = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/999~US",
vehicle_summary={"make": " ", "model": None, "drive": "???", "gearbox": "unknown"},
payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}},
)
self.assertEqual(record.brand, "UNKNOWN")
self.assertEqual(record.model, "UNKNOWN")
self.assertIsNone(record.steering_wheel if record.steering_wheel not in {"LEFT", None} else None)
self.assertEqual(record.drive, "NA")
self.assertEqual(record.gearbox, "NA")
def test_mapper_handles_case_and_spaces(self) -> None:
record = self.mapper.map_to_car_record(
# Нормализация регистра и пробелов.
record2 = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/888~US",
vehicle_summary={"make": "Honda", "model": "Civic", "drive": " Front Wheel Drive ", "gearbox": " AUTOMATIC "},
payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}},
)
self.assertEqual(record2.drive, "FWD")
self.assertEqual(record2.gearbox, "AT")
self.assertEqual(record.drive, "FWD")
self.assertEqual(record.gearbox, "AT")
def test_price_parsing_dirty_formats(self) -> None:
def test_price_and_currency_parsing(self) -> None:
record = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/777~US",
vehicle_summary={"make": "Toyota", "model": "Corolla"},
@@ -76,11 +73,9 @@ class TestCarMapper(unittest.TestCase):
"images": {},
},
)
self.assertEqual(record.price, 5200)
def test_currency_detection_from_symbol(self) -> None:
record = self.mapper.map_to_car_record(
record2 = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/778~US",
vehicle_summary={"make": "Toyota", "model": "Corolla"},
payload_insights={
@@ -91,9 +86,8 @@ class TestCarMapper(unittest.TestCase):
"images": {},
},
)
self.assertEqual(record.currency, "EUR")
self.assertEqual(record.price, 4500)
self.assertEqual(record2.currency, "EUR")
self.assertEqual(record2.price, 4500)
if __name__ == "__main__":

View File

@@ -9,7 +9,7 @@ class TestVehicleParserUnit(unittest.TestCase):
def setUp(self) -> None:
self.parser = VehicleParser()
def test_parse_dom_key_value_pairs_extracts_known_fields(self) -> None:
def test_parse_dom_pairs_and_title(self) -> None:
dom_text = """
Stock #:
45089484
@@ -23,41 +23,27 @@ class TestVehicleParserUnit(unittest.TestCase):
self.assertEqual(result.get("primary_damage"), "Front End")
self.assertEqual(result.get("odometer"), "50,123 mi (Actual)")
def test_parse_title_for_year_make_model(self) -> None:
parsed = self.parser._parse_title_for_year_make_model("2014 TOYOTA CAMRY for sale", "")
self.assertEqual(parsed["year"], "2014")
self.assertEqual(parsed["make"], "TOYOTA")
self.assertEqual(parsed["model"], "CAMRY")
def test_extract_image_urls_filters_other_vehicle(self) -> None:
vehicle_url = "https://www.iaai.com/VehicleDetail/45089484~US"
payloads = [
{
"imageUrls": [
"https://vis.iaai.com/resizer?imageKeys=45089484~SID1&width=845&height=633",
"https://vis.iaai.com/resizer?imageKeys=99999999~SID2&width=845&height=633",
]
}
]
urls = self.parser._extract_image_urls(payloads, "", vehicle_url)
self.assertEqual(len(urls), 1)
self.assertIn("45089484", urls[0])
def test_extract_image_urls_deduplicates_same_url(self) -> None:
def test_extract_image_urls_filters_and_deduplicates(self) -> None:
vehicle_url = "https://www.iaai.com/VehicleDetail/45089484~US"
payloads = [{"imageUrls": [
"https://vis.iaai.com/resizer?imageKeys=45089484~SID1&width=845&height=633",
"https://vis.iaai.com/resizer?imageKeys=45089484~SID1&width=845&height=633",
"https://vis.iaai.com/resizer?imageKeys=99999999~SID2&width=845&height=633",
]}]
urls = self.parser._extract_image_urls(payloads, "", vehicle_url)
self.assertEqual(len(urls), 1)
self.assertIn("45089484", urls[0])
def test_dom_hints_detect_captcha_and_antibot(self) -> None:
def test_dom_hints_and_access_notes_detect_antibot(self) -> None:
hints = self.parser._dom_hints("Please verify you are human. CAPTCHA. Incapsula access denied.")
self.assertTrue(hints["has_captcha_text"])
self.assertTrue(hints["has_antibot_text"])
def test_access_notes_reflect_antibot_signals(self) -> None:
summary = {"vin": "", "image_urls": [], "note": "Incapsula access denied. Verify you are human."}
notes = self.parser._build_access_notes(summary, [])
self.assertTrue(notes["possible_captcha"])

View File

@@ -28,104 +28,129 @@ class TestScraperSync(unittest.TestCase):
s.database.url = "sqlite://"
return IAAIScraper(s)
def test_sync_vehicle_uses_db_record_without_remapping(self) -> None:
def test_sync_vehicle_uses_db_record(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=1)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.upsert_car = MagicMock(return_value={"action": "inserted", "images_upserted": 0})
scraper.scrape_vehicle_detail = MagicMock(return_value={"db_record": make_db_record("111")})
scraper.car_mapper.map_to_car_record = MagicMock(side_effect=AssertionError("should not be called"))
result = scraper.sync_vehicle("https://www.iaai.com/VehicleDetail/111~US")
self.assertEqual(result["status"], "success")
self.assertIn("trace_id", result)
self.assertIn("elapsed_seconds", result)
self.assertEqual(scraper.persistence.upsert_car.call_count, 1)
scraper.persistence.upsert_car.assert_called_once()
def test_sync_listing_uses_db_record_without_remapping(self) -> None:
def test_only_new_routing_legacy_and_streaming(self) -> None:
# Legacy path: only_new без listing_url.
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=2)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.get_existing_urls_and_ids = MagicMock(return_value=(set(), set()))
scraper.collect_listing = MagicMock(return_value={"vehicle_urls": ["https://www.iaai.com/VehicleDetail/222~US"]})
scraper.sync_batch = MagicMock(return_value={
"cars_upserted": 1, "cars_failed": 0, "images_upserted": 1, "failures": [],
scraper.persistence.get_existing_urls_and_ids = MagicMock(return_value=(
{"https://www.iaai.com/VehicleDetail/111~US"}, {"iaai:222"},
))
scraper.collect_listing = MagicMock(return_value={
"vehicle_urls": [
"https://www.iaai.com/VehicleDetail/111~US",
"https://www.iaai.com/VehicleDetail/222~US",
"https://www.iaai.com/VehicleDetail/333~US",
]
})
result = scraper.sync_listing()
scraper.sync_batch = MagicMock(return_value={
"cars_upserted": 1, "cars_failed": 0, "images_upserted": 0, "failures": [],
})
result = scraper.sync_listing(only_new=True)
self.assertEqual(result["skipped_existing"], 2)
self.assertEqual(result["cars_upserted"], 1)
self.assertEqual(result["cars_failed"], 0)
self.assertIn("trace_id", result)
self.assertIn("elapsed_seconds", result)
scraper.sync_batch.assert_called_once()
def test_sync_listing_respects_limit(self) -> None:
# Streaming path: only_new + listing_url.
scraper2 = self._make_scraper()
scraper2.persistence.create_tables = MagicMock()
scraper2.persistence.start_sync_run = MagicMock(return_value=6)
scraper2.persistence.finish_sync_run = MagicMock()
scraper2.collect_listing = MagicMock(side_effect=AssertionError("legacy path should not be used"))
scraper2._sync_listing_streaming = MagicMock(return_value={
"listing": {"vehicles_collected": 10, "early_stopped": False, "truncated_by_time_budget": False},
"total": 10, "skipped_existing": 0, "cars_upserted": 10, "cars_failed": 0,
"images_upserted": 20, "failures": [], "all_listing_origin_urls": set(),
})
result2 = scraper2.sync_listing(
only_new=True,
listing_url="https://www.iaai.com/Vehiclelisting/Cars?Make=TOYOTA",
year_min=2020, year_max=2027,
)
self.assertEqual(result2["cars_upserted"], 10)
scraper2._sync_listing_streaming.assert_called_once()
scraper2.collect_listing.assert_not_called()
def test_segmented_sync_calls_per_segment_and_resumes(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=3)
scraper.persistence.finish_sync_run = MagicMock()
# Проверка пути only_new с limit.
scraper._collect_listing_iterative = MagicMock(return_value=(
["https://www.iaai.com/VehicleDetail/222~US"],
["https://www.iaai.com/VehicleDetail/222~US",
"https://www.iaai.com/VehicleDetail/333~US"],
{"vehicle_urls": [], "pages_collected": 1, "early_stopped": False},
0,
))
scraper.sync_batch = MagicMock(return_value={
"cars_upserted": 1, "cars_failed": 0, "images_upserted": 1, "failures": [],
})
call_args_log: list[dict] = []
def _fake_sync_listing(**kwargs):
call_args_log.append(kwargs)
return {
"status": "success", "cars_upserted": 5, "cars_failed": 0,
"images_upserted": 10, "skipped_existing": 0,
"listing": {"vehicles_collected": 50}, "failures": [],
}
scraper.sync_listing(limit=1)
# Должен уйти только один URL.
scraper.sync_batch.assert_called_once()
batch_urls = scraper.sync_batch.call_args[0][0]
self.assertEqual(len(batch_urls), 1)
def test_sync_listing_only_new_filters_existing_by_url_and_origin_id(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=5)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.get_existing_urls_and_ids = MagicMock(return_value=(
{"https://www.iaai.com/VehicleDetail/111~US"},
{"iaai:222"},
))
scraper.collect_listing = MagicMock(return_value={
"vehicle_urls": [
"https://www.iaai.com/VehicleDetail/111~US", # exists by URL
"https://www.iaai.com/VehicleDetail/222~US", # exists by ID
"https://www.iaai.com/VehicleDetail/333~US", # new
scraper.sync_listing = MagicMock(side_effect=_fake_sync_listing)
segments = [
{"make": "TOYOTA", "year_min": 2020, "year_max": 2027},
{"make": "FORD", "year_min": None, "year_max": None},
{"make": "HONDA", "year_min": None, "year_max": None},
]
})
# Возвращаем результат для одного нового авто.
scraper.sync_batch = MagicMock(return_value={
"cars_upserted": 1, "cars_failed": 0, "images_upserted": 0, "failures": [],
})
result = scraper.sync_listing(only_new=True)
# Resume: пропускаем TOYOTA, начинаем с FORD на стр. 5.
result = scraper.sync_listing_segmented(
segments=segments, start_segment=1, start_page=5,
)
self.assertEqual(result["skipped_existing"], 2)
self.assertEqual(result["cars_upserted"], 1)
# В batch должен попасть только новый URL.
scraper.sync_batch.assert_called_once()
batch_urls = scraper.sync_batch.call_args[0][0]
self.assertEqual(len(batch_urls), 1)
self.assertIn("333", batch_urls[0])
scraper.persistence.get_existing_urls_and_ids.assert_called_once()
self.assertEqual(result["segments_completed"], 2) # FORD + HONDA
self.assertEqual(result["cars_upserted"], 10)
# FORD: start_page=5, HONDA: start_page=1.
self.assertEqual(call_args_log[0]["start_page"], 5)
self.assertEqual(call_args_log[1]["start_page"], 1)
# URL содержит бренд.
self.assertIn("FORD", call_args_log[0]["listing_url"])
self.assertIn("HONDA", call_args_log[1]["listing_url"])
def test_build_segment_listing_url(self) -> None:
base = "https://www.iaai.com/Vehiclelisting/Cars"
self.assertEqual(
IAAIScraper._build_segment_listing_url(base, "TOYOTA"),
"https://www.iaai.com/Vehiclelisting/Cars?Make=TOYOTA",
)
self.assertEqual(
IAAIScraper._build_segment_listing_url(base, "LAND ROVER"),
"https://www.iaai.com/Vehiclelisting/Cars?Make=LAND%20ROVER",
)
self.assertEqual(IAAIScraper._build_segment_listing_url(base, None), base)
self.assertEqual(IAAIScraper._build_segment_listing_url(base, ""), base)
def test_guard_and_protection_detection(self) -> None:
with self.assertRaises(AntiBotDetectedError):
IAAIScraper._raise_if_blocked_or_incomplete(
{"dom_hints": {"has_captcha_text": True, "has_antibot_text": False},
"access_notes": {}, "vehicle_summary": {}},
"https://www.iaai.com/VehicleDetail/999~US",
)
with self.assertRaises(SiteStructureChangedError):
IAAIScraper._raise_if_blocked_or_incomplete(
{"dom_hints": {"has_captcha_text": False, "has_antibot_text": False},
"access_notes": {"possible_captcha": False, "possible_antibot": False},
"vehicle_summary": {}},
"https://www.iaai.com/VehicleDetail/999~US",
)
self.assertTrue(IAAIScraper._is_protection_or_network_error(RuntimeError("NS_ERROR_NET_INTERRUPT")))
self.assertTrue(IAAIScraper._is_protection_or_network_error(RuntimeError("captcha challenge")))
self.assertFalse(IAAIScraper._is_protection_or_network_error(RuntimeError("plain validation error")))
def test_close_resets_browser_state(self) -> None:
scraper = self._make_scraper()
@@ -134,85 +159,34 @@ class TestScraperSync(unittest.TestCase):
scraper.context = MagicMock()
scraper.browser = MagicMock()
scraper.playwright = MagicMock()
scraper.close()
http_pool.clear.assert_called_once()
self.assertIsNone(scraper._http_pool)
self.assertIsNone(scraper.context)
self.assertIsNone(scraper.browser)
self.assertIsNone(scraper.playwright)
def test_guard_raises_on_antibot_signals(self) -> None:
with self.assertRaises(AntiBotDetectedError):
IAAIScraper._raise_if_blocked_or_incomplete(
{
"dom_hints": {"has_captcha_text": True, "has_antibot_text": False},
"access_notes": {},
"vehicle_summary": {},
},
"https://www.iaai.com/VehicleDetail/999~US",
)
def test_guard_raises_on_empty_vehicle_page(self) -> None:
with self.assertRaises(SiteStructureChangedError):
IAAIScraper._raise_if_blocked_or_incomplete(
{
"dom_hints": {"has_captcha_text": False, "has_antibot_text": False},
"access_notes": {"possible_captcha": False, "possible_antibot": False},
"vehicle_summary": {},
},
"https://www.iaai.com/VehicleDetail/999~US",
)
def test_is_protection_or_network_error_detects_known_signals(self) -> None:
self.assertTrue(IAAIScraper._is_protection_or_network_error(RuntimeError("NS_ERROR_NET_INTERRUPT")))
self.assertTrue(IAAIScraper._is_protection_or_network_error(RuntimeError("captcha challenge")))
self.assertFalse(IAAIScraper._is_protection_or_network_error(RuntimeError("plain validation error")))
def test_recover_empty_listing_page_reload_recovers_links(self) -> None:
def test_recover_empty_listing_page(self) -> None:
scraper = self._make_scraper()
page = MagicMock()
page_result = SimpleNamespace(
vehicle_links=[SimpleNamespace(href="https://www.iaai.com/VehicleDetail/123~US")],
)
scraper.listing_collector.collect_current_page = MagicMock(return_value=page_result)
all_raw_urls: list[str] = []
seen_urls: set[str] = set()
recovered_result, recovered_urls = scraper._recover_empty_listing_page(
page,
page_number=5,
all_raw_urls=all_raw_urls,
seen_urls=seen_urls,
)
page.reload.assert_called_once()
self.assertIs(recovered_result, page_result)
self.assertEqual(recovered_urls, ["https://www.iaai.com/VehicleDetail/123~US"])
self.assertEqual(all_raw_urls, ["https://www.iaai.com/VehicleDetail/123~US"])
def test_recover_empty_listing_page_reopens_listing_for_first_page(self) -> None:
scraper = self._make_scraper()
page = MagicMock()
page_result = SimpleNamespace(
vehicle_links=[SimpleNamespace(href="https://www.iaai.com/VehicleDetail/456~US")],
)
scraper.listing_collector.open_cars_listing = MagicMock()
scraper.listing_collector.collect_current_page = MagicMock(return_value=page_result)
recovered_result, recovered_urls = scraper._recover_empty_listing_page(
page,
page_number=1,
all_raw_urls=[],
seen_urls=set(),
# Страница > 1: reload.
_, urls = scraper._recover_empty_listing_page(
page, page_number=5, all_raw_urls=[], seen_urls=set(),
)
page.reload.assert_called_once()
self.assertEqual(len(urls), 1)
scraper.listing_collector.open_cars_listing.assert_called_once_with(page)
# Страница 1: переоткрытие листинга.
page.reset_mock()
_, urls = scraper._recover_empty_listing_page(
page, page_number=1, all_raw_urls=[], seen_urls=set(),
)
scraper.listing_collector.open_cars_listing.assert_called_once()
page.reload.assert_not_called()
self.assertIs(recovered_result, page_result)
self.assertEqual(recovered_urls, ["https://www.iaai.com/VehicleDetail/456~US"])
if __name__ == "__main__":

View File

@@ -3,28 +3,72 @@ from __future__ import annotations
import unittest
from iaai_scraper.core.utils import deep_find_key
from iaai_scraper.core.config import parse_listing_segments, IAAI_DEFAULT_MAKES, _LARGE_MAKES, _YEAR_SPLITS
class TestDeepFindKey(unittest.TestCase):
def test_finds_key_in_nested_structure(self) -> None:
def test_finds_nested_and_respects_depth(self) -> None:
payload = {
"root": {
"target": "a",
"nested": [{"target": "b"}, {"x": 1}],
}
}
self.assertEqual(deep_find_key(payload, {"target"}), ["a", "b"])
found = deep_find_key(payload, {"target"})
self.assertEqual(found, ["a", "b"])
deep_payload = {"l1": {"l2": {"l3": {"target": "value"}}}}
self.assertEqual(deep_find_key(deep_payload, {"target"}, max_depth=2), [])
self.assertEqual(deep_find_key(deep_payload, {"target"}, max_depth=8), ["value"])
def test_respects_max_depth(self) -> None:
payload = {"l1": {"l2": {"l3": {"target": "value"}}}}
found_too_shallow = deep_find_key(payload, {"target"}, max_depth=2)
found_enough_depth = deep_find_key(payload, {"target"}, max_depth=8)
class TestParseListingSegments(unittest.TestCase):
def test_empty_and_invalid_return_empty(self) -> None:
self.assertEqual(parse_listing_segments(""), [])
self.assertEqual(parse_listing_segments(" "), [])
self.assertEqual(parse_listing_segments("invalid"), [])
self.assertEqual(found_too_shallow, [])
self.assertEqual(found_enough_depth, ["value"])
def test_auto_segments_complete_coverage(self) -> None:
"""auto: все бренды покрыты, крупные разбиты по годам, годы без дыр."""
segs = parse_listing_segments("auto")
# Точное число сегментов.
expected = len(_LARGE_MAKES) * len(_YEAR_SPLITS) + (len(IAAI_DEFAULT_MAKES) - len(_LARGE_MAKES))
self.assertEqual(len(segs), expected)
# Все бренды из списка присутствуют.
makes_in_segments = {s["make"] for s in segs}
for make in IAAI_DEFAULT_MAKES:
self.assertIn(make, makes_in_segments)
# Крупные бренды разбиты на 3 сегмента с годами.
toyota = [s for s in segs if s["make"] == "TOYOTA"]
self.assertEqual(len(toyota), 3)
for s in toyota:
self.assertIsNotNone(s["year_min"])
# Мелкие бренды без годов.
lexus = [s for s in segs if s["make"] == "LEXUS"]
self.assertEqual(len(lexus), 1)
self.assertIsNone(lexus[0]["year_min"])
# Годовые диапазоны покрывают 1950-2027.
years = set()
for yr_min, yr_max in _YEAR_SPLITS:
years.update(range(yr_min, yr_max + 1))
for year in range(1950, 2027):
self.assertIn(year, years)
def test_json_input_formats(self) -> None:
# Массив строк.
segs = parse_listing_segments('["toyota", "ford"]')
self.assertEqual(len(segs), 2)
self.assertEqual(segs[0]["make"], "TOYOTA")
# Массив объектов с годами.
segs = parse_listing_segments('[{"make":"BMW","year_min":2020,"year_max":2025}]')
self.assertEqual(segs[0]["make"], "BMW")
self.assertEqual(segs[0]["year_min"], 2020)
self.assertEqual(segs[0]["year_max"], 2025)
if __name__ == "__main__":

View File

@@ -8,32 +8,22 @@ from iaai_scraper.worker import tasks
class TestWorkerTaskLockHelpers(unittest.TestCase):
def test_acquire_lock_returns_true_on_success(self) -> None:
def test_lock_acquire_refresh_release(self) -> None:
redis_client = MagicMock()
# Acquire.
redis_client.set.return_value = True
acquired = tasks._acquire_lock(redis_client, "lock:key", "owner-token", 120)
self.assertTrue(acquired)
self.assertTrue(tasks._acquire_lock(redis_client, "lock:key", "owner-token", 120))
redis_client.set.assert_called_once_with("lock:key", "owner-token", nx=True, ex=120)
def test_refresh_lock_if_owner_extends_ttl(self) -> None:
redis_client = MagicMock()
# Refresh.
redis_client.eval.return_value = 1
self.assertTrue(tasks._refresh_lock_if_owner(redis_client, "lock:key", "owner-token", 120))
refreshed = tasks._refresh_lock_if_owner(redis_client, "lock:key", "owner-token", 120)
self.assertTrue(refreshed)
redis_client.eval.assert_called_once()
def test_release_lock_if_owner_uses_owner_token(self) -> None:
redis_client = MagicMock()
# Release.
redis_client.eval.reset_mock()
tasks._release_lock_if_owner(redis_client, "lock:key", "owner-token")
redis_client.eval.assert_called_once()
args = redis_client.eval.call_args[0]
self.assertEqual(args[1], 1)
self.assertEqual(args[2], "lock:key")
self.assertEqual(args[3], "owner-token")
@@ -62,6 +52,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
patch.object(tasks, "_is_full_scan_done", return_value=True), \
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
patch.object(tasks.sync_listing_task, "update_state"), \
patch.object(tasks, "_run_browser_job", return_value={
"run_id": 7,
"cars_upserted": 2,
@@ -90,7 +81,8 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
heartbeat_thread.join.assert_called_once()
release_lock.assert_called_once()
def test_clear_orphan_sync_listing_lock_deletes_when_no_tasks_running(self) -> None:
def test_clear_orphan_sync_listing_lock(self) -> None:
# Нет запущенных задач → удаляет.
redis_client = MagicMock()
redis_client.get.return_value = "owner-token"
redis_client.ttl.return_value = 120
@@ -101,27 +93,21 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
inspector.scheduled.return_value = {"worker@node": []}
celery_app.control.inspect.return_value = inspector
cleared = tasks._clear_orphan_sync_listing_lock(redis_client, celery_app)
self.assertTrue(cleared)
self.assertTrue(tasks._clear_orphan_sync_listing_lock(redis_client, celery_app))
redis_client.delete.assert_called_once_with(tasks.SYNC_LISTING_LOCK_KEY)
def test_clear_orphan_sync_listing_lock_keeps_when_task_detected(self) -> None:
redis_client = MagicMock()
redis_client.get.return_value = "owner-token"
celery_app = MagicMock()
inspector = MagicMock()
inspector.active.return_value = {
"worker@node": [{"name": tasks.SYNC_LISTING_TASK_NAME}]
}
inspector.reserved.return_value = {"worker@node": []}
inspector.scheduled.return_value = {"worker@node": []}
celery_app.control.inspect.return_value = inspector
# Задача активна → не удаляет.
redis_client2 = MagicMock()
redis_client2.get.return_value = "owner-token"
inspector2 = MagicMock()
inspector2.active.return_value = {"worker@node": [{"name": tasks.SYNC_LISTING_TASK_NAME}]}
inspector2.reserved.return_value = {"worker@node": []}
inspector2.scheduled.return_value = {"worker@node": []}
celery_app2 = MagicMock()
celery_app2.control.inspect.return_value = inspector2
cleared = tasks._clear_orphan_sync_listing_lock(redis_client, celery_app)
self.assertFalse(cleared)
redis_client.delete.assert_not_called()
self.assertFalse(tasks._clear_orphan_sync_listing_lock(redis_client2, celery_app2))
redis_client2.delete.assert_not_called()
def test_sync_listing_task_recovers_orphan_lock_and_runs(self) -> None:
with patch.object(tasks, "_get_persistence") as get_persistence, \
@@ -131,6 +117,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
patch.object(tasks, "_is_full_scan_done", return_value=True), \
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
patch.object(tasks.sync_listing_task, "update_state"), \
patch.object(tasks, "_run_browser_job", return_value={
"run_id": 9,
"cars_upserted": 3,
@@ -159,28 +146,22 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
self.assertEqual(acquire_lock.call_count, 2)
release_lock.assert_called_once()
def test_sync_listing_checkpoint_roundtrip(self) -> None:
def test_sync_listing_checkpoint_save_load_and_resume(self) -> None:
# Roundtrip: save → load.
redis_client = MagicMock()
storage: dict[str, str] = {}
redis_client.set.side_effect = lambda key, value: storage.__setitem__(key, value)
redis_client.get.side_effect = lambda key: storage.get(key)
tasks._save_sync_checkpoint(
redis_client,
task_id="task-1",
page_number=12,
make=None,
model=None,
lane="iaai_cars",
redis_client, task_id="task-1", page_number=12,
make=None, model=None, lane="iaai_cars",
)
checkpoint = tasks._load_sync_checkpoint(redis_client)
self.assertIsNotNone(checkpoint)
self.assertEqual(checkpoint["last_successful_page"], 12)
self.assertEqual(checkpoint["status"], "in_progress")
def test_sync_listing_task_resumes_from_checkpoint_page(self) -> None:
# Resume: задача стартует со страницы checkpoint + 1.
with patch.object(tasks, "_get_persistence") as get_persistence, \
patch.object(tasks, "_get_redis") as get_redis, \
patch.object(tasks, "_acquire_lock", return_value=True), \
@@ -188,32 +169,24 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
patch.object(tasks, "_clear_sync_checkpoint") as clear_checkpoint, \
patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()):
persistence = MagicMock()
get_persistence.return_value = persistence
redis_client = MagicMock()
redis_client.get.side_effect = lambda key: json.dumps({
"status": "in_progress",
"last_successful_page": 9,
"make": None,
"model": None,
"lane": "iaai_cars",
patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \
patch.object(tasks.sync_listing_task, "update_state"), \
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=[]):
get_persistence.return_value = MagicMock()
redis_client2 = MagicMock()
redis_client2.get.side_effect = lambda key: json.dumps({
"status": "in_progress", "last_successful_page": 9,
"make": None, "model": None, "lane": "iaai_cars",
}) if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
get_redis.return_value = redis_client
get_redis.return_value = redis_client2
stop_event = MagicMock()
heartbeat_thread = MagicMock()
start_heartbeat.return_value = (stop_event, heartbeat_thread)
sync_listing_mock = MagicMock(return_value={
"run_id": 11,
"status": "success",
"full_scan_completed": True,
"cars_upserted": 1,
"cars_failed": 0,
"images_upserted": 0,
"skipped_existing": 0,
"elapsed_seconds": 1.0,
"failures": [],
"run_id": 11, "status": "success", "full_scan_completed": True,
"cars_upserted": 1, "cars_failed": 0, "images_upserted": 0,
"skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [],
})
scraper_ctx = MagicMock()
scraper_ctx.__enter__.return_value.sync_listing = sync_listing_mock