Compare commits
6 Commits
6d1702faae
...
db2fa5ec17
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
db2fa5ec17 | ||
|
|
18dab6d48d | ||
|
|
3870cf4d94 | ||
|
|
62aafae793 | ||
|
|
acf30fcc20 | ||
|
|
7b1501008e |
@@ -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"]
|
||||||
|
|||||||
@@ -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:
|
||||||
|
|||||||
@@ -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})"
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -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),
|
||||||
|
|||||||
10
iaai_scraper/core/exceptions.py
Normal file
10
iaai_scraper/core/exceptions.py
Normal 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."""
|
||||||
@@ -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,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -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],
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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(
|
||||||
|
|||||||
@@ -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:
|
||||||
with IAAIScraper() as scraper:
|
def _job():
|
||||||
result = scraper.sync_vehicle(vehicle_url, lane=lane)
|
with IAAIScraper() as scraper:
|
||||||
|
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,18 +120,71 @@ 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:
|
||||||
with IAAIScraper() as scraper:
|
self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id})
|
||||||
result = scraper.sync_listing(
|
def _job():
|
||||||
make=make,
|
with IAAIScraper() as scraper:
|
||||||
model=model,
|
return scraper.sync_listing(
|
||||||
lane=lane,
|
make=make,
|
||||||
limit=limit,
|
model=model,
|
||||||
only_new=only_new,
|
lane=lane,
|
||||||
)
|
limit=limit,
|
||||||
|
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)
|
||||||
|
|||||||
@@ -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()
|
||||||
|
|||||||
@@ -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()
|
||||||
|
|||||||
Reference in New Issue
Block a user