6 Commits

Author SHA1 Message Date
qananasikq
db2fa5ec17 add antibot unit tests 2026-04-10 15:18:57 +03:00
qananasikq
18dab6d48d docker non-root and healthcheck fix 2026-04-10 15:18:57 +03:00
qananasikq
3870cf4d94 add task status endpoint 2026-04-10 15:18:57 +03:00
qananasikq
62aafae793 add Redis lock for sync_listing 2026-04-10 15:18:57 +03:00
qananasikq
acf30fcc20 add antibot guard and retry 2026-04-10 15:18:57 +03:00
qananasikq
7b1501008e fix limit after only_new filter 2026-04-10 15:18:57 +03:00
12 changed files with 255 additions and 24 deletions

View File

@@ -3,12 +3,17 @@ FROM mcr.microsoft.com/playwright/python:v1.58.0-noble
ENV PYTHONDONTWRITEBYTECODE=1 \ ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1 PYTHONUNBUFFERED=1
RUN groupadd --system app && useradd --system --gid app --create-home app
WORKDIR /app WORKDIR /app
# Xvfb for headed Chromium in container (IAAI blocks headless) # Xvfb for headed Chromium in container (IAAI blocks headless)
RUN apt-get update -qq && apt-get install -y --no-install-recommends xvfb \ RUN apt-get update -qq && apt-get install -y --no-install-recommends xvfb \
&& rm -rf /var/lib/apt/lists/* && rm -rf /var/lib/apt/lists/*
# Required for non-root Xvfb runtime in containers
RUN mkdir -p /tmp/.X11-unix && chmod 1777 /tmp/.X11-unix
COPY pyproject.toml uv.lock ./ COPY pyproject.toml uv.lock ./
RUN pip install --no-cache-dir uv \ RUN pip install --no-cache-dir uv \
&& uv export --format requirements-txt --no-dev --no-hashes --no-emit-project --frozen -o /tmp/requirements.txt \ && uv export --format requirements-txt --no-dev --no-hashes --no-emit-project --frozen -o /tmp/requirements.txt \
@@ -17,9 +22,12 @@ RUN pip install --no-cache-dir uv \
COPY . . COPY . .
RUN python -m compileall -q iaai_scraper RUN python -m compileall -q iaai_scraper
RUN chmod +x entrypoint.sh RUN chmod +x entrypoint.sh
RUN chown -R app:app /app
STOPSIGNAL SIGINT STOPSIGNAL SIGINT
USER app
# По умолчанию API, но worker/beat переопределяют CMD в docker-compose # По умолчанию API, но worker/beat переопределяют CMD в docker-compose
ENTRYPOINT ["./entrypoint.sh"] ENTRYPOINT ["./entrypoint.sh"]
CMD ["uvicorn", "iaai_scraper.api.app:app", "--host", "0.0.0.0", "--port", "8000"] CMD ["uvicorn", "iaai_scraper.api.app:app", "--host", "0.0.0.0", "--port", "8000"]

View File

@@ -121,6 +121,7 @@ services:
<<: *worker-service <<: *worker-service
container_name: iaai-worker container_name: iaai-worker
restart: unless-stopped restart: unless-stopped
init: true
depends_on: depends_on:
migrate: migrate:
condition: service_completed_successfully condition: service_completed_successfully
@@ -130,13 +131,14 @@ services:
command: > command: >
celery -A iaai_scraper.worker.celery_app worker celery -A iaai_scraper.worker.celery_app worker
--loglevel=info --concurrency=1 --pool=prefork --loglevel=info --concurrency=1 --pool=prefork
--pidfile=/tmp/celery-worker.pid
-Q scraping --max-tasks-per-child=20 -Q scraping --max-tasks-per-child=20
healthcheck: healthcheck:
test: ["CMD", "celery", "-A", "iaai_scraper.worker.celery_app", "inspect", "ping", "-d", "celery@$$HOSTNAME"] test: ["CMD-SHELL", "test -f /tmp/celery-worker.pid && kill -0 $(cat /tmp/celery-worker.pid)"]
interval: 60s interval: 60s
timeout: 20s timeout: 10s
retries: 3 retries: 3
start_period: 40s start_period: 90s
logging: logging:
driver: json-file driver: json-file
options: options:
@@ -148,6 +150,7 @@ services:
<<: *worker-service <<: *worker-service
container_name: iaai-beat container_name: iaai-beat
restart: unless-stopped restart: unless-stopped
init: true
depends_on: depends_on:
migrate: migrate:
condition: service_completed_successfully condition: service_completed_successfully
@@ -156,12 +159,14 @@ services:
command: > command: >
celery -A iaai_scraper.worker.celery_app beat celery -A iaai_scraper.worker.celery_app beat
--loglevel=info --loglevel=info
--pidfile=/tmp/celery-beat.pid
--schedule=/tmp/celerybeat-schedule
healthcheck: healthcheck:
test: ["CMD", "python", "-c", "import pathlib,sys; p=pathlib.Path('/tmp/celerybeat-schedule'); sys.exit(0 if p.exists() else 1)"] test: ["CMD-SHELL", "test -f /tmp/celery-beat.pid && kill -0 $(cat /tmp/celery-beat.pid)"]
interval: 60s interval: 60s
timeout: 20s timeout: 10s
retries: 3 retries: 3
start_period: 60s start_period: 90s
logging: logging:
driver: json-file driver: json-file
options: options:

View File

@@ -63,15 +63,35 @@ start_xvfb_if_needed() {
# IAAI blocks headless Chromium on Linux; run headed via Xvfb virtual display # IAAI blocks headless Chromium on Linux; run headed via Xvfb virtual display
export DISPLAY=:99 export DISPLAY=:99
export IAAI_HEADLESS=false export IAAI_HEADLESS=false
local xvfb_pid_file="/tmp/xvfb-99.pid"
if [ -f "${xvfb_pid_file}" ]; then
local existing_pid
existing_pid="$(cat "${xvfb_pid_file}")"
if [ -n "${existing_pid}" ] && kill -0 "${existing_pid}" 2>/dev/null; then
echo "[entrypoint] Reusing existing Xvfb on DISPLAY=${DISPLAY} (PID ${existing_pid})"
return 0
fi
echo "[entrypoint] Removing stale Xvfb pid file"
rm -f "${xvfb_pid_file}"
fi
if [ -S /tmp/.X11-unix/X99 ] || [ -f /tmp/.X99-lock ]; then if [ -S /tmp/.X11-unix/X99 ] || [ -f /tmp/.X99-lock ]; then
echo "[entrypoint] Reusing existing Xvfb on DISPLAY=${DISPLAY}" echo "[entrypoint] Removing stale Xvfb socket/lock files for DISPLAY=${DISPLAY}"
return 0 rm -f /tmp/.X11-unix/X99 /tmp/.X99-lock
fi fi
Xvfb :99 -screen 0 1920x1080x24 -nolisten tcp & Xvfb :99 -screen 0 1920x1080x24 -nolisten tcp &
XVFB_PID=$! XVFB_PID=$!
echo "${XVFB_PID}" > "${xvfb_pid_file}"
sleep 0.5 sleep 0.5
if ! kill -0 "${XVFB_PID}" 2>/dev/null; then
echo "[entrypoint] ERROR: Xvfb failed to start on DISPLAY=${DISPLAY}"
rm -f "${xvfb_pid_file}"
exit 1
fi
echo "[entrypoint] Xvfb started (PID ${XVFB_PID}, DISPLAY=${DISPLAY})" echo "[entrypoint] Xvfb started (PID ${XVFB_PID}, DISPLAY=${DISPLAY})"
} }

View File

@@ -12,8 +12,7 @@ from .routes import cars, health, tasks
@asynccontextmanager @asynccontextmanager
async def lifespan(app: FastAPI): async def lifespan(app: FastAPI):
# Жизненный цикл. # Жизненный цикл.
persistence: PersistenceService = app.state.persistence # Таблицы создаются через Alembic-миграции (сервис migrate).
persistence.create_tables()
yield yield

View File

@@ -7,6 +7,7 @@ from sqlalchemy import select, func
from ..deps import get_persistence from ..deps import get_persistence
from ...storage.db import PersistenceService from ...storage.db import PersistenceService
from ...storage.models import SyncRun from ...storage.models import SyncRun
from ...worker.celery_app import celery_app
from ...worker.tasks import sync_vehicle_task, sync_listing_task from ...worker.tasks import sync_vehicle_task, sync_listing_task
router = APIRouter() router = APIRouter()
@@ -64,6 +65,25 @@ def start_sync_listing(
} }
@router.get("/tasks/{task_id}")
def get_task_status(task_id: str):
result = celery_app.AsyncResult(task_id)
payload: dict = {
"task_id": task_id,
"state": result.state,
}
if result.successful():
payload["result"] = result.result
elif result.failed():
payload["error"] = str(result.result)
elif result.info is not None:
payload["meta"] = result.info
return payload
@router.get("/sync-runs") @router.get("/sync-runs")
def list_sync_runs( def list_sync_runs(
page: int = Query(1, ge=1), page: int = Query(1, ge=1),

View File

@@ -0,0 +1,10 @@
class ScraperError(Exception):
"""Base scraper exception."""
class AntiBotDetectedError(ScraperError):
"""Raised when the target website appears to block automation."""
class SiteStructureChangedError(ScraperError):
"""Raised when the page shape changed and required data is missing."""

View File

@@ -7,6 +7,8 @@ from typing import Any
from playwright.sync_api import Error, TimeoutError as PlaywrightTimeoutError from playwright.sync_api import Error, TimeoutError as PlaywrightTimeoutError
from .exceptions import AntiBotDetectedError
logger = logging.getLogger("iaai_scraper.retry") logger = logging.getLogger("iaai_scraper.retry")
# Типы исключений, при которых retry имеет смысл. # Типы исключений, при которых retry имеет смысл.
@@ -16,6 +18,7 @@ RETRYABLE_EXCEPTIONS = (
ConnectionError, ConnectionError,
OSError, OSError,
TimeoutError, TimeoutError,
AntiBotDetectedError,
) )

View File

@@ -236,12 +236,13 @@ class VehicleParser:
def _build_access_notes(self, summary: dict[str, Any], responses: list[dict[str, Any]]) -> dict[str, Any]: def _build_access_notes(self, summary: dict[str, Any], responses: list[dict[str, Any]]) -> dict[str, Any]:
endpoints = [item.get("url", "") for item in responses] endpoints = [item.get("url", "") for item in responses]
dom_hints = self._dom_hints(" ".join(str(value) for value in summary.values() if value is not None))
return { return {
"vin_visible": bool(summary.get("vin")), "vin_visible": bool(summary.get("vin")),
"images_visible": bool(summary.get("image_urls")), "images_visible": bool(summary.get("image_urls")),
"network_json_count": len(responses), "network_json_count": len(responses),
"possible_captcha": False, "possible_captcha": bool(dom_hints.get("has_captcha_text")),
"possible_antibot": False, "possible_antibot": bool(dom_hints.get("has_antibot_text")),
"observed_endpoints": endpoints[:20], "observed_endpoints": endpoints[:20],
} }

View File

@@ -10,6 +10,7 @@ from playwright.sync_api import TimeoutError as PlaywrightTimeoutError
from .browser import BrowserFactory, HumanPacer, NetworkCapture from .browser import BrowserFactory, HumanPacer, NetworkCapture
from .core.config import Settings, settings from .core.config import Settings, settings
from .core.exceptions import AntiBotDetectedError, SiteStructureChangedError
from .core.logs import set_trace_id, setup_logging from .core.logs import set_trace_id, setup_logging
from .core.retry import retryable from .core.retry import retryable
from .core.runtime_config import RuntimeConfig from .core.runtime_config import RuntimeConfig
@@ -26,6 +27,22 @@ VEHICLE_ID_RE = re.compile(r"/VehicleDetail/(\d+)(?:~[A-Z]{2})?", re.IGNORECASE)
class IAAIScraper: class IAAIScraper:
@staticmethod
def _raise_if_blocked_or_incomplete(parsed: dict, vehicle_url: str) -> None:
dom_hints = parsed.get("dom_hints", {}) or {}
access_notes = parsed.get("access_notes", {}) or {}
summary = parsed.get("vehicle_summary", {}) or {}
if dom_hints.get("has_captcha_text") or dom_hints.get("has_antibot_text"):
raise AntiBotDetectedError(f"IAAI anti-bot detected for {vehicle_url}")
if access_notes.get("possible_captcha") or access_notes.get("possible_antibot"):
raise AntiBotDetectedError(f"IAAI blocked or challenged request for {vehicle_url}")
has_identity = bool(summary.get("lot_number") or summary.get("vin") or summary.get("make") or summary.get("model"))
if not has_identity:
raise SiteStructureChangedError(f"Vehicle page returned no recognizable vehicle data: {vehicle_url}")
@staticmethod @staticmethod
def _extract_origin_id_from_url(vehicle_url: str) -> str | None: def _extract_origin_id_from_url(vehicle_url: str) -> str | None:
match = VEHICLE_ID_RE.search(vehicle_url) match = VEHICLE_ID_RE.search(vehicle_url)
@@ -152,6 +169,7 @@ class IAAIScraper:
dom_text = "" dom_text = ""
network_dump = capture.export() network_dump = capture.export()
parsed = self.vehicle_parser.normalize(vehicle_url, html, dom_text, network_dump) parsed = self.vehicle_parser.normalize(vehicle_url, html, dom_text, network_dump)
self._raise_if_blocked_or_incomplete(parsed, vehicle_url)
db_record = self.car_mapper.map_to_car_record( db_record = self.car_mapper.map_to_car_record(
vehicle_url=vehicle_url, vehicle_url=vehicle_url,
@@ -255,8 +273,6 @@ class IAAIScraper:
try: try:
listing = self.collect_listing(make=make, model=model) listing = self.collect_listing(make=make, model=model)
vehicle_urls = list(listing.get("vehicle_urls", [])) vehicle_urls = list(listing.get("vehicle_urls", []))
if limit is not None:
vehicle_urls = vehicle_urls[:max(0, limit)]
effective_only_new = self.settings.sync_only_new if only_new is None else only_new effective_only_new = self.settings.sync_only_new if only_new is None else only_new
if effective_only_new: if effective_only_new:
@@ -279,6 +295,9 @@ class IAAIScraper:
logger.info("Filtering already known vehicles: skipped %d", skipped_existing) logger.info("Filtering already known vehicles: skipped %d", skipped_existing)
vehicle_urls = [url for url in vehicle_urls if url not in known_urls] vehicle_urls = [url for url in vehicle_urls if url not in known_urls]
if limit is not None:
vehicle_urls = vehicle_urls[:max(0, limit)]
total = len(vehicle_urls) total = len(vehicle_urls)
logger.info("Starting sync: %d vehicles to process", total) logger.info("Starting sync: %d vehicles to process", total)
@@ -294,6 +313,7 @@ class IAAIScraper:
# Применяем фильтры из runtime_config (include/exclude/price/mileage/flags) # Применяем фильтры из runtime_config (include/exclude/price/mileage/flags)
if not self.runtime_config.filters.is_empty(): if not self.runtime_config.filters.is_empty():
vehicle_summary = scrape_result.get("vehicle_summary", {}) or {}
filter_values = { filter_values = {
"brand": record.brand, "brand": record.brand,
"model": record.model, "model": record.model,
@@ -302,9 +322,11 @@ class IAAIScraper:
"color": record.color, "color": record.color,
"drive": record.drive, "drive": record.drive,
"gearbox": record.gearbox, "gearbox": record.gearbox,
"location": vehicle_summary.get("location"),
"price": record.price, "price": record.price,
"mileage": record.mileage, "mileage": record.mileage,
"is_damaged": record.is_damaged, "is_damaged": record.is_damaged,
"run_and_drive": vehicle_summary.get("run_and_drive"),
} }
if not self.runtime_config.filters.matches(filter_values): if not self.runtime_config.filters.matches(filter_values):
logger.info( logger.info(

View File

@@ -1,9 +1,11 @@
# Celery-задачи для синхронизации автомобилей и листинга IAAI. # Celery-задачи для синхронизации автомобилей и листинга IAAI.
from concurrent.futures import ThreadPoolExecutor
import json import json
import logging import logging
from celery import shared_task from celery import shared_task
from redis import Redis
from ..core.config import Settings from ..core.config import Settings
from ..scraper import IAAIScraper from ..scraper import IAAIScraper
@@ -11,11 +13,64 @@ from ..storage.db import PersistenceService
logger = logging.getLogger("iaai_scraper.worker.tasks") logger = logging.getLogger("iaai_scraper.worker.tasks")
SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
def _run_browser_job(func, *args, **kwargs):
# Playwright Sync API нельзя запускать в потоке с активным asyncio loop.
# Celery/зависимости могут поднимать loop в worker-процессе, поэтому
# браузерный код выполняем в отдельном thread без loop.
with ThreadPoolExecutor(max_workers=1, thread_name_prefix="iaai-browser") as executor:
future = executor.submit(func, *args, **kwargs)
return future.result()
def _sync_listing_lock_ttl_seconds() -> int:
settings = Settings()
# Небольшой запас к hard time limit задачи, чтобы lock самоснимался после сбоев.
return max(settings.celery.task_time_limit + 120, 300)
def _get_persistence() -> PersistenceService: def _get_persistence() -> PersistenceService:
return PersistenceService(Settings()) return PersistenceService(Settings())
def _get_redis() -> Redis:
settings = Settings()
return Redis.from_url(settings.redis.url, decode_responses=True)
def _release_lock_if_owner(redis_client: Redis, key: str, owner: str) -> None:
try:
current_owner = redis_client.get(key)
if current_owner == owner:
redis_client.delete(key)
except Exception as exc:
logger.warning("Failed to release lock %s: %s", key, exc)
def _has_other_active_sync_listing_task(task) -> bool:
try:
inspector = task.app.control.inspect(timeout=1.0)
active_map = inspector.active() or {}
except Exception as exc:
logger.warning("Failed to inspect active tasks: %s", exc)
return False
current_task_id = task.request.id
for worker_tasks in active_map.values():
for item in worker_tasks or []:
name = str(item.get("name") or "")
task_id = str(item.get("id") or "")
if (
name == "iaai_scraper.worker.tasks.sync_listing_task"
and task_id
and task_id != current_task_id
):
return True
return False
@shared_task( @shared_task(
name="iaai_scraper.worker.tasks.sync_vehicle_task", name="iaai_scraper.worker.tasks.sync_vehicle_task",
bind=True, bind=True,
@@ -29,8 +84,11 @@ def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
persistence.create_tables() persistence.create_tables()
try: try:
def _job():
with IAAIScraper() as scraper: with IAAIScraper() as scraper:
result = scraper.sync_vehicle(vehicle_url, lane=lane) return scraper.sync_vehicle(vehicle_url, lane=lane)
result = _run_browser_job(_job)
logger.info("sync_vehicle_task completed: %s", vehicle_url) logger.info("sync_vehicle_task completed: %s", vehicle_url)
return { return {
"status": "success", "status": "success",
@@ -62,10 +120,59 @@ def sync_listing_task(
# Полный цикл: листинг + sync всех найденных машин. # Полный цикл: листинг + sync всех найденных машин.
persistence = _get_persistence() persistence = _get_persistence()
persistence.create_tables() persistence.create_tables()
task_id = self.request.id or "unknown"
redis_client = _get_redis()
lock_acquired = False
lock_ttl = _sync_listing_lock_ttl_seconds()
try:
lock_acquired = bool(
redis_client.set(
SYNC_LISTING_LOCK_KEY,
task_id,
nx=True,
ex=lock_ttl,
)
)
except Exception as exc:
logger.warning("Failed to acquire sync lock in Redis: %s", exc)
if not lock_acquired:
# Возможен stale lock после рестарта worker. Если активного sync_listing нет —
# снимаем lock и пытаемся взять его заново.
if not _has_other_active_sync_listing_task(self):
try:
stale_owner = redis_client.get(SYNC_LISTING_LOCK_KEY)
if stale_owner:
logger.warning(
"Removing stale sync lock held by task %s",
stale_owner,
)
redis_client.delete(SYNC_LISTING_LOCK_KEY)
lock_acquired = bool(
redis_client.set(
SYNC_LISTING_LOCK_KEY,
task_id,
nx=True,
ex=lock_ttl,
)
)
except Exception as exc:
logger.warning("Failed to recover stale sync lock: %s", exc)
if not lock_acquired:
logger.info("sync_listing_task skipped: another sync is already running")
return {
"status": "skipped",
"reason": "sync_already_running",
"task_id": task_id,
}
try: try:
self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id})
def _job():
with IAAIScraper() as scraper: with IAAIScraper() as scraper:
result = scraper.sync_listing( return scraper.sync_listing(
make=make, make=make,
model=model, model=model,
lane=lane, lane=lane,
@@ -73,7 +180,11 @@ def sync_listing_task(
only_new=only_new, only_new=only_new,
) )
result = _run_browser_job(_job)
summary = { summary = {
"task_id": task_id,
"run_id": result.get("run_id"),
"cars_upserted": result.get("cars_upserted", 0), "cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0), "cars_failed": result.get("cars_failed", 0),
"images_upserted": result.get("images_upserted", 0), "images_upserted": result.get("images_upserted", 0),
@@ -89,3 +200,6 @@ def sync_listing_task(
except Exception as exc: except Exception as exc:
logger.error("sync_listing_task failed: %s", exc) logger.error("sync_listing_task failed: %s", exc)
raise self.retry(exc=exc) raise self.retry(exc=exc)
finally:
if lock_acquired:
_release_lock_if_owner(redis_client, SYNC_LISTING_LOCK_KEY, task_id)

View File

@@ -57,6 +57,12 @@ class TestVehicleParserUnit(unittest.TestCase):
self.assertTrue(hints["has_captcha_text"]) self.assertTrue(hints["has_captcha_text"])
self.assertTrue(hints["has_antibot_text"]) self.assertTrue(hints["has_antibot_text"])
def test_access_notes_reflect_antibot_signals(self) -> None:
summary = {"vin": "", "image_urls": [], "note": "Incapsula access denied. Verify you are human."}
notes = self.parser._build_access_notes(summary, [])
self.assertTrue(notes["possible_captcha"])
self.assertTrue(notes["possible_antibot"])
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()

View File

@@ -4,6 +4,7 @@ import unittest
from unittest.mock import MagicMock from unittest.mock import MagicMock
from iaai_scraper.core.config import Settings from iaai_scraper.core.config import Settings
from iaai_scraper.core.exceptions import AntiBotDetectedError, SiteStructureChangedError
from iaai_scraper.scraper import IAAIScraper from iaai_scraper.scraper import IAAIScraper
from iaai_scraper.storage.schemas import CarRecord from iaai_scraper.storage.schemas import CarRecord
@@ -149,6 +150,28 @@ class TestScraperSync(unittest.TestCase):
self.assertIsNone(scraper.browser) self.assertIsNone(scraper.browser)
self.assertIsNone(scraper.playwright) self.assertIsNone(scraper.playwright)
def test_guard_raises_on_antibot_signals(self) -> None:
with self.assertRaises(AntiBotDetectedError):
IAAIScraper._raise_if_blocked_or_incomplete(
{
"dom_hints": {"has_captcha_text": True, "has_antibot_text": False},
"access_notes": {},
"vehicle_summary": {},
},
"https://www.iaai.com/VehicleDetail/999~US",
)
def test_guard_raises_on_empty_vehicle_page(self) -> None:
with self.assertRaises(SiteStructureChangedError):
IAAIScraper._raise_if_blocked_or_incomplete(
{
"dom_hints": {"has_captcha_text": False, "has_antibot_text": False},
"access_notes": {"possible_captcha": False, "possible_antibot": False},
"vehicle_summary": {},
},
"https://www.iaai.com/VehicleDetail/999~US",
)
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()