6 Commits

26 changed files with 4001 additions and 315 deletions

10
.gitignore vendored
View File

@@ -19,4 +19,12 @@ build/
artifacts/
celerybeat-schedule*
tokens_data/
tokens.json
tokens.json
# Local deployment/debug artifacts
/*.tgz
/*.tar.gz
/*.zip
/NOW()
/_speed_check.sql
/.env.bak*

View File

@@ -13,6 +13,10 @@ RUN pip install --no-cache-dir uv \
&& pip install --no-cache-dir -r /tmp/requirements.txt \
&& pip install --no-cache-dir playwright-stealth
# Speed-up libs: orjson (fast JSON parse), brotli (HTTP br decompress).
# Pinned outside uv.lock to avoid lock-regen churn.
RUN pip install --no-cache-dir "orjson>=3.10.0" "brotli>=1.1.0"
# Install both Chromium and Firefox. Chromium is the default engine in Docker
# because it works more reliably with the current IAAI listing page.
RUN python -m playwright install chromium firefox

View File

@@ -112,7 +112,32 @@ docker compose up -d
`entrypoint.sh` поднимает SOCKS5→HTTP proxy bridge (если задан `SOCKS5_PROXY_HOST`) и запускает `Xvfb` только для `worker` и CLI scraping-команд.
По умолчанию `beat` запускает сбор листинга **раз в 1 час** (`CELERY_BEAT_SYNC_INTERVAL_MINUTES=60`).
По умолчанию `beat` запускает синхронизацию **раз в 1 час** (`CELERY_BEAT_SYNC_INTERVAL_MINUTES=60`).
Текущая стратегия устойчива к багам пагинации IAAI:
- **первый запуск** делает полный bootstrap через `sitemap`;
- если bootstrap не завершился, автоматически продолжается с чекпоинта;
- **после первого полного прохода** каждый час выполняется проход по `sitemap`.
По умолчанию включён режим **rolling refresh**:
- берутся все URL из sitemap,
- новые URL добавляются сразу,
- исчезнувшие URL помечаются как `is_sold=true`,
- уже записанные авто обновляются **батчами по кругу**, а не все разом.
Это даёт более стабильную нагрузку и снижает риск не уложиться в час.
Доступные режимы:
- `IAAI_HOURLY_MODE=rolling_refresh` — рекомендуемый продовый режим;
- `IAAI_HOURLY_MODE=full_refresh` — перепарсить все активные авто за один hourly цикл;
- `IAAI_HOURLY_MODE=diff` — только новые авто + sold.
За счёт этого система:
- не зависит от page 39 / 70 / 100,
- не уходит в бесконечный цикл пагинации,
- даёт стабильный hourly refresh уже записанных авто без монолитного полного hourly прохода,
- остаётся быстрой и устойчивой при hourly sync.
Swagger-документация: `http://localhost:8000/docs`

View File

@@ -0,0 +1,77 @@
# Добавление таблиц для новой архитектуры ingestion: candidates, raw snapshots, parse results, retry queue
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
revision: str = "004_add_ingestion_tables"
down_revision: Union[str, None] = "003_add_composite_indexes"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
# Таблица vehicle_candidates
op.create_table(
"vehicle_candidates",
sa.Column("id", sa.BigInteger().with_variant(sa.Integer(), "sqlite"), primary_key=True, autoincrement=True),
sa.Column("url", sa.String(), nullable=False, unique=True),
sa.Column("discovered_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("status", sa.String(20), nullable=False, server_default="pending"),
sa.Column("priority", sa.Integer(), nullable=False, server_default="0"),
sa.Column("last_attempt_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("attempts", sa.Integer(), nullable=False, server_default="0"),
)
op.create_index("ix_vehicle_candidates_url", "vehicle_candidates", ["url"])
op.create_index("ix_vehicle_candidates_status_discovered", "vehicle_candidates", ["status", "discovered_at"])
# Таблица vehicle_raw_snapshots
op.create_table(
"vehicle_raw_snapshots",
sa.Column("id", sa.BigInteger().with_variant(sa.Integer(), "sqlite"), primary_key=True, autoincrement=True),
sa.Column("candidate_id", sa.Integer(), sa.ForeignKey("vehicle_candidates.id", ondelete="CASCADE"), nullable=False),
sa.Column("captured_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("method", sa.String(20), nullable=False),
sa.Column("success", sa.Boolean(), nullable=False),
sa.Column("raw_data", sa.Text(), nullable=True),
sa.Column("error_message", sa.Text(), nullable=True),
)
op.create_index("ix_vehicle_raw_snapshots_candidate_id", "vehicle_raw_snapshots", ["candidate_id"])
op.create_index("ix_vehicle_raw_snapshots_captured_at", "vehicle_raw_snapshots", ["captured_at"])
# Таблица vehicle_parse_results
op.create_table(
"vehicle_parse_results",
sa.Column("id", sa.BigInteger().with_variant(sa.Integer(), "sqlite"), primary_key=True, autoincrement=True),
sa.Column("snapshot_id", sa.Integer(), sa.ForeignKey("vehicle_raw_snapshots.id", ondelete="CASCADE"), nullable=False),
sa.Column("parsed_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("success", sa.Boolean(), nullable=False),
sa.Column("parsed_data", sa.Text(), nullable=True),
sa.Column("error_message", sa.Text(), nullable=True),
)
op.create_index("ix_vehicle_parse_results_snapshot_id", "vehicle_parse_results", ["snapshot_id"])
op.create_index("ix_vehicle_parse_results_parsed_at", "vehicle_parse_results", ["parsed_at"])
# Таблица vehicle_retry_queue
op.create_table(
"vehicle_retry_queue",
sa.Column("id", sa.BigInteger().with_variant(sa.Integer(), "sqlite"), primary_key=True, autoincrement=True),
sa.Column("candidate_id", sa.Integer(), sa.ForeignKey("vehicle_candidates.id", ondelete="CASCADE"), nullable=False),
sa.Column("reason", sa.String(50), nullable=False),
sa.Column("retry_at", sa.DateTime(timezone=True), nullable=False),
sa.Column("attempts", sa.Integer(), nullable=False, server_default="0"),
sa.Column("max_attempts", sa.Integer(), nullable=False, server_default="3"),
)
op.create_index("ix_vehicle_retry_queue_candidate_id", "vehicle_retry_queue", ["candidate_id"])
op.create_index("ix_vehicle_retry_queue_retry_at", "vehicle_retry_queue", ["retry_at"])
def downgrade() -> None:
op.drop_table("vehicle_retry_queue")
op.drop_table("vehicle_parse_results")
op.drop_table("vehicle_raw_snapshots")
op.drop_table("vehicle_candidates")

60
deploy_vps.sh Normal file
View File

@@ -0,0 +1,60 @@
#!/usr/bin/env bash
set -euo pipefail
APP_DIR="/opt/iaai_scraper_project"
REPO_URL="${1:-}"
if [[ -z "${REPO_URL}" ]]; then
echo "Usage: ./deploy_vps.sh <git-repo-url>"
exit 1
fi
export DEBIAN_FRONTEND=noninteractive
apt-get update
apt-get install -y --no-install-recommends \
ca-certificates \
curl \
git \
gnupg \
lsb-release
install -m 0755 -d /etc/apt/keyrings
if [[ ! -f /etc/apt/keyrings/docker.gpg ]]; then
curl -fsSL https://download.docker.com/linux/ubuntu/gpg | gpg --dearmor -o /etc/apt/keyrings/docker.gpg
fi
chmod a+r /etc/apt/keyrings/docker.gpg
ARCH="$(dpkg --print-architecture)"
CODENAME="$(. /etc/os-release && echo "$VERSION_CODENAME")"
echo \
"deb [arch=${ARCH} signed-by=/etc/apt/keyrings/docker.gpg] https://download.docker.com/linux/ubuntu ${CODENAME} stable" \
> /etc/apt/sources.list.d/docker.list
apt-get update
apt-get install -y --no-install-recommends \
docker-ce \
docker-ce-cli \
containerd.io \
docker-buildx-plugin \
docker-compose-plugin
mkdir -p /opt
if [[ -d "${APP_DIR}/.git" ]]; then
git -C "${APP_DIR}" fetch --all --prune
git -C "${APP_DIR}" reset --hard origin/HEAD
else
rm -rf "${APP_DIR}"
git clone "${REPO_URL}" "${APP_DIR}"
fi
cd "${APP_DIR}"
if [[ ! -f .env ]]; then
cp .env.vps.example .env
fi
docker compose down --remove-orphans || true
docker compose up -d --build postgres redis migrate api worker beat
docker compose ps

View File

@@ -5,13 +5,13 @@ x-app-env: &app-env
CELERY_BROKER_URL: ${CELERY_BROKER_URL:-redis://redis:6379/0}
CELERY_RESULT_BACKEND: ${CELERY_RESULT_BACKEND:-redis://redis:6379/0}
IAAI_DATABASE_POOL_RECYCLE_SECONDS: ${IAAI_DATABASE_POOL_RECYCLE_SECONDS:-1800}
CELERY_TASK_SOFT_TIME_LIMIT: ${CELERY_TASK_SOFT_TIME_LIMIT:-3300}
CELERY_TASK_TIME_LIMIT: ${CELERY_TASK_TIME_LIMIT:-3600}
CELERY_BROKER_VISIBILITY_TIMEOUT: ${CELERY_BROKER_VISIBILITY_TIMEOUT:-7200}
CELERY_TASK_SOFT_TIME_LIMIT: ${CELERY_TASK_SOFT_TIME_LIMIT:-900}
CELERY_TASK_TIME_LIMIT: ${CELERY_TASK_TIME_LIMIT:-1200}
CELERY_BROKER_VISIBILITY_TIMEOUT: ${CELERY_BROKER_VISIBILITY_TIMEOUT:-2400}
# 0 => без лимита (в коде интерпретируется как None).
# После первичного полного прохода hourly-режим должен успевать за всеми новыми авто.
CELERY_BEAT_SYNC_LIMIT: ${CELERY_BEAT_SYNC_LIMIT:-0}
CELERY_WORKER_MAX_TASKS_PER_CHILD: ${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5}
CELERY_WORKER_MAX_TASKS_PER_CHILD: ${CELERY_WORKER_MAX_TASKS_PER_CHILD:-3}
CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-200}
IAAI_PARALLEL_TABS: ${IAAI_PARALLEL_TABS:-40}
IAAI_BLOCK_RESOURCES: ${IAAI_BLOCK_RESOURCES:-true}
@@ -134,6 +134,8 @@ services:
container_name: iaai-worker
restart: unless-stopped
init: true
mem_limit: ${IAAI_WORKER_MEM_LIMIT:-2g}
memswap_limit: ${IAAI_WORKER_MEM_LIMIT:-2g}
depends_on:
migrate:
condition: service_completed_successfully
@@ -142,9 +144,9 @@ services:
stop_grace_period: 60s
command: >
celery -A iaai_scraper.worker.celery_app worker
--loglevel=info --concurrency=${CELERY_WORKER_CONCURRENCY:-4} --pool=prefork
--loglevel=info --concurrency=${CELERY_WORKER_CONCURRENCY:-2} --pool=prefork
--pidfile=/tmp/celery-worker.pid
-Q scraping --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5}
-Q scraping --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-3}
healthcheck:
test: ["CMD-SHELL", "test -f /tmp/celery-worker.pid && kill -0 $(cat /tmp/celery-worker.pid)"]
interval: 60s

View File

@@ -3,13 +3,12 @@ import re
import time
from dataclasses import asdict, dataclass, field
from typing import Any
from urllib.parse import urljoin
from urllib.parse import parse_qsl, urlencode, urljoin, urlparse, urlunparse
from playwright.sync_api import Page
from .pace import HumanPacer
from ..core.config import Settings
from ..core.utils import first_non_empty
logger = logging.getLogger("iaai_scraper.listing")
VEHICLE_HREF_RE = re.compile(r"/VehicleDetail/(\d+)(?:~[A-Z]{2})?", re.IGNORECASE)
@@ -44,6 +43,7 @@ class ListingPageResult:
class ListingCollector:
_DEEP_PAGINATION_DIRECT_ONLY_FROM_PAGE = 40
_NEXT_PAGE_SELECTORS: tuple[str, ...] = (
"a[aria-label*='Next']",
"button[aria-label*='Next']",
@@ -103,7 +103,33 @@ class ListingCollector:
return None
@staticmethod
def _wait_for_navigation_result(page: Page, old_first_href: str, expected_page_number: int | None) -> bool:
def _get_vehicle_link_fingerprint(page: Page, limit: int = 5) -> tuple[str, ...]:
try:
values = page.evaluate(
"""
(limit) => {
return Array.from(document.querySelectorAll("a[href*='/VehicleDetail/'], a[href*='/vehicledetail/'], a[href*='VehicleDetail'], a[href*='vehicledetail']"))
.map((el) => (el.getAttribute('href') || '').trim())
.filter(Boolean)
.slice(0, limit);
}
""",
limit,
)
if not isinstance(values, list):
return tuple()
return tuple(str(value) for value in values if value)
except Exception:
return tuple()
@staticmethod
def _wait_for_navigation_result(
page: Page,
old_first_href: str,
expected_page_number: int | None,
*,
old_fingerprint: tuple[str, ...] = (),
) -> bool:
if old_first_href:
try:
page.wait_for_function(
@@ -117,9 +143,12 @@ class ListingCollector:
except Exception:
pass
new_fingerprint = ListingCollector._get_vehicle_link_fingerprint(page)
if expected_page_number is not None:
current_page = ListingCollector._get_current_page_number(page)
if current_page == expected_page_number:
if current_page == expected_page_number and (
not old_fingerprint or (new_fingerprint and new_fingerprint != old_fingerprint)
):
return True
try:
@@ -127,11 +156,81 @@ class ListingCollector:
except Exception:
pass
if not new_fingerprint:
new_fingerprint = ListingCollector._get_vehicle_link_fingerprint(page)
current_page = ListingCollector._get_current_page_number(page)
if new_fingerprint and old_fingerprint and new_fingerprint != old_fingerprint:
return current_page is None or expected_page_number is None or current_page == expected_page_number
if expected_page_number is not None and current_page == expected_page_number:
return True
return not old_fingerprint
return not old_first_href
return not old_first_href and (not old_fingerprint or bool(new_fingerprint))
def _extract_next_page_href(self, page: Page, expected_page_number: int | None = None) -> str | None:
try:
href = page.evaluate(
"""
(expectedPageNumber) => {
const visible = (el) => !!(el && (el.offsetWidth || el.offsetHeight || el.getClientRects().length));
const disabled = (el) => {
if (!el) return true;
const cls = (el.getAttribute('class') || '').toLowerCase();
const ariaDisabled = (el.getAttribute('aria-disabled') || '').toLowerCase();
return el.hasAttribute('disabled') || ariaDisabled === 'true' || cls.includes('disabled');
};
const controls = Array.from(document.querySelectorAll('a,button,[role="button"]'))
.filter((el) => visible(el) && !disabled(el));
const cleanHref = (el) => {
const href = (el?.getAttribute('href') || '').trim();
if (!href || href.startsWith('javascript:') || href.startsWith('#')) {
return '';
}
return href;
};
if (expectedPageNumber !== null && expectedPageNumber !== undefined) {
const numeric = controls.find((el) => {
const text = (el.textContent || '').trim();
return /^\d+$/.test(text) && parseInt(text, 10) === expectedPageNumber;
});
const numericHref = cleanHref(numeric);
if (numericHref) {
return numericHref;
}
}
const explicitNext = controls.find((el) => {
const text = (el.textContent || '').trim().toLowerCase();
const aria = (el.getAttribute('aria-label') || '').trim().toLowerCase();
const title = (el.getAttribute('title') || '').trim().toLowerCase();
const rel = (el.getAttribute('rel') || '').trim().toLowerCase();
const cls = (el.getAttribute('class') || '').trim().toLowerCase();
const hasRightArrowIcon = !!el.querySelector('img[src*="icon-arrow-right"], img[src*="arrow-right"]');
return rel === 'next' || aria.includes('next') || title.includes('next') || cls.includes('next') || hasRightArrowIcon || ['next', '', '»', '>'].includes(text);
});
return cleanHref(explicitNext);
}
""",
expected_page_number,
)
if not href:
return None
return urljoin(page.url or self.settings.home_url, str(href))
except Exception:
return None
@staticmethod
def _build_listing_page_url(url: str, page_number: int) -> str:
parsed = urlparse(url)
query_items = [
(key, value)
for key, value in parse_qsl(parsed.query, keep_blank_values=True)
if key.lower() not in {"page", "pagenumber", "currentpage"}
]
query_items.append(("page", str(page_number)))
return urlunparse(parsed._replace(query=urlencode(query_items)))
def open_cars_listing(self, page: Page, *, url_override: str | None = None) -> None:
url = url_override or self.settings.listing.cars_url
@@ -202,6 +301,31 @@ class ListingCollector:
# Короткая пауза вместо длинного sleep.
time.sleep(1.0)
def _get_listing_html(self, page: Page, page_number: int) -> str:
try:
return page.locator("html").inner_html(timeout=5_000)
except Exception as exc:
logger.warning("listing html read failed on page %d: %s", page_number, exc)
return ""
def _has_next_page_from_html(self, html: str, current_page_number: int | None = None) -> bool:
if not html:
return False
normalized = html.replace("\\/", "/")
if re.search(
r"rel\s*=\s*['\"]next['\"]|aria-label\s*=\s*['\"][^'\"]*next|title\s*=\s*['\"][^'\"]*next|class\s*=\s*['\"][^'\"]*next|icon-arrow-right|arrow-right|>\s*next\s*<|>\s*[›»>]\s*<",
normalized,
re.IGNORECASE,
):
return True
if current_page_number is not None:
next_page = current_page_number + 1
if re.search(rf">\s*{next_page}\s*<", normalized, re.IGNORECASE):
return True
if re.search(rf"page={next_page}(?:\D|$)", normalized, re.IGNORECASE):
return True
return False
def apply_filters(
self,
page: Page,
@@ -279,74 +403,30 @@ class ListingCollector:
return False
def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult:
# Считываем ссылки одним проходом по DOM.
# Для IAAI HTML/hydration-извлечение стабильнее, чем прямой DOM eval.
self._accept_cookie_banner(page)
self._wait_for_listing_content(page)
try:
raw_items = page.eval_on_selector_all(
"a[href], [data-href], [href]",
"""
(nodes) => nodes.map((node) => ({
href:
node.getAttribute('href') ||
node.getAttribute('data-href') ||
node.getAttribute('data-url') ||
'',
title: node.getAttribute('title') || node.getAttribute('aria-label') || '',
text: (node.textContent || '').trim(),
}))
""",
)
except Exception as exc:
logger.warning("collect_current_page failed on page %d: %s", page_number, exc)
raw_items = []
total = min(len(raw_items), self.settings.listing.page_link_limit)
links: list[ListingVehicleLink] = []
seen: set[str] = set()
for idx in range(total):
item = raw_items[idx] if isinstance(raw_items[idx], dict) else {}
href = str(item.get("href") or "")
match = VEHICLE_HREF_RE.search(href)
if not match:
continue
lot_number = match.group(1)
absolute = urljoin(self.settings.home_url, match.group(0))
html = self._get_listing_html(page, page_number)
for absolute, lot_number in self._extract_vehicle_links_from_html(html):
if absolute in seen:
continue
seen.add(absolute)
title = first_non_empty([item.get("title"), item.get("text"), ""]) or ""
links.append(ListingVehicleLink(href=absolute, title=str(title).strip(), lot_number=lot_number))
links.append(ListingVehicleLink(href=absolute, title="", lot_number=lot_number))
if len(links) >= self.settings.listing.max_vehicles_per_run:
break
# Fallback: на IAAI ссылки иногда не рендерятся как <a>,
# но присутствуют в hydration/inline JSON внутри HTML (часто как \/VehicleDetail\/").
if not links:
try:
page.wait_for_timeout(1_500)
except Exception:
pass
try:
html = page.content()
except Exception as exc:
logger.debug("page.content() failed on page %d: %s", page_number, exc)
html = ""
if links:
logger.info(
"Page %d: recovered %d vehicle links from HTML fallback",
page_number,
len(links),
)
for absolute, lot_number in self._extract_vehicle_links_from_html(html):
if absolute in seen:
continue
seen.add(absolute)
links.append(ListingVehicleLink(href=absolute, title="", lot_number=lot_number))
if len(links) >= self.settings.listing.max_vehicles_per_run:
break
if links:
logger.info(
"Page %d: recovered %d vehicle links from HTML fallback",
page_number,
len(links),
)
next_page_detected = self._has_next_page(page)
next_page_detected = self._has_next_page_from_html(html, current_page_number=page_number)
if not next_page_detected and not html:
next_page_detected = self._has_next_page(page)
return ListingPageResult(source_url=page.url, page_number=page_number, vehicle_links=links, pagination_available=next_page_detected, next_page_detected=next_page_detected)
def _extract_vehicle_links_from_html(self, html: str) -> list[tuple[str, str]]:
@@ -373,12 +453,59 @@ class ListingCollector:
def go_to_next_page(self, page: Page, expected_page_number: int | None = None) -> bool:
# Запоминаем первую ссылку текущей страницы для определения смены контента.
old_first_href = ""
old_fingerprint: tuple[str, ...] = tuple()
try:
first_link = page.locator(VEHICLE_LINK_SELECTOR).first
if first_link.count() > 0:
old_first_href = first_link.get_attribute("href") or ""
except Exception:
pass
old_fingerprint = self._get_vehicle_link_fingerprint(page)
if expected_page_number is not None and expected_page_number > 1:
direct_page_url = self._build_listing_page_url(
page.url or self.settings.listing.cars_url,
expected_page_number,
)
try:
logger.debug("Navigating directly to listing page %d via URL: %s", expected_page_number, direct_page_url)
page.goto(direct_page_url, wait_until="domcontentloaded", timeout=15_000)
if self._wait_for_navigation_result(
page,
old_first_href,
expected_page_number,
old_fingerprint=old_fingerprint,
):
self.pacer.after_page_change()
return True
except Exception as exc:
logger.debug("Direct page-number navigation failed for page %d via %s: %s", expected_page_number, direct_page_url, exc)
next_href = self._extract_next_page_href(page, expected_page_number)
if next_href:
try:
logger.debug("Navigating directly to next listing page: %s", next_href)
page.goto(next_href, wait_until="domcontentloaded", timeout=15_000)
if self._wait_for_navigation_result(
page,
old_first_href,
expected_page_number,
old_fingerprint=old_fingerprint,
):
self.pacer.after_page_change()
return True
except Exception as exc:
logger.debug("Direct next-page navigation failed for %s: %s", next_href, exc)
if (
expected_page_number is not None
and expected_page_number >= self._DEEP_PAGINATION_DIRECT_ONLY_FROM_PAGE
):
logger.warning(
"Deep pagination direct navigation failed for page %d; skipping flaky UI pagination fallbacks",
expected_page_number,
)
return False
for selector in self._NEXT_PAGE_SELECTORS:
locator = page.locator(selector).first
@@ -398,7 +525,12 @@ class ListingCollector:
except Exception:
continue
if self._wait_for_navigation_result(page, old_first_href, expected_page_number):
if self._wait_for_navigation_result(
page,
old_first_href,
expected_page_number,
old_fingerprint=old_fingerprint,
):
self.pacer.after_page_change()
return True
@@ -458,7 +590,12 @@ class ListingCollector:
"""
))
if clicked:
if self._wait_for_navigation_result(page, old_first_href, expected_page_number):
if self._wait_for_navigation_result(
page,
old_first_href,
expected_page_number,
old_fingerprint=old_fingerprint,
):
self.pacer.after_page_change()
return True
except Exception as exc:
@@ -489,7 +626,12 @@ class ListingCollector:
"""
))
if clicked:
if self._wait_for_navigation_result(page, old_first_href, expected_page_number):
if self._wait_for_navigation_result(
page,
old_first_href,
expected_page_number,
old_fingerprint=old_fingerprint,
):
self.pacer.after_page_change()
return True
except Exception as exc:
@@ -574,10 +716,18 @@ class ListingCollector:
len(all_links) >= self.settings.listing.max_vehicles_per_run
or self.settings.listing.collect_current_page_only
or not self.settings.listing.include_pagination
or not page_result.next_page_detected
):
break
if not self.go_to_next_page(page):
blind_page_probe = not page_result.next_page_detected
if not self.go_to_next_page(page, expected_page_number=page_number + 1):
if blind_page_probe:
logger.info(
"Stopping pagination on page %d: direct page probe for %d failed and no next-page control was detected",
page_number,
page_number + 1,
)
break
break
return {
@@ -630,18 +780,45 @@ class ListingCollector:
try:
return bool(page.evaluate(
"""
() => Array.from(document.querySelectorAll('a,button,[role="button"]')).some((el) => {
const text = (el.textContent || '').trim().toLowerCase();
const aria = (el.getAttribute('aria-label') || '').trim().toLowerCase();
const title = (el.getAttribute('title') || '').trim().toLowerCase();
const rel = (el.getAttribute('rel') || '').trim().toLowerCase();
const cls = (el.getAttribute('class') || '').trim().toLowerCase();
const disabled = el.hasAttribute('disabled') || el.getAttribute('aria-disabled') === 'true' || cls.includes('disabled');
const visible = !!(el.offsetWidth || el.offsetHeight || el.getClientRects().length);
const hasRightArrowIcon = !!el.querySelector('img[src*="icon-arrow-right"], img[src*="arrow-right"]');
const looksNumericNext = /^\d+$/.test(text);
return !disabled && visible && (rel === 'next' || aria.includes('next') || title.includes('next') || cls.includes('next') || hasRightArrowIcon || looksNumericNext || ['next', '', '»', '>'].includes(text));
})
() => {
const visible = (el) => !!(el && (el.offsetWidth || el.offsetHeight || el.getClientRects().length));
const controls = Array.from(document.querySelectorAll('a,button,[role="button"]')).filter((el) => {
const cls = (el.getAttribute('class') || '').trim().toLowerCase();
const disabled = el.hasAttribute('disabled') || el.getAttribute('aria-disabled') === 'true' || cls.includes('disabled');
return !disabled && visible(el);
});
const hasExplicitNext = controls.some((el) => {
const text = (el.textContent || '').trim().toLowerCase();
const aria = (el.getAttribute('aria-label') || '').trim().toLowerCase();
const title = (el.getAttribute('title') || '').trim().toLowerCase();
const rel = (el.getAttribute('rel') || '').trim().toLowerCase();
const cls = (el.getAttribute('class') || '').trim().toLowerCase();
const hasRightArrowIcon = !!el.querySelector('img[src*="icon-arrow-right"], img[src*="arrow-right"]');
return rel === 'next' || aria.includes('next') || title.includes('next') || cls.includes('next') || hasRightArrowIcon || ['next', '', '»', '>'].includes(text);
});
if (hasExplicitNext) {
return true;
}
const current = controls.find((el) => {
const text = (el.textContent || '').trim();
const cls = (el.getAttribute('class') || '').toLowerCase();
const ariaCurrent = (el.getAttribute('aria-current') || '').toLowerCase();
return /^\d+$/.test(text) && (ariaCurrent === 'page' || cls.includes('active') || cls.includes('current') || cls.includes('selected'));
});
if (!current) {
return false;
}
const currentPage = parseInt((current.textContent || '').trim(), 10);
return controls.some((el) => {
const text = (el.textContent || '').trim();
return /^\d+$/.test(text) && parseInt(text, 10) === currentPage + 1;
});
}
"""
))
except Exception:

View File

@@ -38,6 +38,17 @@ def build_parser() -> argparse.ArgumentParser:
sync_listing_parser.add_argument("--only-new", choices=["true", "false"], default=None)
sync_listing_parser.add_argument("--output", default=str(default_output_dir / "iaai_sync_listing.json"))
# Новые команды для ingestion pipeline
subparsers.add_parser("discover-vehicles", help="Discover new vehicle URLs from sitemap")
fetch_parser = subparsers.add_parser("fetch-pending", help="Fetch raw data for pending candidates")
fetch_parser.add_argument("--limit", type=int, default=10, help="Max candidates to process")
enrich_parser = subparsers.add_parser("enrich-snapshots", help="Parse and enrich raw snapshots")
enrich_parser.add_argument("--limit", type=int, default=10, help="Max snapshots to process")
subparsers.add_parser("run-pipeline", help="Run full ingestion pipeline (discover -> fetch -> enrich)")
return parser
@@ -64,6 +75,39 @@ def main() -> None:
data = scraper.scrape_vehicle_detail(args.vehicle_url)
elif args.command == "sync-vehicle":
data = scraper.sync_vehicle(args.vehicle_url, lane=args.lane)
elif args.command == "sync-listing":
only_new = None if args.only_new is None else args.only_new == "true"
data = scraper.sync_listing(
make=args.make,
model=args.model,
lane=args.lane,
limit=args.limit,
only_new=only_new,
)
elif args.command == "discover-vehicles":
from .discovery_service import DiscoveryService
discovery = DiscoveryService()
count = discovery.discover_new_vehicles()
print(f"Discovered {count} new vehicle candidates")
return
elif args.command == "fetch-pending":
from .fetch_service import FetchService
fetch = FetchService()
count = fetch.process_pending_candidates(limit=args.limit)
print(f"Successfully fetched {count} candidates")
return
elif args.command == "enrich-snapshots":
from .enrichment_service import EnrichmentService
enrichment = EnrichmentService()
count = enrichment.process_unparsed_snapshots(limit=args.limit)
print(f"Successfully enriched {count} snapshots")
return
elif args.command == "run-pipeline":
from .scheduler_service import SchedulerService
scheduler = SchedulerService()
scheduler.run_full_pipeline()
print("Pipeline completed")
return
else:
only_new = None if args.only_new is None else args.only_new == "true"
data = scraper.sync_listing(

View File

@@ -106,6 +106,7 @@ class ListingConfig:
cars_url: str = _env_str("IAAI_CARS_LISTING_URL", "https://www.iaai.com/Vehiclelisting/Cars")
max_pages_per_run: int = _env_int("IAAI_MAX_PAGES_PER_RUN", 9999)
max_vehicles_per_run: int = _env_int("IAAI_MAX_VEHICLES_PER_RUN", 50000)
resume_max_nav_pages: int = _env_int("IAAI_LISTING_RESUME_MAX_NAV_PAGES", 90)
page_link_limit: int = _env_int("IAAI_PAGE_LINK_LIMIT", 500)
include_pagination: bool = _env_bool("IAAI_INCLUDE_PAGINATION", True)
collect_current_page_only: bool = _env_bool("IAAI_COLLECT_CURRENT_PAGE_ONLY", False)
@@ -118,6 +119,22 @@ class ListingConfig:
listing_segments_json: str = _env_str("IAAI_LISTING_SEGMENTS", "")
@dataclass(slots=True)
class DiscoveryConfig:
mode: str = _env_str("IAAI_DISCOVERY_MODE", "sitemap")
hourly_mode: str = _env_str("IAAI_HOURLY_MODE", "rolling_refresh")
hourly_refresh_batch_size: int = _env_int("IAAI_HOURLY_REFRESH_BATCH_SIZE", 3000)
sitemap_timeout_seconds: int = _env_int("IAAI_SITEMAP_TIMEOUT_SECONDS", 30)
sitemap_retry_attempts: int = _env_int("IAAI_SITEMAP_RETRY_ATTEMPTS", 4)
sitemap_retry_backoff_seconds: float = _env_float("IAAI_SITEMAP_RETRY_BACKOFF_SECONDS", 1.25)
sitemap_direct_probe_limit: int = _env_int("IAAI_SITEMAP_DIRECT_PROBE_LIMIT", 12)
sitemap_direct_probe_stop_after_misses: int = _env_int(
"IAAI_SITEMAP_DIRECT_PROBE_STOP_AFTER_MISSES",
3,
)
sitemap_use_curl_cffi: bool = _env_bool("IAAI_SITEMAP_USE_CURL_CFFI", True)
# Список брендов IAAI для автоматической сегментации.
# Покрывает >99% автомобилей на сайте. Порядок: от крупных к мелким.
IAAI_DEFAULT_MAKES: tuple[str, ...] = (
@@ -204,6 +221,9 @@ class DatabaseConfig:
@dataclass(slots=True)
class RedisConfig:
url: str = _env_str("IAAI_REDIS_URL", "redis://localhost:6379/0")
socket_timeout_seconds: float = _env_float("IAAI_REDIS_SOCKET_TIMEOUT_SECONDS", 10.0)
socket_connect_timeout_seconds: float = _env_float("IAAI_REDIS_SOCKET_CONNECT_TIMEOUT_SECONDS", 5.0)
health_check_interval_seconds: int = _env_int("IAAI_REDIS_HEALTH_CHECK_INTERVAL_SECONDS", 30)
# --- Конфиг Celery (лимиты задач, concurrency, beat-расписание) ---
@@ -214,6 +234,7 @@ class CeleryConfig:
result_backend: str = _env_str("CELERY_RESULT_BACKEND", "")
task_soft_time_limit: int = _env_int("CELERY_TASK_SOFT_TIME_LIMIT", 3300)
task_time_limit: int = _env_int("CELERY_TASK_TIME_LIMIT", 3600)
task_stall_timeout_seconds: int = _env_int("CELERY_TASK_STALL_TIMEOUT_SECONDS", 600)
worker_concurrency: int = _env_int("CELERY_WORKER_CONCURRENCY", 4)
worker_max_tasks_per_child: int = _env_int("CELERY_WORKER_MAX_TASKS_PER_CHILD", 5)
broker_visibility_timeout: int = _env_int("CELERY_BROKER_VISIBILITY_TIMEOUT", 7200)
@@ -221,6 +242,7 @@ class CeleryConfig:
beat_sync_limit: int | None = _env_int("CELERY_BEAT_SYNC_LIMIT", 0) or None
batch_size: int = _env_int("CELERY_BATCH_SIZE", 50)
parallel_tabs: int = _env_int("IAAI_PARALLEL_TABS", 8)
parallel_segments: bool = _env_bool("IAAI_PARALLEL_SEGMENTS", False)
block_resources: bool = _env_bool("IAAI_BLOCK_RESOURCES", True)
@@ -275,6 +297,7 @@ class Settings:
capture: CaptureConfig = field(default_factory=CaptureConfig)
pace: HumanPaceConfig = field(default_factory=HumanPaceConfig)
listing: ListingConfig = field(default_factory=ListingConfig)
discovery: DiscoveryConfig = field(default_factory=DiscoveryConfig)
database: DatabaseConfig = field(default_factory=DatabaseConfig)
redis: RedisConfig = field(default_factory=RedisConfig)
celery: CeleryConfig = field(default_factory=CeleryConfig)

View File

@@ -8,3 +8,7 @@ class AntiBotDetectedError(ScraperError):
class SiteStructureChangedError(ScraperError):
"""Вызывается, когда структура страницы изменилась и данных не хватает."""
class ListingResumeError(ScraperError):
"""Вызывается, когда resume по checkpoint больше недостижим."""

View File

@@ -0,0 +1,17 @@
from .sitemap import (
DEFAULT_SITEMAP_INDEX_URL,
SitemapDiscoveryError,
SitemapBlockedError,
SitemapDiscoveryResult,
discover_vehicle_urls_from_sitemap,
discover_vehicle_urls_from_sitemap_with_stats,
)
__all__ = [
"DEFAULT_SITEMAP_INDEX_URL",
"SitemapBlockedError",
"SitemapDiscoveryError",
"SitemapDiscoveryResult",
"discover_vehicle_urls_from_sitemap",
"discover_vehicle_urls_from_sitemap_with_stats",
]

View File

@@ -0,0 +1,441 @@
from __future__ import annotations
import gzip
import io
import logging
import random
import re
import time
import xml.etree.ElementTree as ET
from dataclasses import dataclass
from typing import Iterable
from urllib.error import HTTPError, URLError
from urllib.parse import urljoin, urlsplit, urlunsplit
from urllib.request import Request, urlopen
from ..core.config import Settings
try:
from curl_cffi import requests as curl_requests
except ImportError: # pragma: no cover - optional runtime dependency guard
curl_requests = None
logger = logging.getLogger("iaai_scraper.discovery.sitemap")
DEFAULT_SITEMAP_INDEX_URL = "https://www.iaai.com/Xj9rDOVMEi0hc38S/sitemap_index.xml"
_LOC_TAG_RE = re.compile(rb"<loc>\s*(.*?)\s*</loc>", re.IGNORECASE | re.DOTALL)
_XML_PREFIX_RE = re.compile(rb"^\s*(<\?xml\b.*?\?>)?\s*<", re.IGNORECASE | re.DOTALL)
_BLOCK_MARKERS = (
b"pardon our interruption",
b"incapsula",
b"please stand by",
b"_incapsula_resource",
b"as you were browsing something about your browser",
)
_HTML_MARKERS = (b"<html", b"<!doctype html", b"<script", b"document.getelementsbyclassname")
_CURL_IMPERSONATE_CHOICES = ("chrome136", "chrome133", "chrome131", "safari184")
class SitemapDiscoveryError(RuntimeError):
pass
class SitemapBlockedError(SitemapDiscoveryError):
pass
@dataclass(slots=True)
class SitemapFetchResult:
url: str
payload: bytes
source: str
blocked: bool = False
status_code: int | None = None
content_type: str | None = None
@dataclass(slots=True)
class SitemapDiscoveryStats:
transport: str
fetched_sitemaps: int = 0
blocked_sitemaps: int = 0
malformed_sitemaps: int = 0
direct_probe_hits: int = 0
direct_probe_misses: int = 0
@dataclass(slots=True)
class SitemapDiscoveryResult:
vehicle_urls: list[str]
stats: SitemapDiscoveryStats
_DEFAULT_HEADERS = {
"Accept": "application/xml,text/xml,application/xhtml+xml,text/html;q=0.9,*/*;q=0.8",
"Accept-Encoding": "gzip, deflate, br",
"Accept-Language": "en-US,en;q=0.9",
"Cache-Control": "no-cache",
"Pragma": "no-cache",
"Upgrade-Insecure-Requests": "1",
}
class _SitemapDownloader:
def __init__(self, settings: Settings) -> None:
self.settings = settings
self.discovery = settings.discovery
self.proxy_url = settings.proxy.server
self.proxy_auth = None
if settings.proxy.username:
password = settings.proxy.password or ""
self.proxy_auth = (settings.proxy.username, password)
self._transport_name = self._resolve_transport_name()
@property
def transport_name(self) -> str:
return self._transport_name
def _resolve_transport_name(self) -> str:
if self.discovery.sitemap_use_curl_cffi and curl_requests is not None:
return "curl_cffi"
return "urllib"
def fetch(self, url: str) -> SitemapFetchResult:
attempts = max(1, int(self.discovery.sitemap_retry_attempts))
last_exc: Exception | None = None
for attempt in range(1, attempts + 1):
try:
result = self._fetch_once(url)
if result.blocked:
raise SitemapBlockedError(f"Anti-bot page returned for {url}")
return result
except SitemapBlockedError as exc:
last_exc = exc
logger.warning(
"Sitemap fetch blocked (attempt %d/%d, transport=%s): %s",
attempt,
attempts,
self.transport_name,
url,
)
except Exception as exc:
last_exc = exc
logger.warning(
"Sitemap fetch failed (attempt %d/%d, transport=%s): %s -> %s",
attempt,
attempts,
self.transport_name,
url,
exc,
)
if attempt >= attempts:
break
time.sleep(self._backoff_delay(attempt))
if last_exc is None:
raise SitemapDiscoveryError(f"Failed to fetch sitemap: {url}")
if isinstance(last_exc, SitemapDiscoveryError):
raise last_exc
raise SitemapDiscoveryError(f"Failed to fetch sitemap {url}: {last_exc}") from last_exc
def _fetch_once(self, url: str) -> SitemapFetchResult:
if self.transport_name == "curl_cffi":
return self._fetch_with_curl_cffi(url)
return self._fetch_with_urllib(url)
def _build_headers(self) -> dict[str, str]:
headers = dict(_DEFAULT_HEADERS)
headers["User-Agent"] = self.settings.fingerprint.user_agent
return headers
def _fetch_with_curl_cffi(self, url: str) -> SitemapFetchResult:
assert curl_requests is not None
response = curl_requests.get(
url,
headers=self._build_headers(),
timeout=self.discovery.sitemap_timeout_seconds,
impersonate=random.choice(_CURL_IMPERSONATE_CHOICES),
proxies={"http": self.proxy_url, "https": self.proxy_url} if self.proxy_url else None,
proxy_auth=self.proxy_auth,
allow_redirects=True,
)
payload = self._decode_payload(
payload=response.content,
encoding=(response.headers.get("Content-Encoding") or ""),
url=url,
)
blocked = _looks_like_block_page(payload)
return SitemapFetchResult(
url=url,
payload=payload,
source="curl_cffi",
blocked=blocked,
status_code=int(response.status_code),
content_type=response.headers.get("Content-Type"),
)
def _fetch_with_urllib(self, url: str) -> SitemapFetchResult:
request = Request(url, headers=self._build_headers())
try:
with urlopen(request, timeout=self.discovery.sitemap_timeout_seconds) as response:
payload = response.read()
decoded = self._decode_payload(
payload=payload,
encoding=(response.headers.get("Content-Encoding") or ""),
url=url,
)
return SitemapFetchResult(
url=url,
payload=decoded,
source="urllib",
blocked=_looks_like_block_page(decoded),
status_code=getattr(response, "status", None),
content_type=response.headers.get("Content-Type"),
)
except HTTPError as exc:
payload = exc.read() if hasattr(exc, "read") else b""
decoded = self._decode_payload(payload=payload, encoding=exc.headers.get("Content-Encoding", ""), url=url)
blocked = exc.code in {403, 429} or _looks_like_block_page(decoded)
if blocked:
raise SitemapBlockedError(f"HTTP {exc.code} for {url}") from exc
raise SitemapDiscoveryError(f"HTTP {exc.code} for {url}") from exc
except URLError as exc:
raise SitemapDiscoveryError(f"Network error for {url}: {exc}") from exc
@staticmethod
def _decode_payload(*, payload: bytes, encoding: str, url: str) -> bytes:
normalized = (encoding or "").lower().strip()
if normalized == "gzip" or url.lower().endswith(".gz"):
try:
return gzip.GzipFile(fileobj=io.BytesIO(payload)).read()
except OSError:
return payload
return payload
def _backoff_delay(self, attempt: int) -> float:
base = max(0.1, float(self.discovery.sitemap_retry_backoff_seconds))
jitter = random.uniform(0.05, 0.35)
return base * (2 ** (attempt - 1)) + jitter
def _is_vehicle_sitemap_url(url: str) -> bool:
lowered = url.strip().lower()
if "sitemapbranches" in lowered or "sitemapauctions" in lowered:
return False
return lowered.endswith(".xml") or lowered.endswith(".xml.gz")
def _normalize_vehicle_url(url: str) -> str:
parts = urlsplit(url.strip())
return urlunsplit((parts.scheme, parts.netloc, parts.path, "", ""))
def _local_name(tag: str) -> str:
if "}" in tag:
return tag.rsplit("}", 1)[1]
return tag
def _looks_like_block_page(payload: bytes) -> bool:
sample = payload[:8192].lower()
if any(marker in sample for marker in _BLOCK_MARKERS):
return True
if any(marker in sample for marker in _HTML_MARKERS) and b"<loc>" not in sample:
return True
return False
def _iter_loc_values_fallback(xml_bytes: bytes) -> Iterable[str]:
for match in _LOC_TAG_RE.finditer(xml_bytes):
try:
value = match.group(1).decode("utf-8", errors="ignore").strip()
except Exception:
continue
if value:
yield value
def _iter_loc_values(xml_bytes: bytes) -> Iterable[str]:
if _looks_like_block_page(xml_bytes):
raise SitemapBlockedError("Anti-bot content returned instead of sitemap XML")
if not _XML_PREFIX_RE.search(xml_bytes[:256]):
fallback_values = list(_iter_loc_values_fallback(xml_bytes))
if fallback_values:
yield from fallback_values
return
try:
root = ET.fromstring(xml_bytes)
except ET.ParseError as exc:
fallback_values = list(_iter_loc_values_fallback(xml_bytes))
if fallback_values:
logger.warning(
"Falling back to regex sitemap loc extraction after XML parse error: %s",
exc,
)
yield from fallback_values
return
raise SitemapDiscoveryError(f"Invalid sitemap XML: {exc}") from exc
for element in root.iter():
if _local_name(element.tag) != "loc":
continue
if not element.text:
continue
value = element.text.strip()
if value:
yield value
def _filter_vehicle_urls(urls: Iterable[str]) -> list[str]:
result: list[str] = []
seen: set[str] = set()
for url in urls:
normalized = _normalize_vehicle_url(url)
if "/VehicleDetail/" not in normalized and "/vehicledetail/" not in normalized:
continue
if normalized in seen:
continue
seen.add(normalized)
result.append(normalized)
return result
def _looks_like_vehicle_detail_sitemap(urls: list[str]) -> bool:
return any("/vehicledetail/" in url.lower() for url in urls)
def _derive_direct_probe_urls(index_url: str, limit: int) -> list[str]:
base_dir = index_url.rsplit("/", 1)[0] + "/"
return [urljoin(base_dir, f"sitemap{idx}.xml") for idx in range(1, max(1, limit) + 1)]
def _discover_sitemap_urls(index_url: str, downloader: _SitemapDownloader, stats: SitemapDiscoveryStats) -> list[str]:
logger.info("Downloading sitemap index: %s", index_url)
try:
index_result = downloader.fetch(index_url)
sitemap_urls = [url for url in _iter_loc_values(index_result.payload) if _is_vehicle_sitemap_url(url)]
if sitemap_urls:
return sitemap_urls
except SitemapBlockedError as exc:
stats.blocked_sitemaps += 1
logger.warning("Sitemap index blocked, switching to direct probe: %s", exc)
except SitemapDiscoveryError as exc:
stats.malformed_sitemaps += 1
logger.warning("Sitemap index unusable, switching to direct probe: %s", exc)
logger.warning("Sitemap index returned no usable sitemap URLs; switching to direct probe")
return _probe_direct_sitemap_urls(index_url, downloader, stats)
def _probe_direct_sitemap_urls(
index_url: str,
downloader: _SitemapDownloader,
stats: SitemapDiscoveryStats,
) -> list[str]:
discovered: list[str] = []
consecutive_misses = 0
probe_urls = _derive_direct_probe_urls(index_url, downloader.discovery.sitemap_direct_probe_limit)
for sitemap_url in probe_urls:
try:
result = downloader.fetch(sitemap_url)
raw_urls = list(_iter_loc_values(result.payload))
except SitemapBlockedError:
stats.blocked_sitemaps += 1
consecutive_misses += 1
stats.direct_probe_misses += 1
if consecutive_misses >= downloader.discovery.sitemap_direct_probe_stop_after_misses:
break
continue
except SitemapDiscoveryError:
consecutive_misses += 1
stats.direct_probe_misses += 1
if consecutive_misses >= downloader.discovery.sitemap_direct_probe_stop_after_misses:
break
continue
if not raw_urls:
consecutive_misses += 1
stats.direct_probe_misses += 1
if consecutive_misses >= downloader.discovery.sitemap_direct_probe_stop_after_misses:
break
continue
consecutive_misses = 0
stats.direct_probe_hits += 1
discovered.append(sitemap_url)
return discovered
def _collect_vehicle_urls_from_sitemap(
sitemap_url: str,
downloader: _SitemapDownloader,
stats: SitemapDiscoveryStats,
) -> list[str]:
logger.info("Downloading sitemap: %s", sitemap_url)
try:
result = downloader.fetch(sitemap_url)
raw_urls = list(_iter_loc_values(result.payload))
except SitemapBlockedError as exc:
stats.blocked_sitemaps += 1
logger.warning("Skipping blocked sitemap %s: %s", sitemap_url, exc)
return []
except SitemapDiscoveryError as exc:
stats.malformed_sitemaps += 1
logger.warning("Skipping malformed sitemap %s: %s", sitemap_url, exc)
return []
stats.fetched_sitemaps += 1
vehicle_urls = _filter_vehicle_urls(raw_urls)
if raw_urls and not vehicle_urls and not _looks_like_vehicle_detail_sitemap(raw_urls):
logger.info("Skipping non-vehicle sitemap %s", sitemap_url)
return []
logger.info("Sitemap %s yielded %d vehicle URLs", sitemap_url, len(vehicle_urls))
return vehicle_urls
def discover_vehicle_urls_from_sitemap(
index_url: str = DEFAULT_SITEMAP_INDEX_URL,
*,
settings: Settings | None = None,
) -> list[str]:
result = discover_vehicle_urls_from_sitemap_with_stats(index_url=index_url, settings=settings)
return result.vehicle_urls
def discover_vehicle_urls_from_sitemap_with_stats(
index_url: str = DEFAULT_SITEMAP_INDEX_URL,
*,
settings: Settings | None = None,
) -> SitemapDiscoveryResult:
runtime_settings = settings or Settings()
downloader = _SitemapDownloader(runtime_settings)
stats = SitemapDiscoveryStats(transport=downloader.transport_name)
sitemap_urls = _discover_sitemap_urls(index_url, downloader, stats)
if not sitemap_urls:
raise SitemapDiscoveryError("Sitemap discovery found no sitemap URLs")
all_vehicle_urls: list[str] = []
for sitemap_url in sitemap_urls:
all_vehicle_urls.extend(_collect_vehicle_urls_from_sitemap(sitemap_url, downloader, stats))
deduped = _filter_vehicle_urls(all_vehicle_urls)
if not deduped:
raise SitemapDiscoveryError(
"No vehicle detail URLs discovered from vehicle sitemaps; likely anti-bot or upstream sitemap issue"
)
logger.info(
"Sitemap discovery done: %d vehicle URLs (transport=%s fetched=%d blocked=%d malformed=%d direct_hits=%d direct_misses=%d)",
len(deduped),
stats.transport,
stats.fetched_sitemaps,
stats.blocked_sitemaps,
stats.malformed_sitemaps,
stats.direct_probe_hits,
stats.direct_probe_misses,
)
return SitemapDiscoveryResult(vehicle_urls=deduped, stats=stats)

View File

@@ -0,0 +1,129 @@
import logging
from datetime import datetime, timezone
from typing import List
from sqlalchemy.orm import Session
from .core.config import settings
from .discovery import discover_vehicle_urls_from_sitemap_with_stats, SitemapDiscoveryError
from .storage.db import get_db_session
from .storage.models import VehicleCandidate
logger = logging.getLogger("iaai_scraper.discovery_service")
class DiscoveryService:
"""Сервис для обнаружения новых автомобилей через sitemap IAAI."""
def __init__(self):
self.db: Session = get_db_session()
def discover_new_vehicles(self, max_urls: int = None) -> int:
"""
Обнаруживает новые URL автомобилей и добавляет их в vehicle_candidates.
Args:
max_urls: Максимальное количество URL для обработки (None = все)
Returns:
Количество добавленных кандидатов
"""
try:
logger.info("Starting vehicle discovery from sitemap")
result = discover_vehicle_urls_from_sitemap_with_stats()
if not result.urls:
logger.info("No new URLs discovered")
return 0
# Ограничиваем количество URL
urls_to_process = result.urls[:max_urls] if max_urls else result.urls
logger.info(f"Discovered {len(urls_to_process)} potential vehicle URLs (from {len(result.urls)} total)")
# Фильтруем уже существующие кандидаты
existing_urls = self._get_existing_candidate_urls(urls_to_process)
new_urls = [url for url in urls_to_process if url not in existing_urls]
if not new_urls:
logger.info("All discovered URLs already exist as candidates")
return 0
# Добавляем новых кандидатов
added_count = self._add_candidates(new_urls)
logger.info(f"Added {added_count} new vehicle candidates")
return added_count
except SitemapDiscoveryError as e:
logger.error(f"Sitemap discovery failed: {e}")
raise
except Exception as e:
logger.error(f"Unexpected error during discovery: {e}")
raise
def _get_existing_candidate_urls(self, urls: List[str]) -> set:
"""Получает множество уже существующих URL кандидатов."""
if not urls:
return set()
# Разбиваем на батчи для эффективности
batch_size = 1000
existing = set()
for i in range(0, len(urls), batch_size):
batch = urls[i:i + batch_size]
result = self.db.query(VehicleCandidate.url).filter(
VehicleCandidate.url.in_(batch)
).all()
existing.update(row[0] for row in result)
return existing
def _add_candidates(self, urls: List[str]) -> int:
"""Добавляет новые кандидаты в базу данных."""
now = datetime.now(timezone.utc)
candidates = [
VehicleCandidate(
url=url,
discovered_at=now,
status="pending",
priority=0
)
for url in urls
]
self.db.add_all(candidates)
self.db.commit()
return len(candidates)
def get_pending_candidates(self, limit: int = 100) -> List[VehicleCandidate]:
"""Получает кандидатов в статусе 'pending' для обработки."""
return self.db.query(VehicleCandidate).filter(
VehicleCandidate.status == "pending"
).order_by(VehicleCandidate.priority.desc(), VehicleCandidate.discovered_at).limit(limit).all()
def mark_candidate_processing(self, candidate_id: int):
"""Помечает кандидата как обрабатываемого."""
self.db.query(VehicleCandidate).filter(
VehicleCandidate.id == candidate_id
).update({
"status": "processing",
"last_attempt_at": datetime.now(timezone.utc),
"attempts": VehicleCandidate.attempts + 1
})
self.db.commit()
def mark_candidate_processed(self, candidate_id: int):
"""Помечает кандидата как обработанного."""
self.db.query(VehicleCandidate).filter(
VehicleCandidate.id == candidate_id
).update({"status": "processed"})
self.db.commit()
def mark_candidate_failed(self, candidate_id: int):
"""Помечает кандидата как неудачного."""
self.db.query(VehicleCandidate).filter(
VehicleCandidate.id == candidate_id
).update({"status": "failed"})
self.db.commit()

View File

@@ -0,0 +1,140 @@
import json
import logging
from datetime import datetime, timezone
from typing import Optional
from sqlalchemy.orm import Session
from .core.config import settings
from .parsing.mapper import CarMapper
from .parsing.parser import VehicleParser
from .storage.db import get_db_session
from .storage.models import VehicleRawSnapshot, VehicleParseResult, Car
from .storage.schemas import CarRecord
logger = logging.getLogger("iaai_scraper.enrichment_service")
class EnrichmentService:
"""Сервис для парсинга сырых данных и обогащения автомобилей."""
def __init__(self):
self.db: Session = get_db_session()
self.parser = VehicleParser()
self.mapper = CarMapper()
def enrich_vehicle(self, snapshot: VehicleRawSnapshot) -> Optional[VehicleParseResult]:
"""
Парсит snapshot и сохраняет результат.
Args:
snapshot: Raw snapshot для парсинга
Returns:
VehicleParseResult если парсинг успешен
"""
logger.info(f"Enriching snapshot {snapshot.id} for candidate {snapshot.candidate_id}")
if not snapshot.success or not snapshot.raw_data:
logger.warning(f"Snapshot {snapshot.id} has no data to parse")
return self._save_parse_result(snapshot.id, False, None, "No raw data available")
try:
# Парсим данные
parsed_data = self._parse_raw_data(snapshot.raw_data)
if not parsed_data:
return self._save_parse_result(snapshot.id, False, None, "Parsing failed")
# Маппим в CarRecord
car_record = self._map_to_car_record(parsed_data)
if not car_record:
return self._save_parse_result(snapshot.id, False, None, "Mapping failed")
# Сохраняем в cars таблицу
self._save_car(car_record)
# Сохраняем успешный результат парсинга
return self._save_parse_result(snapshot.id, True, json.dumps(parsed_data))
except Exception as e:
error_msg = f"Unexpected error during enrichment: {str(e)}"
logger.error(f"Error enriching snapshot {snapshot.id}: {error_msg}")
return self._save_parse_result(snapshot.id, False, None, error_msg)
def _parse_raw_data(self, raw_data: str) -> Optional[dict]:
"""Парсит сырые данные в словарь."""
try:
# Если это JSON от browser capture
if raw_data.strip().startswith('{'):
data = json.loads(raw_data)
# Извлекаем vehicle_summary из scrape result
return data.get("vehicle_summary", {})
# Если HTML, используем VehicleParser
parsed = self.parser.parse_vehicle_page(raw_data, "dummy_url")
return parsed.get("vehicle_summary", {})
except json.JSONDecodeError:
# Если не JSON, пробуем как HTML
try:
parsed = self.parser.parse_vehicle_page(raw_data, "dummy_url")
return parsed.get("vehicle_summary", {})
except Exception:
return None
def _map_to_car_record(self, parsed_data: dict) -> Optional[CarRecord]:
"""Маппит parsed data в CarRecord."""
try:
# Используем CarMapper
payload_insights = {} # TODO: extract from raw data if needed
return self.mapper.map_to_car_record("dummy_url", parsed_data, payload_insights)
except Exception as e:
logger.error(f"Mapping failed: {e}")
return None
def _save_car(self, car_record: CarRecord):
"""Сохраняет CarRecord в базу данных."""
# Используем PersistenceService
from .storage.db import PersistenceService
from .core.config import settings
persistence = PersistenceService(settings)
persistence.create_tables()
upsert_result = persistence.upsert_car(car_record)
logger.info(f"Car upserted: {upsert_result}")
def _save_parse_result(self, snapshot_id: int, success: bool,
parsed_data: Optional[str], error_message: Optional[str] = None) -> VehicleParseResult:
"""Сохраняет результат парсинга."""
result = VehicleParseResult(
snapshot_id=snapshot_id,
success=success,
parsed_data=parsed_data,
error_message=error_message
)
self.db.add(result)
self.db.commit()
self.db.refresh(result)
return result
def process_unparsed_snapshots(self, limit: int = 10) -> int:
"""
Обрабатывает snapshots без результатов парсинга.
Returns:
Количество успешно обогащенных snapshots
"""
# Находим snapshots без parse results
snapshots = self.db.query(VehicleRawSnapshot).outerjoin(
VehicleParseResult, VehicleRawSnapshot.id == VehicleParseResult.snapshot_id
).filter(
VehicleRawSnapshot.success == True,
VehicleParseResult.id.is_(None)
).limit(limit).all()
enriched_count = 0
for snapshot in snapshots:
result = self.enrich_vehicle(snapshot)
if result and result.success:
enriched_count += 1
return enriched_count

View File

@@ -0,0 +1,119 @@
import logging
from datetime import datetime, timezone
from typing import Optional
from sqlalchemy.orm import Session
from .core.config import settings
from .scraper import IAAIScraper
from .storage.db import get_db_session
from .storage.models import VehicleCandidate, VehicleRawSnapshot
logger = logging.getLogger("iaai_scraper.fetch_service")
class FetchService:
"""Сервис для захвата сырых данных автомобилей (HTTP/browser fallback)."""
def __init__(self):
self.db: Session = get_db_session()
self.scraper = IAAIScraper()
def fetch_vehicle_data(self, candidate: VehicleCandidate) -> Optional[VehicleRawSnapshot]:
"""
Захватывает данные для кандидата автомобиля.
Args:
candidate: Кандидат для обработки
Returns:
VehicleRawSnapshot если захват успешен, None если неудача
"""
logger.info(f"Fetching data for candidate {candidate.id}: {candidate.url}")
try:
# Сначала пытаемся HTTP fast-path
raw_data = self._try_http_capture(candidate.url)
if raw_data:
return self._save_snapshot(candidate.id, "http", True, raw_data)
# Если HTTP не сработал, используем browser fallback
logger.info(f"HTTP failed for {candidate.url}, trying browser fallback")
raw_data = self._try_browser_capture(candidate.url)
if raw_data:
return self._save_snapshot(candidate.id, "browser", True, raw_data)
# Оба метода failed
logger.warning(f"Both HTTP and browser capture failed for {candidate.url}")
self._save_snapshot(candidate.id, "browser", False, None, "Both capture methods failed")
return None
except Exception as e:
error_msg = f"Unexpected error during capture: {str(e)}"
logger.error(f"Error fetching {candidate.url}: {error_msg}")
self._save_snapshot(candidate.id, "unknown", False, None, error_msg)
return None
def _try_http_capture(self, url: str) -> Optional[str]:
"""Пытается захватить данные через HTTP."""
try:
# Используем существующий метод из scraper
# Но нам нужно инициализировать scraper с сессией
# Пока заглушка - интегрируем позже
# Для теста вернем None, чтобы использовать browser
return None
except Exception as e:
logger.debug(f"HTTP capture failed for {url}: {e}")
return None
def _try_browser_capture(self, url: str) -> Optional[str]:
"""Пытается захватить данные через browser."""
try:
# Используем scrape_vehicle_detail из scraper
with IAAIScraper() as scraper:
result = scraper.scrape_vehicle_detail(url)
# Возвращаем HTML или JSON данные
return result.get("raw_html") or json.dumps(result)
except Exception as e:
logger.debug(f"Browser capture failed for {url}: {e}")
return None
def _save_snapshot(self, candidate_id: int, method: str, success: bool,
raw_data: Optional[str], error_message: Optional[str] = None) -> VehicleRawSnapshot:
"""Сохраняет snapshot в базу данных."""
snapshot = VehicleRawSnapshot(
candidate_id=candidate_id,
method=method,
success=success,
raw_data=raw_data,
error_message=error_message
)
self.db.add(snapshot)
self.db.commit()
self.db.refresh(snapshot)
return snapshot
def process_pending_candidates(self, limit: int = 10) -> int:
"""
Обрабатывает ожидающих кандидатов.
Returns:
Количество успешно обработанных кандидатов
"""
from .discovery_service import DiscoveryService
discovery = DiscoveryService()
candidates = discovery.get_pending_candidates(limit)
processed_count = 0
for candidate in candidates:
discovery.mark_candidate_processing(candidate.id)
snapshot = self.fetch_vehicle_data(candidate)
if snapshot and snapshot.success:
discovery.mark_candidate_processed(candidate.id)
processed_count += 1
else:
discovery.mark_candidate_failed(candidate.id)
return processed_count

View File

@@ -0,0 +1,69 @@
import logging
import time
from datetime import datetime, timezone, timedelta
from .core.config import settings
from .discovery_service import DiscoveryService
from .fetch_service import FetchService
from .enrichment_service import EnrichmentService
logger = logging.getLogger("iaai_scraper.scheduler_service")
class SchedulerService:
"""Сервис для планирования и координации ingestion pipeline."""
def __init__(self):
self.discovery = DiscoveryService()
self.fetch = FetchService()
self.enrichment = EnrichmentService()
def run_full_pipeline(self):
"""Запускает полный цикл ingestion: discovery -> fetch -> enrichment."""
logger.info("Starting full ingestion pipeline")
try:
# 1. Discovery phase
logger.info("Phase 1: Discovery")
discovered_count = self.discovery.discover_new_vehicles()
logger.info(f"Discovered {discovered_count} new candidates")
# 2. Fetch phase
logger.info("Phase 2: Fetch")
fetched_count = self.fetch.process_pending_candidates(limit=50)
logger.info(f"Successfully fetched {fetched_count} candidates")
# 3. Enrichment phase
logger.info("Phase 3: Enrichment")
enriched_count = self.enrichment.process_unparsed_snapshots(limit=50)
logger.info(f"Successfully enriched {enriched_count} snapshots")
logger.info("Ingestion pipeline completed")
except Exception as e:
logger.error(f"Pipeline failed: {e}")
raise
def run_continuous_pipeline(self, interval_minutes: int = 30):
"""Запускает непрерывный цикл ingestion с интервалом."""
logger.info(f"Starting continuous ingestion pipeline with {interval_minutes}min intervals")
while True:
try:
self.run_full_pipeline()
except Exception as e:
logger.error(f"Pipeline iteration failed: {e}")
logger.info(f"Sleeping for {interval_minutes} minutes")
time.sleep(interval_minutes * 60)
def run_targeted_enrichment(self):
"""Запускает только enrichment для существующих snapshots."""
logger.info("Running targeted enrichment")
enriched_count = self.enrichment.process_unparsed_snapshots(limit=100)
logger.info(f"Enriched {enriched_count} snapshots")
def cleanup_old_data(self, days_to_keep: int = 30):
"""Очищает старые данные (опционально)."""
# TODO: implement if needed
pass

View File

@@ -20,8 +20,9 @@ from playwright.sync_api import BrowserContext, Page, sync_playwright
from playwright.sync_api import TimeoutError as PlaywrightTimeoutError
from .browser import BrowserFactory, HumanPacer, NetworkCapture
from .discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitemap_with_stats
from .core.config import Settings, settings, parse_listing_segments
from .core.exceptions import AntiBotDetectedError, SiteStructureChangedError
from .core.exceptions import AntiBotDetectedError, ListingResumeError, SiteStructureChangedError
from .core.logs import set_trace_id, setup_logging
from .core.retry import retryable
from .core.runtime_config import RuntimeConfig
@@ -120,6 +121,7 @@ class IAAIScraper:
self.car_mapper = CarMapper()
self.persistence = PersistenceService(self.settings)
self._shutdown_requested = False
self._progress_callback: Callable[[str, dict[str, Any]], None] | None = None
proxy_url = self.settings.proxy.server
if proxy_url:
_proxy_kwargs: dict = {
@@ -148,6 +150,18 @@ class IAAIScraper:
set_trace_id(trace_id)
return trace_id
def set_progress_callback(self, callback: Callable[[str, dict[str, Any]], None] | None) -> None:
self._progress_callback = callback
def _emit_progress(self, stage: str, **payload: Any) -> None:
callback = self._progress_callback
if callback is None:
return
try:
callback(stage, payload)
except Exception:
logger.debug("progress callback failed for stage=%s", stage, exc_info=True)
def __enter__(self) -> "IAAIScraper":
if self.playwright is None:
self.playwright = sync_playwright().start()
@@ -237,6 +251,133 @@ class IAAIScraper:
vehicle_urls.append(normalized)
return vehicle_urls
def _should_use_sitemap_discovery(
self,
*,
make: str | None,
model: str | None,
listing_url: str | None,
year_min: int | None,
year_max: int | None,
limit: int | None,
effective_only_new: bool,
) -> bool:
# Полный unfiltered scan всегда должен идти через sitemap.
# Listing-пагинация на IAAI медленная, нестабильная и может зависать
# на resume/recovery даже когда raw HTTP по карточкам работает быстро.
return (
make is None
and model is None
and listing_url is None
and year_min is None
and year_max is None
and (limit is None or limit <= 0)
and not effective_only_new
)
def _sync_listing_via_sitemap(
self,
*,
lane: str,
limit: int | None,
effective_only_new: bool,
) -> dict[str, Any]:
del effective_only_new # sitemap full-scan path is used only for full discovery.
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=self.settings)
vehicle_urls = self._dedupe_urls(discovery_result.vehicle_urls)
if limit is not None and limit > 0:
vehicle_urls = vehicle_urls[:limit]
batch_size = self.settings.celery.batch_size
total = len(vehicle_urls)
cars_upserted = 0
cars_failed = 0
images_upserted = 0
failures: list[dict[str, str]] = []
logger.info("Sitemap-based full discovery found %d vehicle URLs", total)
self._emit_progress(
"sitemap_discovery_done",
urls_found=total,
batch_size=batch_size,
transport=discovery_result.stats.transport,
fetched_sitemaps=discovery_result.stats.fetched_sitemaps,
blocked_sitemaps=discovery_result.stats.blocked_sitemaps,
malformed_sitemaps=discovery_result.stats.malformed_sitemaps,
direct_probe_hits=discovery_result.stats.direct_probe_hits,
direct_probe_misses=discovery_result.stats.direct_probe_misses,
)
for batch_start in range(0, total, batch_size):
batch_urls = vehicle_urls[batch_start:batch_start + batch_size]
self._emit_progress(
"listing_batch_start",
source="sitemap",
batch_offset=batch_start,
batch_size=len(batch_urls),
pending_urls=max(0, total - batch_start),
batch_first_url=batch_urls[0] if batch_urls else None,
)
try:
batch_result = self.sync_batch(batch_urls, lane=lane)
cars_upserted += batch_result.get("cars_upserted", 0)
cars_failed += batch_result.get("cars_failed", 0)
images_upserted += batch_result.get("images_upserted", 0)
failures.extend(batch_result.get("failures", []))
except Exception as batch_exc:
logger.error(
"Sitemap batch %d-%d failed: %s",
batch_start + 1,
batch_start + len(batch_urls),
batch_exc,
)
cars_failed += len(batch_urls)
failures.append({
"vehicle_url": f"sitemap_batch_{batch_start}",
"error": str(batch_exc),
})
self._emit_progress(
"listing_batch_done",
source="sitemap",
batch_offset=batch_start,
cars_upserted=cars_upserted,
cars_failed=cars_failed,
pending_urls=max(0, total - (batch_start + len(batch_urls))),
)
return {
"listing": {
"status": "ok",
"listing_url": self.settings.listing.cars_url,
"applied_filters": {
"make": None,
"model": None,
"year_min": None,
"year_max": None,
},
"pages_collected": 0,
"vehicles_collected": total,
"vehicle_urls": vehicle_urls,
"early_stopped": False,
"truncated_by_time_budget": False,
"pagination_interrupted": False,
"pages": [],
"source": "sitemap",
"transport": discovery_result.stats.transport,
"blocked_sitemaps": discovery_result.stats.blocked_sitemaps,
"malformed_sitemaps": discovery_result.stats.malformed_sitemaps,
"direct_probe_hits": discovery_result.stats.direct_probe_hits,
"direct_probe_misses": discovery_result.stats.direct_probe_misses,
},
"total": total,
"skipped_existing": 0,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"failures": failures,
"all_listing_origin_urls": set(vehicle_urls),
}
def _filter_known_urls(self, vehicle_urls: list[str]) -> tuple[list[str], int]:
url_to_origin_id = {
url: self._extract_db_origin_id_from_url(url)
@@ -310,6 +451,30 @@ class IAAIScraper:
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)
def _open_listing_page(
self,
*,
make: str | None,
model: str | None,
listing_url: str | None = None,
year_min: int | None = None,
year_max: int | None = None,
) -> tuple[Page, dict[str, str | int | None]]:
page = self._get_page_with_warmup()
try:
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,
)
except Exception:
page.close()
raise
return page, applied_filters
def _reopen_listing_and_resume(
self,
*,
@@ -319,29 +484,51 @@ class IAAIScraper:
listing_url: str | None = None,
year_min: int | None = None,
year_max: int | None = None,
max_nav_pages: int = 10,
) -> Page:
max_nav_pages: int | None = None,
) -> tuple[Page, dict[str, str | int | None]]:
if max_nav_pages is None:
max_nav_pages = max(1, self.settings.listing.resume_max_nav_pages)
# Ограничиваем глубину навигации: если до цели > max_nav_pages кликов — не пытаемся.
if target_page_number > max_nav_pages + 1:
raise RuntimeError(
raise ListingResumeError(
f"Cannot resume at page {target_page_number}: "
f"exceeds max navigation depth ({max_nav_pages} pages)"
)
page = self._get_page_with_warmup()
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
page, applied_filters = self._open_listing_page(
make=make,
model=model,
listing_url=listing_url,
year_min=year_min,
year_max=year_max,
)
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):
page.close()
raise RuntimeError(f"Failed to resume listing at page {target_page_number}")
raise ListingResumeError(f"Failed to resume listing at page {target_page_number}")
logger.warning("Listing resumed at page %d after recovery", target_page_number)
return page
return page, applied_filters
def _open_listing_for_stream(
self,
*,
make: str | None,
model: str | None,
listing_url: str | None = None,
year_min: int | None = None,
year_max: int | None = None,
) -> tuple[Page, dict[str, str | int | None]]:
# Всегда открываем с page 1. Page-level resume убран: эфемерные URL и короткие
# сегменты делали pagination-resume хрупким. Bootstrap прогресс сохраняется
# на уровне сегментов в worker/tasks.py (_save_last_completed_segment).
return self._open_listing_page(
make=make,
model=model,
listing_url=listing_url,
year_min=year_min,
year_max=year_max,
)
def collect_listing(
self,
@@ -377,8 +564,6 @@ class IAAIScraper:
limit: int | None,
effective_only_new: bool,
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,
@@ -397,6 +582,7 @@ class IAAIScraper:
failures: list[dict[str, str]] = []
early_stopped = False
truncated_by_time_budget = False
pagination_interrupted = False
known_origin_ids: set[str] | None = None
threshold = self.settings.listing.early_stop_threshold
@@ -518,27 +704,66 @@ class IAAIScraper:
failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)})
return True
page = self._get_page_with_warmup() if start_page <= 1 else None
page = None
try:
if start_page <= 1:
assert page is not None
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, "year_min": year_min, "year_max": year_max}
page, applied_filters = self._open_listing_for_stream(
make=make,
model=model,
listing_url=listing_url,
year_min=year_min,
year_max=year_max,
)
self._emit_progress(
"listing_opened",
make=make,
model=model,
listing_url=listing_url,
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)
for page_number in range(1, self.settings.listing.max_pages_per_run + 1):
self._emit_progress(
"listing_collect_start",
page_number=page_number,
pending_urls=len(pending_urls),
total_new=total,
)
# Защита от crashes при сборе данных страницы
page_result = None
collect_attempts = 0
max_collect_attempts = 3
while collect_attempts < max_collect_attempts and page_result is None:
collect_attempts += 1
try:
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
except (PlaywrightError, PlaywrightTimeoutError, Exception) as collect_exc:
logger.warning("Exception during page collection (page %d, attempt %d/%d): %s",
page_number, collect_attempts, max_collect_attempts, collect_exc)
if collect_attempts < max_collect_attempts:
logger.info("Retrying page collection after 5 seconds...")
time.sleep(5)
try:
page.reload(wait_until="domcontentloaded", timeout=30_000)
except Exception as reload_exc:
logger.debug("Page reload failed during collection retry: %s", reload_exc)
else:
logger.error("All collection attempts failed for page %d — stopping pagination", page_number)
pagination_interrupted = True
page = None
break
if page_result is None:
break
self._emit_progress(
"listing_collect_done",
page_number=page_number,
links_found=len(page_result.vehicle_links),
next_page_detected=page_result.next_page_detected,
)
pages_info.append({
"page_number": page_result.page_number,
"links_found": len(page_result.vehicle_links),
@@ -552,6 +777,7 @@ class IAAIScraper:
)
if not page_urls:
self._emit_progress("listing_empty_page_recovery", page_number=page_number)
page_result, page_urls = self._recover_empty_listing_page(
page,
page_number=page_number,
@@ -562,9 +788,10 @@ class IAAIScraper:
)
if not page_urls and page_number > 1:
logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number)
self._emit_progress("listing_resume_start", page_number=page_number)
page.close()
try:
page = self._reopen_listing_and_resume(
page, _ = self._reopen_listing_and_resume(
target_page_number=page_number,
make=make,
model=model,
@@ -579,12 +806,13 @@ class IAAIScraper:
seen_urls,
all_listing_origin_urls,
)
except RuntimeError as resume_exc:
except ListingResumeError 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()
pagination_interrupted = True
page = None
page_urls = []
if not page_urls:
logger.info("Page %d: 0 new links, stopping pagination", page_number)
@@ -648,30 +876,55 @@ class IAAIScraper:
while len(pending_urls) >= batch_size:
batch_urls = pending_urls[:batch_size]
self._emit_progress(
"listing_batch_start",
page_number=page_number,
pending_urls=len(pending_urls),
batch_size=len(batch_urls),
batch_first_url=batch_urls[0] if batch_urls else None,
)
if not _process_pending_batch(batch_urls, cars_upserted + cars_failed):
pending_urls = []
break
pending_urls = pending_urls[batch_size:]
if progress_callback is not None:
try:
progress_callback(page_number)
except Exception:
logger.warning("Failed to persist progress callback for page %d", page_number, exc_info=True)
self._emit_progress(
"listing_batch_done",
page_number=page_number,
pending_urls=len(pending_urls),
cars_upserted=cars_upserted,
cars_failed=cars_failed,
)
if limit is not None and limit > 0 and total >= limit:
break
if (
self.settings.listing.collect_current_page_only
or not self.settings.listing.include_pagination
or not page_result.next_page_detected
):
break
if not self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1):
self._emit_progress("listing_next_page_start", from_page=page_number, to_page=page_number + 1)
# Обработка навигации к следующей странице с защитой от crashes
navigation_success = False
try:
navigation_success = self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1)
except (PlaywrightError, PlaywrightTimeoutError, Exception) as nav_exc:
logger.warning("Exception during page navigation (page %d%d): %s", page_number, page_number + 1, nav_exc)
navigation_success = False
if not navigation_success:
if not page_result.next_page_detected:
logger.info(
"Failed to navigate to page %d with no next-page control detected — treating pagination as exhausted",
page_number + 1,
)
break
logger.warning("Failed to navigate to page %d — reopening listing and resuming", page_number + 1)
self._emit_progress("listing_next_page_failed", from_page=page_number, to_page=page_number + 1)
page.close()
page = None
try:
page = self._reopen_listing_and_resume(
page, _ = self._reopen_listing_and_resume(
target_page_number=page_number + 1,
make=make,
model=model,
@@ -679,16 +932,29 @@ class IAAIScraper:
year_min=year_min,
year_max=year_max,
)
except RuntimeError as resume_exc:
except ListingResumeError as resume_exc:
logger.warning(
"Cannot resume at page %d (%s) — stopping pagination for this segment",
"Cannot resume at page %d (%s) — treating pagination as exhausted",
page_number + 1, resume_exc,
)
page = self._get_page_with_warmup()
pagination_interrupted = True
break
else:
self._emit_progress("listing_next_page_done", from_page=page_number, to_page=page_number + 1)
if pending_urls:
self._emit_progress(
"listing_final_batch_start",
pending_urls=len(pending_urls),
batch_first_url=pending_urls[0] if pending_urls else None,
)
_process_pending_batch(pending_urls, cars_upserted + cars_failed)
self._emit_progress(
"listing_final_batch_done",
pending_urls=0,
cars_upserted=cars_upserted,
cars_failed=cars_failed,
)
logger.warning(
"sync_listing final progress: pages=%d, discovered=%d, upserted=%d, failed=%d",
@@ -711,6 +977,7 @@ class IAAIScraper:
"vehicle_urls": all_raw_urls,
"early_stopped": early_stopped,
"truncated_by_time_budget": truncated_by_time_budget,
"pagination_interrupted": pagination_interrupted,
"pages": pages_info,
},
"total": total,
@@ -1084,6 +1351,26 @@ class IAAIScraper:
]
return "; ".join(pairs)
def _ensure_http_fast_path_session(self) -> None:
if self.context is not None:
try:
if self._cookie_header_for_context():
return
except Exception:
logger.debug("Failed to inspect existing context cookies", exc_info=True)
warmup_page: Page | None = None
try:
warmup_page = self._get_page_with_warmup()
except Exception as exc:
logger.warning("HTTP fast-path warmup failed: %s", exc)
finally:
if warmup_page is not None:
try:
warmup_page.close()
except Exception:
pass
def _scrape_via_raw_http(
self,
vehicle_url: str,
@@ -1177,7 +1464,8 @@ class IAAIScraper:
_HTTP_MAX_RETRIES = 2 # Повторов на URL до перехода в браузер
_HTTP_RETRY_DELAYS = (0.4, 1.0) # Задержки между повторами
_FALLBACK_PARALLEL_PAGES = 4 # Параллельных вкладок для fallback
_FALLBACK_NAV_TIMEOUT_MS = 8000 # Таймаут навигации fallback
# _FALLBACK_NAV_TIMEOUT_MS теперь берётся из settings.fallback_navigation_timeout_ms
# (см. использование ниже). Хардкод 8000 убран — слишком мало при rate-limit от Imperva.
def _browser_fallback_parallel(
self,
@@ -1202,6 +1490,12 @@ class IAAIScraper:
page: Page | None = None
try:
for url, global_idx in url_slice:
self._emit_progress(
"browser_fallback_open_start",
vehicle_url=url,
vehicle_index=global_idx,
vehicle_total=total,
)
if page is None:
page = self._get_page()
if self.settings.block_resources:
@@ -1211,7 +1505,7 @@ class IAAIScraper:
page.goto(
url,
wait_until="commit",
timeout=self._FALLBACK_NAV_TIMEOUT_MS,
timeout=self.settings.fallback_navigation_timeout_ms,
)
except Exception as exc:
if self._is_protection_or_network_error(exc):
@@ -1227,6 +1521,12 @@ class IAAIScraper:
continue
try:
self._emit_progress(
"browser_fallback_parse_start",
vehicle_url=url,
vehicle_index=global_idx,
vehicle_total=total,
)
js_data: dict = {}
try:
js_data = page.evaluate(self._JS_EXTRACT) or {}
@@ -1284,6 +1584,15 @@ class IAAIScraper:
record = CarRecord.model_validate(db_record.model_dump(mode="json"))
local_records.append(record)
self._emit_progress(
"browser_fallback_parse_done",
vehicle_url=url,
vehicle_index=global_idx,
vehicle_total=total,
brand=record.brand,
model=record.model,
year=record.year,
)
logger.debug(
"[%d/%d] Parsed %s %s %s (fallback)",
global_idx, total, record.brand, record.model, record.year or "?",
@@ -1293,6 +1602,13 @@ class IAAIScraper:
local_protection += 1
local_failed += 1
local_failures.append({"vehicle_url": url, "error": str(exc)})
self._emit_progress(
"browser_fallback_parse_failed",
vehicle_url=url,
vehicle_index=global_idx,
vehicle_total=total,
error=str(exc),
)
logger.error("[%d/%d] Failed %s: %s", global_idx, total, url, exc)
finally:
if page is not None:
@@ -1354,7 +1670,15 @@ class IAAIScraper:
total = len(vehicle_urls)
logger.info("sync_batch: %d vehicles, %d parallel workers", total, num_workers)
self._emit_progress(
"sync_batch_start",
batch_size=total,
parallel_workers=num_workers,
first_vehicle_url=vehicle_urls[0] if vehicle_urls else None,
last_vehicle_url=vehicle_urls[-1] if vehicle_urls else None,
)
self._ensure_http_fast_path_session()
cookie_header = self._cookie_header_for_context()
user_agent = (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:128.0) "
@@ -1402,14 +1726,36 @@ class IAAIScraper:
"HTTP phase done: %d OK, %d fallback (%.1fs)",
http_successes, http_fallbacks, time.perf_counter() - started_at,
)
self._emit_progress(
"sync_batch_http_done",
batch_size=total,
http_successes=http_successes,
http_fallbacks=http_fallbacks,
)
# ── Фаза 2: браузерный fallback ──
if fallback_urls:
logger.info(
"Fallback browser mode: %d/%d URLs, single page",
len(fallback_urls), total,
fallback_pages = max(
1,
min(
self._FALLBACK_PARALLEL_PAGES,
num_workers,
len(fallback_urls),
),
)
fb_results = self._browser_fallback_parallel(fallback_urls, total, 1)
logger.info(
"Fallback browser mode: %d/%d URLs, %d page(s)",
len(fallback_urls), total,
fallback_pages,
)
self._emit_progress(
"sync_batch_fallback_start",
fallback_count=len(fallback_urls),
batch_size=total,
first_fallback_url=fallback_urls[0][0] if fallback_urls else None,
fallback_pages=fallback_pages,
)
fb_results = self._browser_fallback_parallel(fallback_urls, total, fallback_pages)
records.extend(fb_results["records"])
cars_failed += fb_results["cars_failed"]
protection_events += fb_results["protection_events"]
@@ -1437,6 +1783,16 @@ class IAAIScraper:
cars_failed += len(records)
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
self._emit_progress(
"sync_batch_done",
batch_size=total,
cars_upserted=cars_upserted,
cars_failed=cars_failed,
images_upserted=images_upserted,
http_successes=http_successes,
http_fallbacks=http_fallbacks,
status=status,
)
return {
"trace_id": trace_id,
"status": status,
@@ -1457,11 +1813,10 @@ class IAAIScraper:
lane: str = "iaai_cars",
limit: int | None = None,
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,
skip_mark_sold: bool = False,
):
# runtime_config — дефолты; CLI/API аргументы приоритетнее.
rc = self.runtime_config.sync
@@ -1488,19 +1843,36 @@ class IAAIScraper:
try:
effective_only_new = self.settings.sync_only_new if only_new is None else only_new
stream_result = self._sync_listing_streaming(
if self._should_use_sitemap_discovery(
make=make,
model=model,
lane=lane,
limit=limit,
effective_only_new=effective_only_new,
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,
)
limit=limit,
effective_only_new=effective_only_new,
):
try:
stream_result = self._sync_listing_via_sitemap(
lane=lane,
limit=limit,
effective_only_new=effective_only_new,
)
except SitemapDiscoveryError as exc:
logger.error("Sitemap-only full scan failed: %s", exc)
raise
else:
stream_result = self._sync_listing_streaming(
make=make,
model=model,
lane=lane,
limit=limit,
effective_only_new=effective_only_new,
started_at=started_at,
listing_url=listing_url,
year_min=year_min,
year_max=year_max,
)
listing = stream_result["listing"]
total = stream_result["total"]
skipped_existing = stream_result["skipped_existing"]
@@ -1518,8 +1890,11 @@ class IAAIScraper:
or (limit is not None and limit > 0)
or listing.get("early_stopped", False)
or listing.get("truncated_by_time_budget", False)
or listing.get("pagination_interrupted", False)
)
if all_listing_origin_urls and not is_partial_scan:
if skip_mark_sold:
logger.debug("Skipping mark_sold: caller requested (parallel segment mode)")
elif all_listing_origin_urls and not is_partial_scan:
try:
sold_count = self.persistence.mark_sold_not_in_listing_by_urls(all_listing_origin_urls)
if sold_count:
@@ -1575,19 +1950,20 @@ class IAAIScraper:
only_new: bool | None = None,
start_segment: int = 0,
start_page: int = 1,
progress_callback: Callable[[int, int], None] | None = None,
progress_callback: Callable[[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) после каждой страницы.
start_page: не используется (остаётся для обратной совместимости API).
progress_callback: вызывается (segment_index) после каждого ПОЛНОСТЬЮ пройденного сегмента.
"""
trace_id = self._new_trace_id("sync-segmented")
started_at = time.perf_counter()
base_url = self.settings.listing.cars_url
_ = start_page # обратная совместимость — page-level resume удалён
total_cars_upserted = 0
total_cars_failed = 0
@@ -1599,8 +1975,8 @@ class IAAIScraper:
completed_all = True
logger.warning(
"Starting segmented sync: %d segments, resume from segment=%d page=%d",
len(segments), start_segment, start_page,
"Starting segmented sync: %d segments, resume from segment=%d",
len(segments), start_segment,
)
for seg_idx in range(start_segment, len(segments)):
@@ -1610,28 +1986,21 @@ class IAAIScraper:
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,
"Segment %d/%d: %s",
seg_idx + 1, len(segments), seg_label,
)
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,
@@ -1646,11 +2015,26 @@ class IAAIScraper:
"segment": seg,
"segment_index": seg_idx,
"status": result.get("status"),
"full_scan_completed": bool(result.get("full_scan_completed", False)),
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"vehicles_collected": result.get("listing", {}).get("vehicles_collected", 0),
})
segment_done = bool(result.get("full_scan_completed", False))
if not segment_done:
completed_all = False
# Сегмент полностью пройден — фиксируем чекпоинт, даже если часть машин failed.
if segment_done and progress_callback is not None:
try:
progress_callback(seg_idx)
except Exception:
logger.warning(
"Failed to persist segment checkpoint for segment=%d",
seg_idx, exc_info=True,
)
logger.warning(
"Segment %d/%d done: %s → upserted=%d, failed=%d, collected=%d",
seg_idx + 1, len(segments), seg_label,

View File

@@ -14,6 +14,13 @@ from .schemas import CarRecord
logger = logging.getLogger("iaai_scraper.db")
def get_db_session() -> Session:
"""Получить новую сессию базы данных."""
settings = Settings()
persistence = PersistenceService(settings)
return persistence.session_factory()
CAR_DB_FIELDS = {
col.key for col in Car.__table__.columns
if col.key not in ("id",)
@@ -517,21 +524,38 @@ class PersistenceService:
is_postgres = "postgresql" in self.settings.database.url
with self.session_scope() as session:
# Страховка от ложного mark_sold при битом/неполном full-scan:
# если текущий список активных URL аномально мал относительно уже активных машин в БД,
# ничего не помечаем проданным.
active_db_count = int(session.execute(
select(func.count())
.select_from(Car)
.where(Car.is_sold == False) # noqa: E712
.where(Car.origin_id.like(f"{lane}:%"))
).scalar_one() or 0)
current_active_count = len(active_origin_urls)
if active_db_count >= 1000 and current_active_count < max(500, int(active_db_count * 0.25)):
logger.warning(
"Skipping mark_sold: suspiciously small active set (%d URLs vs %d active in DB)",
current_active_count,
active_db_count,
)
return 0
if is_postgres:
# Создаём временную таблицу с активными URL.
session.execute(text("CREATE TEMP TABLE IF NOT EXISTS _active_urls (url TEXT NOT NULL) ON COMMIT DROP"))
session.execute(text("CREATE TEMP TABLE IF NOT EXISTS _active_urls (url TEXT PRIMARY KEY) ON COMMIT DROP"))
session.execute(text("TRUNCATE _active_urls"))
# Вставляем активные URL чанками.
# Вставляем активные URL чанками через executemany.
url_list = list(active_origin_urls)
for i in range(0, len(url_list), _IN_CHUNK_SIZE):
chunk = url_list[i:i + _IN_CHUNK_SIZE]
values = ",".join(f"(:{f'u{j}'})" for j in range(len(chunk)))
params = {f"u{j}": url for j, url in enumerate(chunk)}
session.execute(text(f"INSERT INTO _active_urls (url) VALUES {values}"), params)
# Создаём индекс на временной таблице для ускорения JOIN.
session.execute(text("CREATE INDEX IF NOT EXISTS _ix_active_urls ON _active_urls (url)"))
session.execute(
text("INSERT INTO _active_urls (url) VALUES (:url) ON CONFLICT DO NOTHING"),
[{"url": url} for url in chunk],
)
# Массовая пометка проданных в PostgreSQL.
result = session.execute(text("""
@@ -543,10 +567,10 @@ class PersistenceService:
LEFT JOIN _active_urls a ON c.origin_url = a.url
WHERE a.url IS NULL
AND c.is_sold = FALSE
AND c.origin_id LIKE 'iaai:%%'
AND c.origin_id LIKE :lane_prefix
) sub
WHERE cars.id = sub.id
"""))
"""), {"lane_prefix": f"{lane}:%"})
count = result.rowcount or 0
else:
# Упрощённый путь для SQLite.
@@ -554,7 +578,7 @@ class PersistenceService:
update(Car)
.where(Car.origin_url.notin_(active_origin_urls))
.where(Car.is_sold == False) # noqa: E712
.where(Car.origin_id.like("iaai:%"))
.where(Car.origin_id.like(f"{lane}:%"))
.values(is_sold=True)
)
result = session.execute(stmt)
@@ -573,4 +597,43 @@ class PersistenceService:
result = session.execute(
select(Car.origin_id).where(Car.origin_id.like(f"{prefix}%")).execution_options(yield_per=10000)
)
return {str(row[0]) for row in result if row and row[0]}
return {str(row[0]) for row in result if row and row[0]}
def get_all_active_origin_urls_for_lane(self, prefix: str = "iaai:") -> set[str]:
"""Возвращает все активные (не sold) origin_url для указанного lane/prefix."""
with self.session_scope() as session:
result = session.execute(
select(Car.origin_url)
.where(Car.origin_id.like(f"{prefix}%"))
.where(Car.is_sold == False) # noqa: E712
.execution_options(yield_per=10000)
)
return {str(row[0]) for row in result if row and row[0]}
def get_active_origin_urls_batch_for_refresh(
self,
*,
prefix: str = "iaai:",
offset: int = 0,
limit: int = 3000,
) -> list[str]:
"""Возвращает батч активных origin_url для циклического hourly refresh."""
with self.session_scope() as session:
rows = session.execute(
select(Car.origin_url)
.where(Car.origin_id.like(f"{prefix}%"))
.where(Car.is_sold == False) # noqa: E712
.order_by(Car.last_seen_at.asc(), Car.id.asc())
.offset(max(0, int(offset)))
.limit(max(1, int(limit)))
).all()
return [str(row[0]) for row in rows if row and row[0]]
def count_active_cars_for_lane(self, prefix: str = "iaai:") -> int:
with self.session_scope() as session:
return int(session.execute(
select(func.count())
.select_from(Car)
.where(Car.origin_id.like(f"{prefix}%"))
.where(Car.is_sold == False) # noqa: E712
).scalar_one() or 0)

View File

@@ -80,3 +80,65 @@ class SyncRun(Base):
cars_failed: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
images_upserted: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
error_summary: Mapped[str | None] = mapped_column(Text, nullable=True)
class VehicleCandidate(Base):
__tablename__ = "vehicle_candidates"
__table_args__ = (
Index("ix_vehicle_candidates_url", "url"),
Index("ix_vehicle_candidates_status_discovered", "status", "discovered_at"),
)
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
url: Mapped[str] = mapped_column(String(), nullable=False, unique=True, index=True)
discovered_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), index=True)
status: Mapped[str] = mapped_column(String(20), nullable=False, default="pending", index=True) # pending, processing, processed, failed
priority: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
last_attempt_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
attempts: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
raw_snapshots: Mapped[list["VehicleRawSnapshot"]] = relationship("VehicleRawSnapshot", back_populates="candidate", cascade="all, delete-orphan")
class VehicleRawSnapshot(Base):
__tablename__ = "vehicle_raw_snapshots"
__table_args__ = (
Index("ix_vehicle_raw_snapshots_candidate_id", "candidate_id"),
Index("ix_vehicle_raw_snapshots_captured_at", "captured_at"),
)
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
candidate_id: Mapped[int] = mapped_column(Integer, ForeignKey("vehicle_candidates.id", ondelete="CASCADE"), nullable=False, index=True)
candidate: Mapped[VehicleCandidate] = relationship("VehicleCandidate", back_populates="raw_snapshots")
captured_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), index=True)
method: Mapped[str] = mapped_column(String(20), nullable=False) # http, browser
success: Mapped[bool] = mapped_column(Boolean, nullable=False)
raw_data: Mapped[str | None] = mapped_column(Text, nullable=True) # HTML or JSON content
error_message: Mapped[str | None] = mapped_column(Text, nullable=True)
parse_results: Mapped[list["VehicleParseResult"]] = relationship("VehicleParseResult", back_populates="snapshot", cascade="all, delete-orphan")
class VehicleParseResult(Base):
__tablename__ = "vehicle_parse_results"
__table_args__ = (
Index("ix_vehicle_parse_results_snapshot_id", "snapshot_id"),
Index("ix_vehicle_parse_results_parsed_at", "parsed_at"),
)
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
snapshot_id: Mapped[int] = mapped_column(Integer, ForeignKey("vehicle_raw_snapshots.id", ondelete="CASCADE"), nullable=False, index=True)
snapshot: Mapped[VehicleRawSnapshot] = relationship("VehicleRawSnapshot", back_populates="parse_results")
parsed_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), index=True)
success: Mapped[bool] = mapped_column(Boolean, nullable=False)
parsed_data: Mapped[str | None] = mapped_column(Text, nullable=True) # JSON with parsed vehicle data
error_message: Mapped[str | None] = mapped_column(Text, nullable=True)
class VehicleRetryQueue(Base):
__tablename__ = "vehicle_retry_queue"
__table_args__ = (
Index("ix_vehicle_retry_queue_candidate_id", "candidate_id"),
Index("ix_vehicle_retry_queue_retry_at", "retry_at"),
)
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
candidate_id: Mapped[int] = mapped_column(Integer, ForeignKey("vehicle_candidates.id", ondelete="CASCADE"), nullable=False, index=True)
reason: Mapped[str] = mapped_column(String(50), nullable=False) # parse_failed, capture_failed, etc.
retry_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, index=True)
attempts: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
max_attempts: Mapped[int] = mapped_column(Integer, nullable=False, default=3)

View File

@@ -4,11 +4,13 @@ import logging
from celery import Celery
from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging
from redis import Redis
from ..core.config import settings
from ..core.logs import setup_logging
logger = logging.getLogger("iaai_scraper.worker.celery_app")
STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched"
@celery_setup_logging.connect
@@ -96,6 +98,24 @@ celery_app.autodiscover_tasks(["iaai_scraper.worker"])
def _on_worker_ready(**kwargs):
"""Сразу при старте worker отправляем первую задачу sync_listing,
чтобы не ждать час до первого beat-цикла."""
try:
redis_client = Redis.from_url(
settings.redis.url,
decode_responses=True,
socket_connect_timeout=settings.redis.socket_connect_timeout_seconds,
socket_timeout=settings.redis.socket_timeout_seconds,
health_check_interval=settings.redis.health_check_interval_seconds,
retry_on_timeout=True,
)
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
except Exception:
logger.warning("Worker ready startup sync dedupe check failed; skipping immediate dispatch", exc_info=True)
return
if not should_dispatch:
logger.info("Worker ready immediate sync already dispatched recently; skipping duplicate enqueue")
return
logger.info("Worker ready — dispatching initial sync_listing task")
celery_app.send_task(
"iaai_scraper.worker.tasks.sync_listing_task",

File diff suppressed because it is too large Load Diff

View File

@@ -11,6 +11,7 @@ requires-python = ">=3.11"
dependencies = [
"alembic>=1.14.0",
"celery>=5.4.0",
"curl-cffi>=0.11.0",
"fastapi>=0.115.0",
"playwright>=1.53.0",
"psycopg2-binary>=2.9.9",
@@ -20,6 +21,8 @@ dependencies = [
"sqlalchemy>=2.0.32",
"uvicorn>=0.34.0",
"urllib3>=2.0.0",
"orjson>=3.10.0",
"brotli>=1.1.0",
]
[project.optional-dependencies]

View File

@@ -8,26 +8,83 @@ from iaai_scraper.browser.listing import ListingCollector
class _FakePage:
def __init__(self, counts: dict[str, int]) -> None:
def __init__(self, counts: dict[str, int], attrs: dict[str, dict[str, str]] | None = None) -> None:
self._counts = counts
self._attrs = attrs or {}
self._evaluate_result = False
self._evaluate_values: list[object] = []
self._html = ""
self.url = "https://www.iaai.com/Vehiclelisting/Cars"
self._current_page_number: int | None = None
self._vehicle_hrefs: list[str] = []
self._wait_for_function_error: Exception | None = None
self._clicks: dict[str, int] = {}
self._goto_calls: list[str] = []
class _Locator:
def __init__(self, count_value: int) -> None:
def __init__(self, page: "_FakePage", selector: str, count_value: int) -> None:
self._page = page
self._selector = selector
self._count_value = count_value
@property
def first(self) -> "_FakePage._Locator":
return self
def count(self) -> int:
return self._count_value
def locator(self, selector: str) -> "_FakePage._Locator":
return _FakePage._Locator(self._counts.get(selector, 0))
def is_visible(self, timeout: int | None = None) -> bool:
_ = timeout
return False
def evaluate(self, _script: str):
def get_attribute(self, name: str, timeout: int | None = None) -> str | None:
_ = timeout
return self._page._attrs.get(self._selector, {}).get(name)
def click(self, timeout: int | None = None) -> None:
_ = timeout
self._page._clicks[self._selector] = self._page._clicks.get(self._selector, 0) + 1
return None
class _HtmlLocator:
def __init__(self, html: str) -> None:
self._html = html
def inner_html(self, timeout: int | None = None) -> str:
_ = timeout
return self._html
def locator(self, selector: str):
if selector == "html":
return _FakePage._HtmlLocator(self._html)
return _FakePage._Locator(self, selector, self._counts.get(selector, 0))
def evaluate(self, script: str, *args):
_ = args
if self._evaluate_values:
return self._evaluate_values.pop(0)
if "querySelectorAll" in script and "VehicleDetail" in script:
return list(self._vehicle_hrefs)
if "document.body?.innerText" in script and "aria-current" in script and "parseInt" in script:
return self._current_page_number if self._current_page_number is not None else self._evaluate_result
return self._evaluate_result
def wait_for_selector(self, selector: str, timeout: int | None = None) -> None:
_ = selector, timeout
return None
def wait_for_function(self, script: str, timeout: int | None = None) -> None:
_ = script, timeout
if self._wait_for_function_error is not None:
raise self._wait_for_function_error
return None
def goto(self, url: str, wait_until: str | None = None, timeout: int | None = None) -> None:
_ = wait_until, timeout
self._goto_calls.append(url)
self.url = url
class TestListingUnit(unittest.TestCase):
def test_pagination_detection_and_page_number(self) -> None:
@@ -35,7 +92,7 @@ class TestListingUnit(unittest.TestCase):
self.assertTrue(ListingCollector._has_next_page(_FakePage({"a[aria-label*='Next']": 1})))
# Нет селекторов → False.
self.assertFalse(ListingCollector._has_next_page(_FakePage({})))
# Числовая пагинация через JS → True.
# Числовая пагинация через JS → True только если evaluate явно нашёл next.
page = _FakePage({})
page._evaluate_result = True
self.assertTrue(ListingCollector._has_next_page(page))
@@ -63,6 +120,67 @@ class TestListingUnit(unittest.TestCase):
],
)
def test_has_next_page_from_html(self) -> None:
collector = ListingCollector(Settings(), HumanPacer(Settings()))
html = '<a class="pagination-next" href="/Vehiclelisting/Cars?page=43">Next</a>'
self.assertTrue(collector._has_next_page_from_html(html, current_page_number=42))
self.assertFalse(collector._has_next_page_from_html("<div>done</div>", current_page_number=42))
def test_build_listing_page_url_replaces_existing_page_parameter(self) -> None:
collector = ListingCollector(Settings(), HumanPacer(Settings()))
url = collector._build_listing_page_url(
"https://www.iaai.com/Vehiclelisting/Cars?Make=TOYOTA&page=43&foo=bar",
44,
)
self.assertEqual(
url,
"https://www.iaai.com/Vehiclelisting/Cars?Make=TOYOTA&foo=bar&page=44",
)
def test_collect_current_page_uses_html_first(self) -> None:
collector = ListingCollector(Settings(), HumanPacer(Settings()))
page = _FakePage({})
page._html = '<a href="/VehicleDetail/555~US">Car 555</a><a class="pagination-next" href="/Vehiclelisting/Cars?page=2">Next</a>'
result = collector.collect_current_page(page, page_number=1)
self.assertEqual(len(result.vehicle_links), 1)
self.assertEqual(result.vehicle_links[0].lot_number, "555")
self.assertTrue(result.next_page_detected)
def test_wait_for_navigation_result_rejects_duplicate_content_even_when_page_number_matches(self) -> None:
page = _FakePage({})
page._wait_for_function_error = RuntimeError("same content")
page._current_page_number = 41
page._vehicle_hrefs = ["/VehicleDetail/111~US", "/VehicleDetail/222~US"]
self.assertFalse(
ListingCollector._wait_for_navigation_result(
page,
"/VehicleDetail/111~US",
41,
old_fingerprint=("/VehicleDetail/111~US", "/VehicleDetail/222~US"),
)
)
def test_go_to_next_page_skips_flaky_ui_fallbacks_on_deep_pages(self) -> None:
collector = ListingCollector(Settings(), HumanPacer(Settings()))
page = _FakePage(
{ListingCollector._NEXT_PAGE_SELECTORS[0]: 1, "a[href*='/VehicleDetail/'], a[href*='/vehicledetail/'], a[href*='VehicleDetail'], a[href*='vehicledetail']": 1},
attrs={
"a[href*='/VehicleDetail/'], a[href*='/vehicledetail/'], a[href*='VehicleDetail'], a[href*='vehicledetail']": {"href": "/VehicleDetail/111~US"},
},
)
page._wait_for_function_error = RuntimeError("same content")
page._current_page_number = 41
page._vehicle_hrefs = ["/VehicleDetail/111~US", "/VehicleDetail/222~US"]
self.assertFalse(collector.go_to_next_page(page, expected_page_number=41))
self.assertGreaterEqual(len(page._goto_calls), 1)
self.assertEqual(page._clicks.get(ListingCollector._NEXT_PAGE_SELECTORS[0], 0), 0)
if __name__ == "__main__":
unittest.main()

View File

@@ -6,6 +6,8 @@ from unittest.mock import MagicMock
from iaai_scraper.core.config import Settings
from iaai_scraper.core.exceptions import AntiBotDetectedError, SiteStructureChangedError
from iaai_scraper.core.exceptions import ListingResumeError
from iaai_scraper.discovery import SitemapDiscoveryError
from iaai_scraper.scraper import IAAIScraper
from iaai_scraper.storage.schemas import CarRecord
@@ -85,6 +87,51 @@ class TestScraperSync(unittest.TestCase):
scraper2._sync_listing_streaming.assert_called_once()
scraper2.collect_listing.assert_not_called()
def test_full_scan_uses_sitemap_discovery_path(self) -> None:
scraper = self._make_scraper()
scraper.settings.discovery.mode = "listing"
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=88)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.mark_sold_not_in_listing_by_urls = MagicMock(return_value=0)
scraper._sync_listing_via_sitemap = MagicMock(return_value={
"listing": {
"vehicles_collected": 123,
"early_stopped": False,
"truncated_by_time_budget": False,
"pagination_interrupted": False,
"source": "sitemap",
},
"total": 123,
"skipped_existing": 0,
"cars_upserted": 120,
"cars_failed": 3,
"images_upserted": 500,
"failures": [{"vehicle_url": "v", "error": "e"}],
"all_listing_origin_urls": {"https://www.iaai.com/VehicleDetail/111~US"},
})
scraper._sync_listing_streaming = MagicMock(side_effect=AssertionError("pagination path should not be used"))
result = scraper.sync_listing()
scraper._sync_listing_via_sitemap.assert_called_once()
scraper._sync_listing_streaming.assert_not_called()
self.assertEqual(result["cars_upserted"], 120)
self.assertTrue(result["full_scan_completed"])
def test_full_scan_does_not_fallback_to_pagination_when_sitemap_fails(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=89)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.mark_sold_not_in_listing_by_urls = MagicMock(return_value=0)
scraper._sync_listing_via_sitemap = MagicMock(side_effect=SitemapDiscoveryError("sitemap down"))
result = scraper.sync_listing()
scraper._sync_listing_via_sitemap.assert_called_once()
self.assertEqual(result["status"], "failed")
self.assertEqual(result["cars_upserted"], 0)
def test_segmented_sync_calls_per_segment_and_resumes(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
@@ -107,20 +154,52 @@ class TestScraperSync(unittest.TestCase):
{"make": "HONDA", "year_min": None, "year_max": None},
]
# Resume: пропускаем TOYOTA, начинаем с FORD на стр. 5.
# Resume: пропускаем TOYOTA, начинаем сразу с FORD.
result = scraper.sync_listing_segmented(
segments=segments, start_segment=1, start_page=5,
segments=segments, start_segment=1,
)
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_segmented_sync_does_not_mark_full_scan_completed_when_segment_failed(self) -> None:
scraper = self._make_scraper()
scraper.sync_listing = MagicMock(side_effect=[
{
"status": "failed",
"full_scan_completed": False,
"cars_upserted": 0,
"cars_failed": 0,
"images_upserted": 0,
"skipped_existing": 0,
"listing": {"vehicles_collected": 0},
"failures": [{"vehicle_url": "segment", "error": "resume failed"}],
},
{
"status": "success",
"full_scan_completed": True,
"cars_upserted": 1,
"cars_failed": 0,
"images_upserted": 0,
"skipped_existing": 0,
"listing": {"vehicles_collected": 10},
"failures": [],
},
])
result = scraper.sync_listing_segmented(
segments=[
{"make": "EAGLE", "year_min": None, "year_max": None},
{"make": "FORD", "year_min": None, "year_max": None},
]
)
self.assertFalse(result["full_scan_completed"])
self.assertEqual(result["status"], "partial_success")
def test_build_segment_listing_url(self) -> None:
base = "https://www.iaai.com/Vehiclelisting/Cars"
self.assertEqual(
@@ -188,6 +267,174 @@ class TestScraperSync(unittest.TestCase):
scraper.listing_collector.open_cars_listing.assert_called_once()
page.reload.assert_not_called()
def test_sync_listing_always_starts_from_page_one(self) -> None:
scraper = self._make_scraper()
scraper.settings.celery.batch_size = 1
page = MagicMock()
page_result = SimpleNamespace(
page_number=1,
vehicle_links=[SimpleNamespace(href="https://www.iaai.com/VehicleDetail/999~US", lot_number="999")],
next_page_detected=False,
)
scraper._get_page_with_warmup = MagicMock(return_value=page)
# Pagination-resume убран: _reopen_listing_and_resume не должен вызываться.
scraper._reopen_listing_and_resume = MagicMock(
side_effect=AssertionError("pagination resume must not be used")
)
scraper.listing_collector.open_cars_listing = MagicMock()
scraper.listing_collector.apply_filters = MagicMock(return_value={
"make": None,
"model": None,
"year_min": None,
"year_max": None,
})
scraper.listing_collector.collect_current_page = MagicMock(return_value=page_result)
scraper._extract_page_urls = MagicMock(return_value=["https://www.iaai.com/VehicleDetail/999~US"])
scraper.sync_batch = MagicMock(return_value={
"cars_upserted": 1,
"cars_failed": 0,
"images_upserted": 0,
"failures": [],
})
result = scraper._sync_listing_streaming(
make=None,
model=None,
lane="iaai_cars",
limit=None,
effective_only_new=False,
started_at=0.0,
listing_url="https://www.iaai.com/Vehiclelisting/Cars?Make=EAGLE",
)
scraper.listing_collector.open_cars_listing.assert_called_once()
scraper._reopen_listing_and_resume.assert_not_called()
scraper.sync_batch.assert_called_once()
self.assertEqual(result["cars_upserted"], 1)
self.assertEqual(result["total"], 1)
def test_sync_listing_marks_pagination_interrupted_when_next_page_resume_fails(self) -> None:
scraper = self._make_scraper()
scraper.settings.celery.batch_size = 1000
page = MagicMock()
page_result = SimpleNamespace(
page_number=1,
vehicle_links=[SimpleNamespace(href="https://www.iaai.com/VehicleDetail/999~US", lot_number="999")],
next_page_detected=True,
)
scraper._get_page_with_warmup = MagicMock(return_value=page)
scraper.listing_collector.open_cars_listing = MagicMock()
scraper.listing_collector.apply_filters = MagicMock(return_value={
"make": None,
"model": None,
"year_min": None,
"year_max": None,
})
scraper.listing_collector.collect_current_page = MagicMock(return_value=page_result)
scraper.listing_collector.go_to_next_page = MagicMock(return_value=False)
scraper._extract_page_urls = MagicMock(return_value=["https://www.iaai.com/VehicleDetail/999~US"])
scraper._reopen_listing_and_resume = MagicMock(side_effect=ListingResumeError("resume failed"))
scraper.sync_batch = MagicMock(return_value={
"cars_upserted": 0,
"cars_failed": 0,
"images_upserted": 0,
"failures": [],
})
result = scraper._sync_listing_streaming(
make=None,
model=None,
lane="iaai_cars",
limit=None,
effective_only_new=False,
started_at=0.0,
listing_url="https://www.iaai.com/Vehiclelisting/Cars?Make=TOYOTA",
)
self.assertTrue(result["listing"]["pagination_interrupted"])
def test_sync_listing_attempts_next_page_even_without_detected_next_control(self) -> None:
scraper = self._make_scraper()
scraper.settings.celery.batch_size = 1000
page = MagicMock()
first_page = SimpleNamespace(
page_number=1,
vehicle_links=[SimpleNamespace(href="https://www.iaai.com/VehicleDetail/111~US", lot_number="111")],
next_page_detected=False,
)
second_page = SimpleNamespace(
page_number=2,
vehicle_links=[SimpleNamespace(href="https://www.iaai.com/VehicleDetail/222~US", lot_number="222")],
next_page_detected=False,
)
scraper._get_page_with_warmup = MagicMock(return_value=page)
scraper.listing_collector.open_cars_listing = MagicMock()
scraper.listing_collector.apply_filters = MagicMock(return_value={
"make": None,
"model": None,
"year_min": None,
"year_max": None,
})
scraper.listing_collector.collect_current_page = MagicMock(side_effect=[first_page, second_page])
scraper.listing_collector.go_to_next_page = MagicMock(side_effect=[True, False])
scraper._extract_page_urls = MagicMock(side_effect=[
["https://www.iaai.com/VehicleDetail/111~US"],
["https://www.iaai.com/VehicleDetail/222~US"],
])
scraper.sync_batch = MagicMock(return_value={
"cars_upserted": 2,
"cars_failed": 0,
"images_upserted": 0,
"failures": [],
})
result = scraper._sync_listing_streaming(
make=None,
model=None,
lane="iaai_cars",
limit=None,
effective_only_new=False,
started_at=0.0,
listing_url="https://www.iaai.com/Vehiclelisting/Cars?Make=TOYOTA",
)
self.assertEqual(scraper.listing_collector.go_to_next_page.call_count, 2)
scraper.listing_collector.go_to_next_page.assert_any_call(page, expected_page_number=2)
self.assertEqual(result["listing"]["pages_collected"], 2)
def test_sync_listing_treats_pagination_interrupted_as_partial_scan(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=77)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.mark_sold_not_in_listing_by_urls = MagicMock()
scraper._sync_listing_streaming = MagicMock(return_value={
"listing": {
"vehicles_collected": 10,
"early_stopped": False,
"truncated_by_time_budget": False,
"pagination_interrupted": True,
},
"total": 10,
"skipped_existing": 0,
"cars_upserted": 10,
"cars_failed": 0,
"images_upserted": 0,
"failures": [],
"all_listing_origin_urls": {"https://www.iaai.com/VehicleDetail/999~US"},
})
result = scraper.sync_listing(listing_url="https://www.iaai.com/Vehiclelisting/Cars?Make=TOYOTA")
self.assertTrue(result["full_scan_completed"] is False)
scraper.persistence.mark_sold_not_in_listing_by_urls.assert_not_called()
if __name__ == "__main__":
unittest.main()

View File

@@ -0,0 +1,217 @@
from __future__ import annotations
import unittest
from unittest.mock import patch
from iaai_scraper.core.config import Settings
from iaai_scraper.discovery.sitemap import (
SitemapBlockedError,
SitemapDiscoveryError,
SitemapFetchResult,
_filter_vehicle_urls,
discover_vehicle_urls_from_sitemap,
discover_vehicle_urls_from_sitemap_with_stats,
)
class TestSitemapDiscovery(unittest.TestCase):
@staticmethod
def _fetch_result(url: str, payload: bytes, *, source: str = "curl_cffi") -> SitemapFetchResult:
return SitemapFetchResult(url=url, payload=payload, source=source)
def test_filter_vehicle_urls_dedupes_and_normalizes(self) -> None:
urls = _filter_vehicle_urls([
"https://www.iaai.com/VehicleDetail/111~US?foo=1",
"https://www.iaai.com/VehicleDetail/111~US?bar=2",
"https://www.iaai.com/VehicleDetail/222~US",
"https://www.iaai.com/about",
])
self.assertEqual(
urls,
[
"https://www.iaai.com/VehicleDetail/111~US",
"https://www.iaai.com/VehicleDetail/222~US",
],
)
def test_discover_vehicle_urls_from_sitemap(self) -> None:
index_xml = b"""
<sitemapindex xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\">
<sitemap><loc>https://www.iaai.com/sitemap-a.xml</loc></sitemap>
<sitemap><loc>https://www.iaai.com/sitemap-b.xml</loc></sitemap>
<sitemap><loc>https://www.iaai.com/sitemapauctions1.xml</loc></sitemap>
</sitemapindex>
"""
sitemap_a = b"""
<urlset xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\">
<url><loc>https://www.iaai.com/VehicleDetail/111~US</loc></url>
<url><loc>https://www.iaai.com/VehicleDetail/222~US?x=1</loc></url>
</urlset>
"""
sitemap_b = b"""
<urlset xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\">
<url><loc>https://www.iaai.com/VehicleDetail/222~US?y=2</loc></url>
<url><loc>https://www.iaai.com/VehicleDetail/333~US</loc></url>
</urlset>
"""
def _fake_fetch(url: str) -> SitemapFetchResult:
if url.endswith("sitemap_index.xml"):
return self._fetch_result(url, index_xml)
if url.endswith("sitemap-a.xml"):
return self._fetch_result(url, sitemap_a)
if url.endswith("sitemap-b.xml"):
return self._fetch_result(url, sitemap_b)
raise AssertionError(f"unexpected url: {url}")
with patch("iaai_scraper.discovery.sitemap._SitemapDownloader.fetch", side_effect=_fake_fetch):
result = discover_vehicle_urls_from_sitemap("https://www.iaai.com/sitemap_index.xml")
self.assertEqual(
result,
[
"https://www.iaai.com/VehicleDetail/111~US",
"https://www.iaai.com/VehicleDetail/222~US",
"https://www.iaai.com/VehicleDetail/333~US",
],
)
def test_discover_vehicle_urls_raises_on_empty_index(self) -> None:
empty_index = b"<sitemapindex xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\"></sitemapindex>"
with patch(
"iaai_scraper.discovery.sitemap._SitemapDownloader.fetch",
return_value=self._fetch_result("https://www.iaai.com/sitemap_index.xml", empty_index),
):
with self.assertRaises(SitemapDiscoveryError):
discover_vehicle_urls_from_sitemap("https://www.iaai.com/sitemap_index.xml")
def test_discover_vehicle_urls_skips_malformed_and_non_vehicle_sitemaps(self) -> None:
sitemap_vehicle = b"""
<urlset xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\">
<url><loc>https://www.iaai.com/VehicleDetail/111~US</loc></url>
</urlset>
"""
sitemap_non_vehicle = b"""
<urlset xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\">
<url><loc>https://www.iaai.com/SalesList/111~US/04222026</loc></url>
</urlset>
"""
index_xml = b"""
<sitemapindex xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\">
<sitemap><loc>https://www.iaai.com/sitemap1.xml</loc></sitemap>
<sitemap><loc>https://www.iaai.com/sitemapbranches1.xml</loc></sitemap>
<sitemap><loc>https://www.iaai.com/sitemap2.xml</loc></sitemap>
<sitemap><loc>https://www.iaai.com/sitemap3.xml</loc></sitemap>
</sitemapindex>
"""
def _fake_fetch(url: str) -> SitemapFetchResult:
if url.endswith("sitemap_index.xml"):
return self._fetch_result(url, index_xml)
if url.endswith("sitemap1.xml"):
return self._fetch_result(url, sitemap_vehicle)
if url.endswith("sitemap2.xml"):
return self._fetch_result(url, sitemap_non_vehicle)
if url.endswith("sitemap3.xml"):
raise SitemapDiscoveryError("broken sitemap")
raise AssertionError(f"unexpected url: {url}")
with patch("iaai_scraper.discovery.sitemap._SitemapDownloader.fetch", side_effect=_fake_fetch):
result = discover_vehicle_urls_from_sitemap("https://www.iaai.com/sitemap_index.xml")
self.assertEqual(result, ["https://www.iaai.com/VehicleDetail/111~US"])
def test_discover_vehicle_urls_falls_back_to_regex_loc_extraction(self) -> None:
malformed_index = (
b"<sitemapindex>"
b"<sitemap><loc>https://www.iaai.com/sitemap1.xml</loc></sitemap>"
b"<broken"
)
sitemap_vehicle = b"""
<urlset xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\">
<url><loc>https://www.iaai.com/VehicleDetail/111~US</loc></url>
</urlset>
"""
def _fake_fetch(url: str) -> SitemapFetchResult:
if url.endswith("sitemap_index.xml"):
return self._fetch_result(url, malformed_index)
if url.endswith("sitemap1.xml"):
return self._fetch_result(url, sitemap_vehicle)
raise AssertionError(f"unexpected url: {url}")
with patch("iaai_scraper.discovery.sitemap._SitemapDownloader.fetch", side_effect=_fake_fetch):
result = discover_vehicle_urls_from_sitemap("https://www.iaai.com/sitemap_index.xml")
self.assertEqual(result, ["https://www.iaai.com/VehicleDetail/111~US"])
def test_discovery_uses_direct_probe_when_index_has_no_urls(self) -> None:
index_xml = b"<sitemapindex xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\"></sitemapindex>"
sitemap_one = b"""
<urlset xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\">
<url><loc>https://www.iaai.com/VehicleDetail/111~US</loc></url>
</urlset>
"""
settings = Settings()
settings.discovery.sitemap_direct_probe_limit = 2
settings.discovery.sitemap_direct_probe_stop_after_misses = 1
def _fake_fetch(url: str) -> SitemapFetchResult:
if url.endswith("sitemap_index.xml"):
return self._fetch_result(url, index_xml)
if url.endswith("sitemap1.xml"):
return self._fetch_result(url, sitemap_one)
raise SitemapDiscoveryError("not found")
with patch("iaai_scraper.discovery.sitemap._SitemapDownloader.fetch", side_effect=_fake_fetch):
result = discover_vehicle_urls_from_sitemap_with_stats(
"https://www.iaai.com/sitemap_index.xml",
settings=settings,
)
self.assertEqual(result.vehicle_urls, ["https://www.iaai.com/VehicleDetail/111~US"])
self.assertEqual(result.stats.direct_probe_hits, 1)
def test_discovery_stats_report_direct_probe_hits(self) -> None:
index_xml = b"<sitemapindex xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\"></sitemapindex>"
sitemap1 = b"""
<urlset xmlns=\"http://www.sitemaps.org/schemas/sitemap/0.9\">
<url><loc>https://www.iaai.com/VehicleDetail/111~US</loc></url>
</urlset>
"""
def _fake_fetch(url: str):
if url.endswith("sitemap_index.xml"):
return self._fetch_result(url, index_xml)
if url.endswith("sitemap1.xml"):
return self._fetch_result(url, sitemap1)
raise SitemapDiscoveryError("not found")
with patch("iaai_scraper.discovery.sitemap._SitemapDownloader.fetch", side_effect=_fake_fetch):
result = discover_vehicle_urls_from_sitemap_with_stats("https://www.iaai.com/sitemap_index.xml")
self.assertEqual(result.vehicle_urls, ["https://www.iaai.com/VehicleDetail/111~US"])
self.assertEqual(result.stats.transport, "curl_cffi")
self.assertEqual(result.stats.direct_probe_hits, 1)
self.assertGreaterEqual(result.stats.direct_probe_misses, 1)
def test_block_page_raises_discovery_error(self) -> None:
settings = Settings()
settings.discovery.sitemap_direct_probe_limit = 1
settings.discovery.sitemap_direct_probe_stop_after_misses = 1
with patch(
"iaai_scraper.discovery.sitemap._SitemapDownloader.fetch",
side_effect=SitemapBlockedError("Anti-bot page returned for https://www.iaai.com/sitemap_index.xml"),
):
with self.assertRaises(SitemapDiscoveryError):
discover_vehicle_urls_from_sitemap(
"https://www.iaai.com/sitemap_index.xml",
settings=settings,
)
if __name__ == "__main__":
unittest.main()

View File

@@ -1,13 +1,215 @@
from __future__ import annotations
import json
import unittest
from unittest.mock import MagicMock, patch
from iaai_scraper.worker import tasks
def make_settings_stub(*, parallel_segments: bool = False) -> MagicMock:
settings_stub = MagicMock()
settings_stub.celery.parallel_segments = parallel_segments
settings_stub.celery.task_soft_time_limit = 3300
settings_stub.celery.task_time_limit = 3600
settings_stub.celery.task_stall_timeout_seconds = 600
settings_stub.listing.listing_segments_json = "ignored"
settings_stub.discovery.mode = "listing"
settings_stub.discovery.hourly_mode = "rolling_refresh"
return settings_stub
class TestWorkerTaskLockHelpers(unittest.TestCase):
def test_sync_listing_task_uses_hourly_sitemap_even_when_mode_is_not_sitemap(self) -> None:
settings_stub = make_settings_stub(parallel_segments=False)
settings_stub.discovery.mode = "listing"
settings_stub.discovery.hourly_mode = "rolling_refresh"
with patch.object(tasks, "_get_persistence") as get_persistence, \
patch.object(tasks, "_get_redis") as get_redis, \
patch.object(tasks, "Settings", return_value=settings_stub), \
patch.object(tasks, "_acquire_lock", return_value=True), \
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, "_hourly_sitemap_rolling_refresh_sync", return_value={
"status": "success",
"run_id": None,
"cars_upserted": 25,
"cars_failed": 0,
"images_upserted": 50,
"skipped_existing": 100,
"elapsed_seconds": None,
"failures": [],
"hourly_mode": "sitemap_rolling_refresh",
"discovered_urls": 105,
"new_urls": 5,
"refresh_urls": 20,
"sold_marked": 2,
}) as hourly_sync, \
patch.object(tasks.sync_listing_task, "update_state"):
get_persistence.return_value = MagicMock()
redis_client = MagicMock()
redis_client.get.return_value = None
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
tasks.sync_listing_task.push_request(id="task-hourly-force-sitemap")
try:
result = tasks.sync_listing_task.run()
finally:
tasks.sync_listing_task.pop_request()
self.assertEqual(result["status"], "success")
self.assertEqual(result["hourly_mode"], "sitemap_rolling_refresh")
hourly_sync.assert_called_once()
release_lock.assert_called_once()
def test_sync_listing_task_uses_hourly_sitemap_rolling_refresh_after_bootstrap(self) -> None:
settings_stub = make_settings_stub(parallel_segments=False)
settings_stub.discovery.mode = "sitemap"
settings_stub.discovery.hourly_mode = "rolling_refresh"
with patch.object(tasks, "_get_persistence") as get_persistence, \
patch.object(tasks, "_get_redis") as get_redis, \
patch.object(tasks, "Settings", return_value=settings_stub), \
patch.object(tasks, "_acquire_lock", return_value=True), \
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, "_hourly_sitemap_rolling_refresh_sync", return_value={
"status": "success",
"run_id": None,
"cars_upserted": 25,
"cars_failed": 0,
"images_upserted": 50,
"skipped_existing": 100,
"elapsed_seconds": None,
"failures": [],
"hourly_mode": "sitemap_rolling_refresh",
"discovered_urls": 105,
"new_urls": 5,
"refresh_urls": 20,
"sold_marked": 2,
}) as hourly_sync, \
patch.object(tasks.sync_listing_task, "update_state"):
get_persistence.return_value = MagicMock()
redis_client = MagicMock()
redis_client.get.return_value = None
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
tasks.sync_listing_task.push_request(id="task-hourly")
try:
result = tasks.sync_listing_task.run()
finally:
tasks.sync_listing_task.pop_request()
self.assertEqual(result["status"], "success")
self.assertEqual(result["hourly_mode"], "sitemap_rolling_refresh")
hourly_sync.assert_called_once()
release_lock.assert_called_once()
def test_sync_listing_task_can_use_hourly_sitemap_diff_when_configured(self) -> None:
settings_stub = make_settings_stub(parallel_segments=False)
settings_stub.discovery.mode = "sitemap"
settings_stub.discovery.hourly_mode = "diff"
with patch.object(tasks, "_get_persistence") as get_persistence, \
patch.object(tasks, "_get_redis") as get_redis, \
patch.object(tasks, "Settings", return_value=settings_stub), \
patch.object(tasks, "_acquire_lock", return_value=True), \
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, "_hourly_sitemap_diff_sync", return_value={
"status": "success",
"run_id": None,
"cars_upserted": 5,
"cars_failed": 0,
"images_upserted": 10,
"skipped_existing": 100,
"elapsed_seconds": None,
"failures": [],
"hourly_mode": "sitemap_diff",
"discovered_urls": 105,
"new_urls": 5,
"sold_marked": 2,
}) as hourly_sync, \
patch.object(tasks.sync_listing_task, "update_state"):
get_persistence.return_value = MagicMock()
redis_client = MagicMock()
redis_client.get.return_value = None
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
tasks.sync_listing_task.push_request(id="task-hourly-diff")
try:
result = tasks.sync_listing_task.run()
finally:
tasks.sync_listing_task.pop_request()
self.assertEqual(result["status"], "success")
self.assertEqual(result["hourly_mode"], "sitemap_diff")
hourly_sync.assert_called_once()
release_lock.assert_called_once()
def test_sync_listing_task_forces_sitemap_full_scan_in_bootstrap(self) -> None:
settings_stub = make_settings_stub(parallel_segments=False)
settings_stub.discovery.mode = "listing"
with patch.object(tasks, "_get_persistence") as get_persistence, \
patch.object(tasks, "_get_redis") as get_redis, \
patch.object(tasks, "Settings", return_value=settings_stub), \
patch.object(tasks, "_acquire_lock", return_value=True), \
patch.object(tasks, "_is_full_scan_done", return_value=False), \
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
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=[{"make": "HONDA"}]):
get_persistence.return_value = MagicMock()
redis_client = MagicMock()
redis_client.get.return_value = None
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
sync_listing_mock = MagicMock(return_value={
"run_id": 99,
"status": "success",
"full_scan_completed": True,
"cars_upserted": 10,
"cars_failed": 0,
"images_upserted": 20,
"skipped_existing": 0,
"elapsed_seconds": 1.0,
"failures": [],
})
sync_listing_segmented_mock = MagicMock(side_effect=AssertionError("segmented path should not be used"))
scraper = MagicMock()
scraper.sync_listing = sync_listing_mock
scraper.sync_listing_segmented = sync_listing_segmented_mock
scraper_ctx = MagicMock()
scraper_ctx.__enter__.return_value = scraper
scraper_ctx.__exit__.return_value = None
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
tasks.sync_listing_task.push_request(id="task-bootstrap-sitemap")
try:
result = tasks.sync_listing_task.run()
finally:
tasks.sync_listing_task.pop_request()
self.assertEqual(result["status"], "success")
sync_listing_mock.assert_called_once_with(
make=None,
model=None,
lane="iaai_cars",
limit=None,
only_new=False,
)
sync_listing_segmented_mock.assert_not_called()
release_lock.assert_called_once()
def test_lock_acquire_refresh_release(self) -> None:
redis_client = MagicMock()
@@ -146,22 +348,191 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
self.assertEqual(acquire_lock.call_count, 2)
release_lock.assert_called_once()
def test_sync_listing_checkpoint_save_load_and_resume(self) -> None:
# Roundtrip: save → load.
redis_client = MagicMock()
def test_checkpoint_save_load_roundtrip(self) -> None:
storage: dict[str, str] = {}
redis_client.set.side_effect = lambda key, value: storage.__setitem__(key, value)
def _fake_set(key, value, ex=None):
storage[key] = value
return True
redis_client = MagicMock()
redis_client.set.side_effect = _fake_set
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",
)
checkpoint = tasks._load_sync_checkpoint(redis_client)
self.assertIsNotNone(checkpoint)
self.assertEqual(checkpoint["last_successful_page"], 12)
tasks._save_last_completed_segment(redis_client, 12)
# Resume: задача стартует со страницы checkpoint + 1.
self.assertEqual(tasks._load_last_completed_segment(redis_client), 12)
self.assertEqual(redis_client.set.call_args.kwargs["ex"], tasks.SYNC_LISTING_CHECKPOINT_TTL_SECONDS)
def test_load_checkpoint_clears_invalid_payload(self) -> None:
redis_client = MagicMock()
redis_client.get.return_value = "{not-a-number"
self.assertIsNone(tasks._load_last_completed_segment(redis_client))
redis_client.delete.assert_called_once_with(tasks.SYNC_LISTING_CHECKPOINT_KEY)
def test_load_checkpoint_empty_when_missing(self) -> None:
redis_client = MagicMock()
redis_client.get.return_value = None
self.assertIsNone(tasks._load_last_completed_segment(redis_client))
redis_client.delete.assert_not_called()
def test_sync_listing_bootstrap_prefers_sitemap_over_segment_resume(self) -> None:
segments = [
{"make": "ACURA"},
{"make": "AUDI"},
{"make": "BMW"},
{"make": "EAGLE"},
]
settings_stub = make_settings_stub()
with patch.object(tasks, "_get_persistence") as get_persistence, \
patch.object(tasks, "_get_redis") as get_redis, \
patch.object(tasks, "Settings", return_value=settings_stub), \
patch.object(tasks, "_acquire_lock", return_value=True), \
patch.object(tasks, "_is_full_scan_done", return_value=False), \
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
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=segments):
get_persistence.return_value = MagicMock()
redis_client = MagicMock()
redis_client.get.side_effect = lambda key: (
"1" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
)
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
sync_segmented_mock = MagicMock(side_effect=AssertionError("segmented path should not be used"))
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": [],
})
scraper_ctx = MagicMock()
scraper_ctx.__enter__.return_value.sync_listing_segmented = sync_segmented_mock
scraper_ctx.__enter__.return_value.sync_listing = sync_listing_mock
scraper_ctx.__exit__.return_value = None
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
tasks.sync_listing_task.push_request(id="task-resume-seg")
try:
result = tasks.sync_listing_task.run()
finally:
tasks.sync_listing_task.pop_request()
self.assertEqual(result["status"], "success")
sync_listing_mock.assert_called_once_with(
make=None,
model=None,
lane="iaai_cars",
limit=None,
only_new=False,
)
sync_segmented_mock.assert_not_called()
release_lock.assert_called_once()
def test_sync_listing_hourly_path_ignores_checkpoint_after_full_scan_completed(self) -> None:
segments = [{"make": "ACURA"}, {"make": "AUDI"}]
settings_stub = make_settings_stub()
with patch.object(tasks, "_get_persistence") as get_persistence, \
patch.object(tasks, "_get_redis") as get_redis, \
patch.object(tasks, "Settings", return_value=settings_stub), \
patch.object(tasks, "_acquire_lock", return_value=True), \
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, "_clear_sync_checkpoint") as clear_checkpoint, \
patch.object(tasks, "_hourly_sitemap_rolling_refresh_sync", return_value={
"status": "success",
"run_id": None,
"cars_upserted": 1,
"cars_failed": 0,
"images_upserted": 0,
"skipped_existing": 0,
"elapsed_seconds": None,
"failures": [],
"hourly_mode": "sitemap_rolling_refresh",
"discovered_urls": 10,
"new_urls": 1,
"refresh_urls": 1,
"sold_marked": 0,
}) as hourly_sync, \
patch.object(tasks.sync_listing_task, "update_state"), \
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=segments):
get_persistence.return_value = MagicMock()
redis_client = MagicMock()
# Оставшийся чекпоинт не должен использоваться.
redis_client.get.side_effect = lambda key: (
"0" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
)
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
tasks.sync_listing_task.push_request(id="task-792")
try:
result = tasks.sync_listing_task.run()
finally:
tasks.sync_listing_task.pop_request()
self.assertEqual(result["status"], "success")
self.assertEqual(result["hourly_mode"], "sitemap_rolling_refresh")
hourly_sync.assert_called_once()
clear_checkpoint.assert_called()
release_lock.assert_called_once()
def test_sync_listing_bootstrap_ignores_beyond_segment_checkpoint(self) -> None:
segments = [{"make": "ACURA"}, {"make": "AUDI"}]
settings_stub = make_settings_stub()
with patch.object(tasks, "_get_persistence") as get_persistence, \
patch.object(tasks, "_get_redis") as get_redis, \
patch.object(tasks, "Settings", return_value=settings_stub), \
patch.object(tasks, "_acquire_lock", return_value=True), \
patch.object(tasks, "_is_full_scan_done", return_value=False), \
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
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=segments):
get_persistence.return_value = MagicMock()
redis_client = MagicMock()
redis_client.get.side_effect = lambda key: (
"99" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
)
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
sync_segmented_mock = MagicMock(side_effect=AssertionError("segmented path should not be used"))
sync_listing_mock = MagicMock(return_value={
"run_id": 15, "status": "success", "full_scan_completed": True,
"cars_upserted": 0, "cars_failed": 0, "images_upserted": 0,
"skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [],
})
scraper_ctx = MagicMock()
scraper_ctx.__enter__.return_value.sync_listing_segmented = sync_segmented_mock
scraper_ctx.__enter__.return_value.sync_listing = sync_listing_mock
scraper_ctx.__exit__.return_value = None
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
tasks.sync_listing_task.push_request(id="task-beyond")
try:
result = tasks.sync_listing_task.run()
finally:
tasks.sync_listing_task.pop_request()
self.assertEqual(result["status"], "success")
sync_listing_mock.assert_called_once_with(
make=None,
model=None,
lane="iaai_cars",
limit=None,
only_new=False,
)
sync_segmented_mock.assert_not_called()
release_lock.assert_called_once()
def test_sync_listing_task_clears_checkpoint_on_hourly_run(self) -> None:
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), \
@@ -169,40 +540,121 @@ 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()), \
patch.object(tasks.sync_listing_task, "update_state"), \
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=[]):
patch.object(tasks, "_run_browser_job", return_value={
"run_id": 13,
"status": "partial_success",
"full_scan_completed": True,
"cars_upserted": 2,
"cars_failed": 1,
"images_upserted": 3,
"skipped_existing": 0,
"elapsed_seconds": 4.0,
"failures": [{"vehicle_url": "v", "error": "e"}],
}), \
patch.object(tasks.sync_listing_task, "update_state"):
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_client2
stop_event = MagicMock()
heartbeat_thread = MagicMock()
start_heartbeat.return_value = (stop_event, heartbeat_thread)
redis_client = MagicMock()
redis_client.get.return_value = None
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
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": [],
})
scraper_ctx = MagicMock()
scraper_ctx.__enter__.return_value.sync_listing = sync_listing_mock
scraper_ctx.__exit__.return_value = None
tasks.sync_listing_task.push_request(id="task-791")
try:
result = tasks.sync_listing_task.run(make="Toyota")
finally:
tasks.sync_listing_task.pop_request()
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
tasks.sync_listing_task.push_request(id="task-789")
self.assertEqual(result["status"], "partial_success")
# Hourly-ветка (full_scan_done=True) всегда удаляет оставшийся чекпоинт
# до браузерного job + после. Главное — вызов произошёл.
clear_checkpoint.assert_called()
release_lock.assert_called_once()
def test_try_set_followup_pending_deduplicates(self) -> None:
redis_client = MagicMock()
redis_client.set.side_effect = [True, False]
self.assertTrue(tasks._try_set_followup_pending(redis_client, ttl_seconds=120))
self.assertFalse(tasks._try_set_followup_pending(redis_client, ttl_seconds=120))
def test_bump_bootstrap_failure_streak_opens_circuit_breaker(self) -> None:
redis_client = MagicMock()
redis_client.incr.return_value = tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT
streak, should_enqueue = tasks._bump_bootstrap_failure_streak(
redis_client,
reason="max_retries_exceeded",
)
self.assertEqual(streak, tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT)
self.assertFalse(should_enqueue)
redis_client.expire.assert_called_once_with(
tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY,
tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS,
)
def test_bump_bootstrap_failure_streak_allows_retry_before_limit(self) -> None:
redis_client = MagicMock()
redis_client.incr.return_value = 1
streak, should_enqueue = tasks._bump_bootstrap_failure_streak(
redis_client,
reason="bootstrap_not_completed",
)
self.assertEqual(streak, 1)
self.assertTrue(should_enqueue)
def test_sync_listing_task_does_not_enqueue_followup_after_bootstrap_error_limit(self) -> None:
settings_stub = make_settings_stub(parallel_segments=False)
with patch.object(tasks, "_get_persistence") as get_persistence, \
patch.object(tasks, "_get_redis") as get_redis, \
patch.object(tasks, "Settings", return_value=settings_stub), \
patch.object(tasks, "_acquire_lock", return_value=True), \
patch.object(tasks, "_is_full_scan_done", return_value=False), \
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
patch.object(tasks, "_try_set_followup_pending", return_value=True), \
patch.object(
tasks,
"_bump_bootstrap_failure_streak",
return_value=(tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT, False),
) as bump_streak, \
patch.object(tasks, "_clear_followup_pending") as clear_pending, \
patch.object(tasks.sync_listing_task, "update_state"), \
patch.object(
tasks,
"_run_browser_job",
return_value={
"run_id": 99,
"status": "failed",
"full_scan_completed": False,
"cars_upserted": 0,
"cars_failed": 0,
"images_upserted": 0,
"skipped_existing": 0,
"elapsed_seconds": 1.0,
"failures": [{"vehicle_url": "listing", "error": "bad resume"}],
"listing": {"vehicles_collected": 0},
},
):
get_persistence.return_value = MagicMock()
redis_client = MagicMock()
redis_client.get.return_value = None
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
with patch.object(tasks.sync_listing_task, "app", new=MagicMock()) as task_app:
tasks.sync_listing_task.push_request(id="task-900")
try:
result = tasks.sync_listing_task.run()
finally:
tasks.sync_listing_task.pop_request()
self.assertEqual(result["status"], "success")
self.assertEqual(sync_listing_mock.call_args.kwargs["start_page"], 10)
clear_checkpoint.assert_called()
release_lock.assert_called_once()
self.assertEqual(result["status"], "failed")
bump_streak.assert_called_once()
clear_pending.assert_called()
task_app.send_task.assert_not_called()
if __name__ == "__main__":