1 Commits

23 changed files with 3436 additions and 518 deletions

14
.gitignore vendored
View File

@@ -21,10 +21,10 @@ celerybeat-schedule*
tokens_data/
tokens.json
# Ad-hoc debug scripts and one-off downloads
_*.sh
_*.tar.gz
burst_*.ps1
burst_*.txt
sitemap*.xml
vd_*.html
# Local deployment/debug artifacts
/*.tgz
/*.tar.gz
/*.zip
/NOW()
/_speed_check.sql
/.env.bak*

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,81 @@
# Добавление таблиц для новой архитектуры 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, index=True),
sa.Column("discovered_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now(), index=True),
sa.Column("status", sa.String(20), nullable=False, server_default="pending", index=True),
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"),
)
# Индексы для vehicle_candidates
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, index=True),
sa.Column("captured_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now(), index=True),
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),
)
# Индексы для vehicle_raw_snapshots
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, index=True),
sa.Column("parsed_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now(), index=True),
sa.Column("success", sa.Boolean(), nullable=False),
sa.Column("parsed_data", sa.Text(), nullable=True),
sa.Column("error_message", sa.Text(), nullable=True),
)
# Индексы для vehicle_parse_results
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, index=True),
sa.Column("reason", sa.String(50), nullable=False),
sa.Column("retry_at", sa.DateTime(timezone=True), nullable=False, index=True),
sa.Column("attempts", sa.Integer(), nullable=False, server_default="0"),
sa.Column("max_attempts", sa.Integer(), nullable=False, server_default="3"),
)
# Индексы для vehicle_retry_queue
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

@@ -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 {

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, ...] = (
@@ -217,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)
@@ -224,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)
@@ -278,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

@@ -17,7 +17,8 @@ def first_non_empty(values: Iterable[Any]) -> Any | None:
return None
# Регулярные выражения для lot и price.
# Регулярные выражения для VIN, lot и price.
VIN_RE = re.compile(r"\b([A-HJ-NPR-Z0-9]{17})\b", re.IGNORECASE)
LOT_RE = re.compile(r"\b(\d{7,10})\b")
PRICE_RE = re.compile(r"\$\s?([\d,]+(?:\.\d{1,2})?)")

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

File diff suppressed because it is too large Load Diff

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

@@ -1,7 +1,10 @@
# Задачи Celery для синхронизации автомобилей и листинга IAAI.
from concurrent.futures import ThreadPoolExecutor
import json
import logging
import os
import signal
from threading import Event, Thread
import time
import uuid
@@ -13,6 +16,7 @@ from redis import Redis
from ..core.config import Settings, parse_listing_segments
from ..scraper import IAAIScraper
from ..storage.db import PersistenceService
from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitemap_with_stats
logger = logging.getLogger("iaai_scraper.worker.tasks")
@@ -25,6 +29,13 @@ SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT = 3
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS = 24 * 60 * 60
SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task"
SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}"
SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress"
SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total"
SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}"
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
@@ -91,6 +102,283 @@ def _sync_listing_lock_ttl_seconds() -> int:
return max(effective_hard + 120, 300)
def _sync_segment_lock_ttl_seconds() -> int:
return 300
def _task_progress_key(task_id: str) -> str:
return TASK_PROGRESS_KEY_FMT.format(task_id=task_id)
def _update_task_progress(
redis_client: Redis,
*,
task_id: str,
stage: str,
ttl_seconds: int,
**payload,
) -> None:
try:
data = {
"task_id": task_id,
"stage": stage,
"ts": int(time.time()),
**payload,
}
redis_client.set(
_task_progress_key(task_id),
json.dumps(data, ensure_ascii=False),
ex=max(60, int(ttl_seconds)),
)
except Exception:
logger.warning("Failed to update task progress for %s", task_id, exc_info=True)
def _clear_task_progress(redis_client: Redis, task_id: str) -> None:
try:
redis_client.delete(_task_progress_key(task_id))
except Exception:
logger.warning("Failed to clear task progress for %s", task_id, exc_info=True)
def _hourly_sitemap_diff_sync(*, lane: str, limit: int | None, only_new: bool | None) -> dict[str, object]:
del limit, only_new
settings = Settings()
persistence = _get_persistence()
persistence.create_tables()
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
discovered_urls = discovery_result.vehicle_urls
active_urls = set(discovered_urls)
existing_urls = persistence.get_all_active_origin_urls_for_lane("iaai:")
new_urls = [url for url in discovered_urls if url not in existing_urls]
sold_count = 0
if active_urls:
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
cars_upserted = 0
cars_failed = 0
images_upserted = 0
failures: list[dict[str, str]] = []
if new_urls:
with IAAIScraper(settings) as scraper:
batch_size = settings.celery.batch_size
for batch_start in range(0, len(new_urls), batch_size):
batch_urls = new_urls[batch_start:batch_start + batch_size]
batch_result = scraper.sync_batch(batch_urls, lane=lane)
cars_upserted += int(batch_result.get("cars_upserted", 0))
cars_failed += int(batch_result.get("cars_failed", 0))
images_upserted += int(batch_result.get("images_upserted", 0))
failures.extend(batch_result.get("failures", []))
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
return {
"status": status,
"run_id": None,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"skipped_existing": len(discovered_urls) - len(new_urls),
"elapsed_seconds": None,
"failures": failures,
"full_scan_completed": True,
"hourly_mode": "sitemap_diff",
"discovered_urls": len(discovered_urls),
"new_urls": len(new_urls),
"sold_marked": sold_count,
"transport": discovery_result.stats.transport,
}
def _hourly_sitemap_full_refresh_sync(*, lane: str, limit: int | None, only_new: bool | None) -> dict[str, object]:
del limit, only_new
settings = Settings()
persistence = _get_persistence()
persistence.create_tables()
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
discovered_urls = discovery_result.vehicle_urls
active_urls = set(discovered_urls)
sold_count = 0
if active_urls:
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
cars_upserted = 0
cars_failed = 0
images_upserted = 0
failures: list[dict[str, str]] = []
with IAAIScraper(settings) as scraper:
batch_size = settings.celery.batch_size
for batch_start in range(0, len(discovered_urls), batch_size):
batch_urls = discovered_urls[batch_start:batch_start + batch_size]
batch_result = scraper.sync_batch(batch_urls, lane=lane)
cars_upserted += int(batch_result.get("cars_upserted", 0))
cars_failed += int(batch_result.get("cars_failed", 0))
images_upserted += int(batch_result.get("images_upserted", 0))
failures.extend(batch_result.get("failures", []))
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
return {
"status": status,
"run_id": None,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"skipped_existing": 0,
"elapsed_seconds": None,
"failures": failures,
"full_scan_completed": True,
"hourly_mode": "sitemap_full_refresh",
"discovered_urls": len(discovered_urls),
"new_urls": None,
"sold_marked": sold_count,
"transport": discovery_result.stats.transport,
}
def _hourly_sitemap_rolling_refresh_sync(
*,
redis_client: Redis,
lane: str,
limit: int | None,
only_new: bool | None,
) -> dict[str, object]:
del limit, only_new
settings = Settings()
persistence = _get_persistence()
persistence.create_tables()
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
discovered_urls = discovery_result.vehicle_urls
active_urls = set(discovered_urls)
sold_count = 0
if active_urls:
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
existing_urls = persistence.get_all_active_origin_urls_for_lane("iaai:")
new_urls = [url for url in discovered_urls if url not in existing_urls]
total_active = persistence.count_active_cars_for_lane("iaai:")
batch_size = max(1, int(settings.discovery.hourly_refresh_batch_size))
try:
offset = int(redis_client.get(SITEMAP_HOURLY_REFRESH_OFFSET_KEY) or 0)
except Exception:
offset = 0
refresh_urls = persistence.get_active_origin_urls_batch_for_refresh(
prefix="iaai:",
offset=offset,
limit=batch_size,
)
if not refresh_urls and total_active > 0:
offset = 0
refresh_urls = persistence.get_active_origin_urls_batch_for_refresh(
prefix="iaai:",
offset=0,
limit=batch_size,
)
next_offset = 0
if total_active > 0:
next_offset = offset + len(refresh_urls)
if next_offset >= total_active:
next_offset = 0
try:
redis_client.set(SITEMAP_HOURLY_REFRESH_OFFSET_KEY, str(next_offset))
except Exception:
logger.warning("Failed to persist hourly rolling refresh offset", exc_info=True)
seen: set[str] = set()
target_urls: list[str] = []
for url in new_urls + refresh_urls:
if url in seen:
continue
seen.add(url)
target_urls.append(url)
cars_upserted = 0
cars_failed = 0
images_upserted = 0
failures: list[dict[str, str]] = []
if target_urls:
with IAAIScraper(settings) as scraper:
worker_batch_size = settings.celery.batch_size
for batch_start in range(0, len(target_urls), worker_batch_size):
batch_urls = target_urls[batch_start:batch_start + worker_batch_size]
batch_result = scraper.sync_batch(batch_urls, lane=lane)
cars_upserted += int(batch_result.get("cars_upserted", 0))
cars_failed += int(batch_result.get("cars_failed", 0))
images_upserted += int(batch_result.get("images_upserted", 0))
failures.extend(batch_result.get("failures", []))
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
return {
"status": status,
"run_id": None,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"skipped_existing": max(0, len(discovered_urls) - len(new_urls)),
"elapsed_seconds": None,
"failures": failures,
"full_scan_completed": True,
"hourly_mode": "sitemap_rolling_refresh",
"discovered_urls": len(discovered_urls),
"new_urls": len(new_urls),
"refresh_urls": len(refresh_urls),
"sold_marked": sold_count,
"transport": discovery_result.stats.transport,
"refresh_offset": offset,
"refresh_next_offset": next_offset,
"active_total": total_active,
}
def _start_stall_watchdog(
redis_client: Redis,
*,
task_id: str,
stall_timeout_seconds: int,
) -> tuple[Event, Thread]:
stop_event = Event()
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
def _watchdog() -> None:
key = _task_progress_key(task_id)
while not stop_event.wait(interval_seconds):
try:
raw = redis_client.get(key)
if not raw:
continue
data = json.loads(raw)
last_ts = int(data.get("ts") or 0)
if not last_ts:
continue
age = int(time.time()) - last_ts
if age < stall_timeout_seconds:
continue
logger.error(
"Task %s stalled for %ss at stage=%s payload=%s; killing worker process for redelivery",
task_id,
age,
data.get("stage"),
data,
)
except Exception:
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
continue
os.kill(os.getpid(), signal.SIGKILL)
thread = Thread(target=_watchdog, name=f"task-stall-watchdog-{task_id[:8]}", daemon=True)
thread.start()
return stop_event, thread
def _get_persistence() -> PersistenceService:
settings = Settings()
persistence = PersistenceService(settings)
@@ -360,6 +648,204 @@ def _start_lock_heartbeat(
return stop_event, thread
def _reset_segments_progress(redis_client: Redis, total: int) -> None:
try:
pipe = redis_client.pipeline()
pipe.delete(SYNC_SEGMENTS_PROGRESS_KEY)
pipe.set(SYNC_SEGMENTS_TOTAL_KEY, str(int(total)), ex=SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
pipe.execute()
except Exception:
logger.warning("Failed to reset segments progress", exc_info=True)
def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[int, int]:
"""Помечает сегмент завершённым. Возвращает (completed_count, total)."""
try:
pipe = redis_client.pipeline()
pipe.sadd(SYNC_SEGMENTS_PROGRESS_KEY, str(int(segment_index)))
pipe.expire(SYNC_SEGMENTS_PROGRESS_KEY, SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
pipe.scard(SYNC_SEGMENTS_PROGRESS_KEY)
pipe.get(SYNC_SEGMENTS_TOTAL_KEY)
results = pipe.execute()
completed = int(results[2] or 0)
total = int(results[3] or 0) if results[3] else 0
return completed, total
except Exception:
logger.warning("Failed to mark segment %d completed", segment_index, exc_info=True)
return 0, 0
@shared_task(
name="iaai_scraper.worker.tasks.sync_segment_task",
bind=True,
max_retries=2,
default_retry_delay=60,
acks_late=True,
)
def sync_segment_task(
self,
segment_index: int,
segment: dict,
lane: str = "iaai_cars",
only_new: bool | None = None,
is_bootstrap: bool = False,
):
"""Обработка одного сегмента листинга. Запускается параллельно несколькими воркерами."""
persistence = _get_persistence()
persistence.create_tables()
task_id = self.request.id or "unknown"
redis_client = _get_redis()
settings = Settings()
# Per-segment lock — защита от случайного дубля
seg_lock_key = SYNC_SEGMENT_LOCK_KEY_FMT.format(idx=int(segment_index))
owner_token = f"{task_id}:{uuid.uuid4().hex}"
lock_ttl = _sync_segment_lock_ttl_seconds()
lock_acquired = _acquire_lock(redis_client, seg_lock_key, owner_token, lock_ttl)
if not lock_acquired:
logger.info("sync_segment_task[%d] skipped: already running", segment_index)
return {"status": "skipped", "segment_index": segment_index, "reason": "duplicate"}
seg_make = segment.get("make")
seg_year_min = segment.get("year_min")
seg_year_max = segment.get("year_max")
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})"
heartbeat_stop: Event | None = None
heartbeat_thread: Thread | None = None
watchdog_stop: Event | None = None
watchdog_thread: Thread | None = None
stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds))
progress_ttl = max(lock_ttl + 120, stall_timeout + 120)
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_started",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
)
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(redis_client, seg_lock_key, owner_token, lock_ttl)
watchdog_stop, watchdog_thread = _start_stall_watchdog(
redis_client,
task_id=task_id,
stall_timeout_seconds=stall_timeout,
)
try:
def _job():
with IAAIScraper() as scraper:
scraper.set_progress_callback(
lambda stage, meta: _update_task_progress(
redis_client,
task_id=task_id,
stage=stage,
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
**meta,
)
)
base_url = scraper.settings.listing.cars_url
seg_url = scraper._build_segment_listing_url(base_url, seg_make) if seg_make else None
return scraper.sync_listing(
make=None if seg_url else seg_make,
model=None,
lane=lane,
only_new=only_new,
listing_url=seg_url,
year_min=seg_year_min,
year_max=seg_year_max,
skip_mark_sold=True,
)
result = _run_browser_job(_job)
segment_done = bool(result.get("full_scan_completed", False))
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_completed",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
status=result.get("status", "success"),
cars_upserted=result.get("cars_upserted", 0),
cars_failed=result.get("cars_failed", 0),
)
if is_bootstrap and segment_done:
completed, total = _mark_segment_completed(redis_client, segment_index)
logger.warning(
"Segment %d (%s) bootstrap done: %d/%d completed",
segment_index, seg_label, completed, total,
)
if total > 0 and completed >= total:
_set_full_scan_done(redis_client, True)
_clear_sync_checkpoint(redis_client)
_clear_bootstrap_failure_streak(redis_client)
logger.warning("All %d segments completed; bootstrap full scan done", total)
return {
"status": result.get("status", "success"),
"segment_index": segment_index,
"segment_label": seg_label,
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"vehicles_collected": result.get("listing", {}).get("vehicles_collected", 0),
"full_scan_completed": segment_done,
}
except SoftTimeLimitExceeded:
logger.warning("sync_segment_task[%d] soft timeout — partial progress saved", segment_index)
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_soft_timeout",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
)
return {
"status": "timed_out",
"segment_index": segment_index,
"segment_label": seg_label,
}
except Exception as exc:
logger.error("sync_segment_task[%d] failed: %s%s", segment_index, seg_label, exc, exc_info=True)
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_failed",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
error=str(exc),
)
try:
raise self.retry(exc=exc)
except self.MaxRetriesExceededError:
return {
"status": "failed",
"segment_index": segment_index,
"segment_label": seg_label,
"error": str(exc),
}
finally:
if watchdog_stop is not None:
watchdog_stop.set()
if watchdog_thread is not None:
watchdog_thread.join(timeout=1)
if heartbeat_stop is not None:
heartbeat_stop.set()
if heartbeat_thread is not None:
heartbeat_thread.join(timeout=1)
_clear_task_progress(redis_client, task_id)
_release_lock_if_owner(redis_client, seg_lock_key, owner_token)
@shared_task(
name="iaai_scraper.worker.tasks.sync_vehicle_task",
bind=True,
@@ -477,10 +963,20 @@ def sync_listing_task(
}
try:
settings = Settings()
full_scan_done_before_run = _is_full_scan_done(redis_client)
force_bootstrap_full_scan = not full_scan_done_before_run
effective_limit = None if force_bootstrap_full_scan else limit
effective_only_new = False if force_bootstrap_full_scan else only_new
hourly_mode = settings.discovery.hourly_mode.strip().lower()
discovery_mode = settings.discovery.mode.strip().lower()
prefer_sitemap_mainline = (
make is None
and model is None
and effective_limit is None
and not effective_only_new
)
use_hourly_sitemap_sync = full_scan_done_before_run and prefer_sitemap_mainline
# Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
# Используется только во время bootstrap для пропуска уже обработанных сегментов.
@@ -503,6 +999,58 @@ def sync_listing_task(
effective_limit,
)
if use_hourly_sitemap_sync:
if hourly_mode == "diff":
logger.info("Hourly mode: running sitemap diff sync instead of full listing traversal")
result = _hourly_sitemap_diff_sync(
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
)
elif hourly_mode == "full_refresh":
logger.info("Hourly mode: running sitemap full refresh of all active vehicles")
result = _hourly_sitemap_full_refresh_sync(
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
)
else:
logger.info("Hourly mode: running sitemap rolling refresh of active vehicles")
result = _hourly_sitemap_rolling_refresh_sync(
redis_client=redis_client,
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
)
try:
redis_client.set(SITEMAP_HOURLY_LAST_COUNT_KEY, str(result.get("discovered_urls", 0)))
except Exception:
logger.warning("Failed to persist hourly sitemap count", exc_info=True)
summary = {
"task_id": task_id,
"run_id": result.get("run_id"),
"status": result.get("status", "success"),
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"images_upserted": result.get("images_upserted", 0),
"skipped_existing": result.get("skipped_existing", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
"failures_count": len(result.get("failures") or []),
"hourly_mode": result.get("hourly_mode"),
"discovered_urls": result.get("discovered_urls", 0),
"new_urls": result.get("new_urls", 0),
"sold_marked": result.get("sold_marked", 0),
}
logger.info(
"sync_listing_task hourly diff completed: status=%s, new=%d, sold=%d, skipped=%d",
summary["status"],
summary["new_urls"],
summary["sold_marked"],
summary["skipped_existing"],
)
return summary
_clear_followup_pending(redis_client)
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
@@ -514,9 +1062,91 @@ def sync_listing_task(
self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id})
# Определяем сегменты из конфига.
settings = Settings()
segments = parse_listing_segments(settings.listing.listing_segments_json)
use_segmented = bool(segments) and make is None and model is None
use_segmented = (
bool(segments)
and make is None
and model is None
and not prefer_sitemap_mainline
and discovery_mode != "sitemap"
)
if prefer_sitemap_mainline:
if discovery_mode != "sitemap":
logger.warning(
"Unfiltered full scan forcing sitemap discovery despite IAAI_DISCOVERY_MODE=%s",
discovery_mode or "unset",
)
if segments:
logger.info("Ignoring configured listing segments for unfiltered sitemap full scan")
# --- Параллельный диспатч сегментов: dispatch & exit ---
if use_segmented and settings.celery.parallel_segments:
# Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET.
already_completed: set[int] = set()
if force_bootstrap_full_scan:
try:
raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set()
already_completed = {int(x) for x in raw if str(x).strip().lstrip("-").isdigit()}
except Exception:
already_completed = set()
pending = [
(idx, seg) for idx, seg in enumerate(segments)
if idx not in already_completed
]
if not pending:
# Всё уже сделано — фиксируем bootstrap done.
if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, True)
_clear_sync_checkpoint(redis_client)
_clear_bootstrap_failure_streak(redis_client)
logger.warning("Parallel segments: nothing to dispatch (all completed)")
return {
"status": "success",
"task_id": task_id,
"mode": "parallel_segments",
"segments_total": len(segments),
"segments_dispatched": 0,
"segments_already_completed": len(already_completed),
}
# При первом запуске bootstrap фиксируем total, чтобы знать когда остановиться.
if force_bootstrap_full_scan and not already_completed:
_reset_segments_progress(redis_client, len(segments))
dispatched = 0
for idx, seg in pending:
try:
self.app.send_task(
"iaai_scraper.worker.tasks.sync_segment_task",
kwargs={
"segment_index": idx,
"segment": seg,
"lane": lane,
"only_new": effective_only_new,
"is_bootstrap": force_bootstrap_full_scan,
},
queue="scraping",
)
dispatched += 1
except Exception:
logger.warning("Failed to dispatch segment %d", idx, exc_info=True)
logger.warning(
"Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s)",
dispatched, len(segments), len(already_completed), force_bootstrap_full_scan,
)
return {
"status": "success",
"task_id": task_id,
"mode": "parallel_segments",
"segments_total": len(segments),
"segments_dispatched": dispatched,
"segments_already_completed": len(already_completed),
}
# --- конец параллельной ветки ---
resume_from_segment = 0
if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None:
@@ -608,10 +1238,7 @@ def sync_listing_task(
)
if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, False)
# SoftTimeLimit считаем failure для circuit breaker:
# если сегмент систематически не укладывается в time limit,
# нельзя крутить followup бесконечно.
_enqueue_bootstrap_followup("soft_time_limit_exceeded", count_as_failure=True)
_enqueue_bootstrap_followup("soft_time_limit_exceeded", count_as_failure=False)
# Partial progress уже записан в БД через finish_sync_run.
# Не retry — следующий запуск продолжит обработку по расписанию.
return {
@@ -646,3 +1273,89 @@ def sync_listing_task(
heartbeat_thread.join(timeout=max(1.0, min(5.0, lock_ttl / 10)))
if lock_acquired:
_release_lock_if_owner(redis_client, SYNC_LISTING_LOCK_KEY, owner_token)
# Новые задачи для ingestion pipeline
@shared_task(
name="iaai_scraper.worker.tasks.discover_vehicles_task",
bind=True,
max_retries=2,
default_retry_delay=60,
acks_late=True,
)
def discover_vehicles_task(self, max_urls: int | None = None):
"""Задача для обнаружения новых URL автомобилей."""
from ..discovery_service import DiscoveryService
try:
discovery = DiscoveryService()
count = discovery.discover_new_vehicles(max_urls=max_urls)
logger.info("discover_vehicles_task completed: %d candidates added", count)
return {"status": "success", "candidates_added": count}
except Exception as exc:
logger.error("discover_vehicles_task failed: %s", exc, exc_info=True)
raise self.retry(exc=exc)
@shared_task(
name="iaai_scraper.worker.tasks.fetch_pending_candidates_task",
bind=True,
max_retries=2,
default_retry_delay=30,
acks_late=True,
)
def fetch_pending_candidates_task(self, limit: int = 10):
"""Задача для захвата данных ожидающих кандидатов."""
from ..fetch_service import FetchService
try:
fetch = FetchService()
count = fetch.process_pending_candidates(limit=limit)
logger.info("fetch_pending_candidates_task completed: %d candidates processed", count)
return {"status": "success", "candidates_processed": count}
except Exception as exc:
logger.error("fetch_pending_candidates_task failed: %s", exc, exc_info=True)
raise self.retry(exc=exc)
@shared_task(
name="iaai_scraper.worker.tasks.enrich_snapshots_task",
bind=True,
max_retries=2,
default_retry_delay=30,
acks_late=True,
)
def enrich_snapshots_task(self, limit: int = 10):
"""Задача для парсинга и обогащения snapshots."""
from ..enrichment_service import EnrichmentService
try:
enrichment = EnrichmentService()
count = enrichment.process_unparsed_snapshots(limit=limit)
logger.info("enrich_snapshots_task completed: %d snapshots enriched", count)
return {"status": "success", "snapshots_enriched": count}
except Exception as exc:
logger.error("enrich_snapshots_task failed: %s", exc, exc_info=True)
raise self.retry(exc=exc)
@shared_task(
name="iaai_scraper.worker.tasks.run_ingestion_pipeline_task",
bind=True,
max_retries=2,
default_retry_delay=120,
acks_late=True,
)
def run_ingestion_pipeline_task(self):
"""Задача для запуска полного ingestion pipeline."""
from ..scheduler_service import SchedulerService
try:
scheduler = SchedulerService()
scheduler.run_full_pipeline()
logger.info("run_ingestion_pipeline_task completed")
return {"status": "success"}
except Exception as exc:
logger.error("run_ingestion_pipeline_task failed: %s", exc, exc_info=True)
raise self.retry(exc=exc)

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",

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:
@@ -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()
@@ -268,6 +315,126 @@ class TestScraperSync(unittest.TestCase):
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

@@ -6,7 +6,210 @@ 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()
@@ -175,15 +378,17 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
self.assertIsNone(tasks._load_last_completed_segment(redis_client))
redis_client.delete.assert_not_called()
def test_sync_listing_resumes_from_next_segment_during_bootstrap(self) -> None:
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, \
@@ -199,13 +404,15 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
sync_segmented_mock = MagicMock(return_value={
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):
@@ -216,21 +423,42 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
tasks.sync_listing_task.pop_request()
self.assertEqual(result["status"], "success")
# last_completed=1 → start_segment=2 (AUDI завершён, возобновляем с BMW).
self.assertEqual(sync_segmented_mock.call_args.kwargs["start_segment"], 2)
self.assertEqual(sync_segmented_mock.call_args.kwargs["start_page"], 1)
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_ignores_checkpoint_after_full_scan_completed(self) -> None:
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, "_run_browser_job", side_effect=lambda fn: fn()), \
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()
@@ -242,32 +470,24 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
sync_segmented_mock = MagicMock(return_value={
"run_id": 14, "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.__exit__.return_value = None
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
tasks.sync_listing_task.push_request(id="task-792")
try:
result = tasks.sync_listing_task.run()
finally:
tasks.sync_listing_task.pop_request()
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(sync_segmented_mock.call_args.kwargs["start_segment"], 0)
self.assertIsNone(sync_segmented_mock.call_args.kwargs["progress_callback"])
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_checkpoint_beyond_segments_restarts_from_zero(self) -> None:
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, \
@@ -283,13 +503,15 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
get_redis.return_value = redis_client
start_heartbeat.return_value = (MagicMock(), MagicMock())
sync_segmented_mock = MagicMock(return_value={
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):
@@ -300,7 +522,14 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
tasks.sync_listing_task.pop_request()
self.assertEqual(result["status"], "success")
self.assertEqual(sync_segmented_mock.call_args.kwargs["start_segment"], 0)
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:
@@ -377,28 +606,38 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
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, "_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},
}):
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