Compare commits
3 Commits
prod-prepa
...
3edd706110
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3edd706110 | ||
|
|
fd1d1d3c9f | ||
|
|
b0f44efdb8 |
@@ -29,8 +29,8 @@ IAAI_BEFORE_VEHICLE_OPEN_MIN_S=1.5
|
||||
IAAI_BEFORE_VEHICLE_OPEN_MAX_S=3.0
|
||||
IAAI_AFTER_VEHICLE_OPEN_MIN_S=1.0
|
||||
IAAI_AFTER_VEHICLE_OPEN_MAX_S=2.5
|
||||
IAAI_BETWEEN_VEHICLES_MIN_S=2.0
|
||||
IAAI_BETWEEN_VEHICLES_MAX_S=4.0
|
||||
IAAI_BETWEEN_VEHICLES_MIN_S=3.0
|
||||
IAAI_BETWEEN_VEHICLES_MAX_S=6.0
|
||||
IAAI_AFTER_PAGE_CHANGE_MIN_S=3.0
|
||||
IAAI_AFTER_PAGE_CHANGE_MAX_S=6.0
|
||||
|
||||
|
||||
6
.gitattributes
vendored
Normal file
6
.gitattributes
vendored
Normal file
@@ -0,0 +1,6 @@
|
||||
# Force LF for shell scripts (critical for Docker on Windows)
|
||||
*.sh text eol=lf
|
||||
entrypoint.sh text eol=lf
|
||||
|
||||
# Default for other files
|
||||
* text=auto
|
||||
@@ -11,5 +11,7 @@ RUN pip install --no-cache-dir -r requirements.txt
|
||||
COPY . .
|
||||
RUN chmod +x entrypoint.sh
|
||||
|
||||
STOPSIGNAL SIGINT
|
||||
|
||||
ENTRYPOINT ["./entrypoint.sh"]
|
||||
CMD ["python", "main.py", "--help"]
|
||||
CMD ["python", "main.py", "run-daemon"]
|
||||
|
||||
15
README.md
15
README.md
@@ -141,10 +141,23 @@ pytest tests -q
|
||||
## Docker
|
||||
|
||||
```bash
|
||||
# собрать образ
|
||||
docker compose build
|
||||
docker compose run --rm iaai-scraper python main.py --help
|
||||
|
||||
# запустить daemon (по умолчанию run-daemon)
|
||||
docker compose up -d
|
||||
|
||||
# посмотреть логи
|
||||
docker compose logs -f
|
||||
|
||||
# одноразовая команда
|
||||
docker compose run --rm iaai-scraper python main.py sync-listing --limit 5
|
||||
|
||||
# остановить (graceful shutdown)
|
||||
docker compose down
|
||||
```
|
||||
|
||||
Контейнер автоматически перезапускается при крашах (`restart: unless-stopped`).
|
||||
Для Docker прокси также задаются через `.env`.
|
||||
|
||||
## Защита от дублей
|
||||
|
||||
@@ -3,10 +3,13 @@ services:
|
||||
build: .
|
||||
container_name: iaai-scraper
|
||||
env_file:
|
||||
- .env
|
||||
- path: .env
|
||||
required: false
|
||||
volumes:
|
||||
- ./:/app
|
||||
- /app/.venv
|
||||
- /app/__pycache__
|
||||
working_dir: /app
|
||||
command: python main.py --help
|
||||
restart: unless-stopped
|
||||
stop_grace_period: 30s
|
||||
command: python main.py run-daemon
|
||||
|
||||
@@ -22,10 +22,10 @@ class NetworkCapture:
|
||||
_seen_resp: set[str] = field(default_factory=set)
|
||||
_origin: str | None = None
|
||||
|
||||
def attach(self, page: Page) -> None:
|
||||
def attach(self, page: Page, origin_url: str | None = None) -> None:
|
||||
# Подписываемся на request/response события страницы.
|
||||
try:
|
||||
self._origin = urlparse(page.url).netloc.lower() or None
|
||||
self._origin = urlparse(origin_url or page.url).netloc.lower() or None
|
||||
except Exception:
|
||||
self._origin = None
|
||||
page.on("request", self._on_request)
|
||||
@@ -47,7 +47,7 @@ class NetworkCapture:
|
||||
if len(self.requests) >= self.settings.gentle.max_requests:
|
||||
return
|
||||
key = f"{request.method}:{request.url}:{request.post_data or ''}"
|
||||
# Дедупликация одинаковых запросов.
|
||||
# Дедупликация одинаковых запросов.
|
||||
if key in self._seen_req:
|
||||
return
|
||||
self._seen_req.add(key)
|
||||
@@ -95,7 +95,7 @@ class NetworkCapture:
|
||||
|
||||
@staticmethod
|
||||
def _categorize(url: str) -> str:
|
||||
# Простая эвристика для разбивки ответов по смыслу.
|
||||
# Простая эвристика для разбивки ответов по смыслу.
|
||||
low = url.lower()
|
||||
mapping = {
|
||||
"images": ["image", "media", "photos", "gallery"],
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import argparse
|
||||
from pathlib import Path
|
||||
|
||||
from .core.config import Settings
|
||||
from .core.utils import save_to_json
|
||||
from .scraper import IAAIScraper
|
||||
|
||||
@@ -75,8 +76,6 @@ def main() -> None:
|
||||
|
||||
runtime_settings: Settings | None = None
|
||||
if args.headless is not None or args.debug:
|
||||
from .core.config import Settings
|
||||
|
||||
runtime_settings = Settings()
|
||||
if args.headless is not None:
|
||||
runtime_settings.headless = args.headless == "true"
|
||||
@@ -85,7 +84,6 @@ def main() -> None:
|
||||
|
||||
if args.command == "run-daemon":
|
||||
# daemon: бесконечный цикл
|
||||
from .core.config import Settings
|
||||
runtime_settings = runtime_settings or Settings()
|
||||
if args.interval is not None:
|
||||
runtime_settings.scheduler_interval_minutes = args.interval
|
||||
|
||||
@@ -130,7 +130,10 @@ class CarMapper:
|
||||
"brand": brand, "model": model, "year": year, "price": price, "mileage": mileage,
|
||||
"color": color, "drive": drive, "gearbox": gearbox, "body_type": body_type,
|
||||
"engine_volume": engine_volume, "is_damaged": is_damaged, "is_sold": is_sold,
|
||||
"image_count": len(images_records),
|
||||
"country": country, "selling_type": "AUCTION", "one_owner": one_owner,
|
||||
"new_car": new_car, "evaluation": evaluation, "non_smoking": non_smoking,
|
||||
"rental": rental, "repair_history": repair_history,
|
||||
"images": [image.fullres_image for image in images_records],
|
||||
}, sort_keys=True, default=str).encode()).hexdigest()
|
||||
|
||||
return CarRecord(
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import logging
|
||||
import signal
|
||||
import time
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
@@ -39,6 +40,7 @@ class IAAIScraper:
|
||||
self.vehicle_parser = VehicleParser()
|
||||
self.car_mapper = CarMapper()
|
||||
self.persistence = PersistenceService(self.settings)
|
||||
self._shutdown_requested = False
|
||||
|
||||
# browser lifecycle
|
||||
|
||||
@@ -170,7 +172,7 @@ class IAAIScraper:
|
||||
trace_id = self._new_trace_id("scrape")
|
||||
started_at = time.perf_counter()
|
||||
capture = NetworkCapture(self.settings)
|
||||
capture.attach(page)
|
||||
capture.attach(page, origin_url=vehicle_url)
|
||||
|
||||
page.goto(vehicle_url, wait_until="domcontentloaded", timeout=60_000)
|
||||
try:
|
||||
@@ -264,24 +266,26 @@ class IAAIScraper:
|
||||
trace_id = self._new_trace_id("sync-listing")
|
||||
started_at = time.perf_counter()
|
||||
self.persistence.create_tables()
|
||||
listing = self.collect_listing(make=make, model=model)
|
||||
vehicle_urls = list(listing.get("vehicle_urls", []))
|
||||
if limit is not None:
|
||||
vehicle_urls = vehicle_urls[:max(0, limit)]
|
||||
run_id = self.persistence.start_sync_run(lane=lane)
|
||||
cars_upserted = 0
|
||||
cars_failed = 0
|
||||
images_upserted = 0
|
||||
total = 0
|
||||
failures: list[dict[str, str]] = []
|
||||
listing: dict = {}
|
||||
|
||||
total = len(vehicle_urls)
|
||||
logger.info("Starting sync: %d vehicles to process", total)
|
||||
|
||||
page = self._get_page()
|
||||
try:
|
||||
listing = self.collect_listing(make=make, model=model)
|
||||
vehicle_urls = list(listing.get("vehicle_urls", []))
|
||||
if limit is not None:
|
||||
vehicle_urls = vehicle_urls[:max(0, limit)]
|
||||
|
||||
total = len(vehicle_urls)
|
||||
logger.info("Starting sync: %d vehicles to process", total)
|
||||
|
||||
for index, vehicle_url in enumerate(vehicle_urls, start=1):
|
||||
page = self._get_page()
|
||||
try:
|
||||
# Повторно используем страницу, чтобы не создавать лишний overhead.
|
||||
logger.info("[%d/%d] Scraping %s", index, total, vehicle_url)
|
||||
scrape_result = self._scrape_on_page(page, vehicle_url)
|
||||
db_record = scrape_result.get("db_record")
|
||||
@@ -289,7 +293,8 @@ class IAAIScraper:
|
||||
raise RuntimeError("Scrape result does not contain db_record")
|
||||
record = CarRecord.model_validate(db_record)
|
||||
upsert = self.persistence.upsert_car(record)
|
||||
cars_upserted += 1
|
||||
if upsert.get("action") != "skipped":
|
||||
cars_upserted += 1
|
||||
img_count = int(upsert.get("images_upserted", 0))
|
||||
images_upserted += img_count
|
||||
logger.info(
|
||||
@@ -299,35 +304,35 @@ class IAAIScraper:
|
||||
record.year or "?", record.price or "N/A", img_count,
|
||||
)
|
||||
except Exception as exc:
|
||||
# После сбоя пересоздаём page, чтобы не остаться в битом состоянии.
|
||||
cars_failed += 1
|
||||
failures.append({"vehicle_url": vehicle_url, "error": str(exc)})
|
||||
logger.error("[%d/%d] Failed %s: %s", index, total, vehicle_url, exc)
|
||||
finally:
|
||||
try:
|
||||
page.close()
|
||||
except Exception:
|
||||
pass
|
||||
page = self._get_page()
|
||||
|
||||
if index < total:
|
||||
self.pacer.between_vehicles()
|
||||
except Exception as exc:
|
||||
# collect_listing itself failed — count as total failure
|
||||
if not failures:
|
||||
failures.append({"vehicle_url": "collect_listing", "error": str(exc)})
|
||||
logger.error("sync_listing failed: %s", exc)
|
||||
finally:
|
||||
try:
|
||||
page.close()
|
||||
except Exception:
|
||||
pass
|
||||
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
|
||||
error_summary = "; ".join(item["error"] for item in failures[:10]) if failures else None
|
||||
self.persistence.finish_sync_run(
|
||||
run_id,
|
||||
status=status,
|
||||
ids_fetched=total,
|
||||
cars_upserted=cars_upserted,
|
||||
cars_failed=cars_failed,
|
||||
images_upserted=images_upserted,
|
||||
error_summary=error_summary,
|
||||
)
|
||||
|
||||
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
|
||||
error_summary = "; ".join(item["error"] for item in failures[:10]) if failures else None
|
||||
self.persistence.finish_sync_run(
|
||||
run_id,
|
||||
status=status,
|
||||
ids_fetched=total,
|
||||
cars_upserted=cars_upserted,
|
||||
cars_failed=cars_failed,
|
||||
images_upserted=images_upserted,
|
||||
error_summary=error_summary,
|
||||
)
|
||||
logger.info(
|
||||
"Sync run #%d finished: %d/%d upserted, %d failed, %d images",
|
||||
run_id, cars_upserted, total, cars_failed, images_upserted,
|
||||
@@ -348,17 +353,26 @@ class IAAIScraper:
|
||||
|
||||
def run_scheduled(self) -> None:
|
||||
# Простой бесконечный цикл без внешнего планировщика.
|
||||
# Graceful shutdown по SIGINT/SIGTERM.
|
||||
def _handle_shutdown(signum, frame):
|
||||
logger.info("Received signal %s, shutting down gracefully...", signum)
|
||||
self._shutdown_requested = True
|
||||
|
||||
signal.signal(signal.SIGINT, _handle_shutdown)
|
||||
signal.signal(signal.SIGTERM, _handle_shutdown)
|
||||
|
||||
interval = self.settings.scheduler_interval_minutes * 60
|
||||
logger.info(
|
||||
"Scheduler started: syncing every %d minutes",
|
||||
self.settings.scheduler_interval_minutes,
|
||||
)
|
||||
cycle = 0
|
||||
while True:
|
||||
while not self._shutdown_requested:
|
||||
cycle += 1
|
||||
logger.info("=== Scheduler cycle #%d starting ===", cycle)
|
||||
start = time.time()
|
||||
try:
|
||||
# Сбрасываем контекст перед каждым циклом.
|
||||
if self.context:
|
||||
try:
|
||||
self.context.close()
|
||||
@@ -366,6 +380,12 @@ class IAAIScraper:
|
||||
pass
|
||||
self.context = None
|
||||
|
||||
# Проверяем, что browser/playwright живы; пересоздаём при необходимости.
|
||||
if self.browser is None or self.playwright is None:
|
||||
logger.info("Browser/Playwright not available, re-initializing...")
|
||||
self.close()
|
||||
self.__enter__()
|
||||
|
||||
result = self.sync_listing()
|
||||
elapsed = time.time() - start
|
||||
logger.info(
|
||||
@@ -377,17 +397,31 @@ class IAAIScraper:
|
||||
except Exception as exc:
|
||||
elapsed = time.time() - start
|
||||
logger.error("Cycle #%d failed after %.1fs: %s", cycle, elapsed, exc)
|
||||
if self.context:
|
||||
try:
|
||||
self.context.close()
|
||||
except PlaywrightError:
|
||||
pass
|
||||
self.context = None
|
||||
# Полный сброс при любой ошибке цикла — следующий цикл пересоздаст всё.
|
||||
try:
|
||||
self.close()
|
||||
except Exception:
|
||||
pass
|
||||
# Пересоздаём browser для следующего цикла.
|
||||
try:
|
||||
self.__enter__()
|
||||
except Exception as reinit_exc:
|
||||
logger.error("Failed to re-initialize browser: %s", reinit_exc)
|
||||
|
||||
if self._shutdown_requested:
|
||||
break
|
||||
|
||||
sleep_time = max(0, interval - (time.time() - start))
|
||||
if sleep_time > 0:
|
||||
logger.info("Sleeping %.0f seconds until next cycle...", sleep_time)
|
||||
time.sleep(sleep_time)
|
||||
# Прерываемый sleep — проверяем shutdown каждые 5 секунд.
|
||||
slept = 0.0
|
||||
while slept < sleep_time and not self._shutdown_requested:
|
||||
chunk = min(5.0, sleep_time - slept)
|
||||
time.sleep(chunk)
|
||||
slept += chunk
|
||||
|
||||
logger.info("Scheduler stopped gracefully after %d cycles.", cycle)
|
||||
|
||||
# helpers
|
||||
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import json
|
||||
import logging
|
||||
from contextlib import contextmanager
|
||||
from datetime import datetime, timezone
|
||||
@@ -13,6 +14,41 @@ from .schemas import CarRecord
|
||||
logger = logging.getLogger("iaai_scraper.db")
|
||||
|
||||
|
||||
CAR_DB_FIELDS = {
|
||||
"parser_id",
|
||||
"brand",
|
||||
"model",
|
||||
"year",
|
||||
"price",
|
||||
"currency",
|
||||
"mileage",
|
||||
"country",
|
||||
"is_sold",
|
||||
"color",
|
||||
"drive",
|
||||
"gearbox",
|
||||
"steering_wheel",
|
||||
"body_type",
|
||||
"engine_volume",
|
||||
"selling_type",
|
||||
"one_owner",
|
||||
"new_car",
|
||||
"is_hidden",
|
||||
"origin",
|
||||
"origin_url",
|
||||
"origin_id",
|
||||
"is_damaged",
|
||||
"evaluation",
|
||||
"non_smoking",
|
||||
"rental",
|
||||
"repair_history",
|
||||
"slug",
|
||||
"last_seen_at",
|
||||
"content_hash",
|
||||
"raw_attributes",
|
||||
}
|
||||
|
||||
|
||||
class PersistenceService:
|
||||
|
||||
def __init__(self, settings: Settings) -> None:
|
||||
@@ -64,15 +100,21 @@ class PersistenceService:
|
||||
for image_payload in images:
|
||||
session.add(Image(fullres_image=str(image_payload["fullres_image"]), preview_image=str(image_payload["preview_image"]), order_index=int(image_payload.get("order_index", 0)), car_id=car_id))
|
||||
|
||||
@staticmethod
|
||||
def _car_payload(record: CarRecord) -> dict[str, object]:
|
||||
payload = record.model_dump(mode="python")
|
||||
result = {key: value for key, value in payload.items() if key in CAR_DB_FIELDS}
|
||||
# Serialize raw_attributes dict to JSON string for Text column.
|
||||
if "raw_attributes" in result and isinstance(result["raw_attributes"], dict):
|
||||
result["raw_attributes"] = json.dumps(result["raw_attributes"], ensure_ascii=False, default=str)
|
||||
return result
|
||||
|
||||
def upsert_car(self, record: CarRecord):
|
||||
"""Insert/update/skip по content_hash."""
|
||||
# Готовим payload отдельно от вложенных изображений и служебных полей.
|
||||
payload = record.model_dump(mode="python")
|
||||
images = payload.pop("images", [])
|
||||
payload.pop("raw_attributes", None)
|
||||
payload.pop("mapping_notes", None)
|
||||
content_hash = payload.pop("content_hash", "")
|
||||
payload["content_hash"] = content_hash
|
||||
# В БД отправляем только поля, реально существующие в финальной схеме cars.
|
||||
payload = self._car_payload(record)
|
||||
images = [image.model_dump(mode="python") for image in record.images]
|
||||
content_hash = str(payload.get("content_hash") or "")
|
||||
with self.session_scope() as session:
|
||||
# поиск по origin_id
|
||||
car = session.execute(select(Car).where(Car.origin_id == record.origin_id)).scalar_one_or_none()
|
||||
|
||||
@@ -54,6 +54,7 @@ class Car(Base):
|
||||
slug: Mapped[str] = mapped_column(String(), nullable=False)
|
||||
last_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now())
|
||||
content_hash: Mapped[str] = mapped_column(String(64), nullable=False, default="", index=True)
|
||||
raw_attributes: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
images: Mapped[list["Image"]] = relationship("Image", back_populates="car", cascade="all, delete-orphan")
|
||||
|
||||
|
||||
|
||||
@@ -104,6 +104,18 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
|
||||
self.assertEqual(result["action"], "skipped")
|
||||
self.assertEqual(result["images_upserted"], 0)
|
||||
|
||||
def test_persistence_ignores_non_db_fields(self) -> None:
|
||||
record = self._record("1000", content_hash="schema-test")
|
||||
record.raw_attributes = {"vin": "123"}
|
||||
record.mapping_notes = ["note"]
|
||||
|
||||
result = self.persistence.upsert_car(record)
|
||||
|
||||
self.assertEqual(result["action"], "inserted")
|
||||
with self.persistence.session_scope() as session:
|
||||
car = session.execute(select(Car).where(Car.origin_id == "1000")).scalar_one()
|
||||
self.assertEqual(car.origin_id, "1000")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -25,6 +25,30 @@ class TestCarMapper(unittest.TestCase):
|
||||
# Длина hex-представления SHA-256
|
||||
self.assertEqual(len(record.content_hash), 64)
|
||||
|
||||
def test_content_hash_changes_when_image_set_changes(self) -> None:
|
||||
record_one = self.mapper.map_to_car_record(
|
||||
vehicle_url="https://www.iaai.com/VehicleDetail/45089484~US",
|
||||
vehicle_summary={
|
||||
"make": "Toyota",
|
||||
"model": "Camry",
|
||||
"year": "2014",
|
||||
"image_urls": ["https://vis.iaai.com/resizer?imageKeys=1&width=845&height=633"],
|
||||
},
|
||||
payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}},
|
||||
)
|
||||
record_two = self.mapper.map_to_car_record(
|
||||
vehicle_url="https://www.iaai.com/VehicleDetail/45089484~US",
|
||||
vehicle_summary={
|
||||
"make": "Toyota",
|
||||
"model": "Camry",
|
||||
"year": "2014",
|
||||
"image_urls": ["https://vis.iaai.com/resizer?imageKeys=2&width=845&height=633"],
|
||||
},
|
||||
payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}},
|
||||
)
|
||||
|
||||
self.assertNotEqual(record_one.content_hash, record_two.content_hash)
|
||||
|
||||
def test_deduplicates_images_by_image_key(self) -> None:
|
||||
urls = [
|
||||
"https://vis.iaai.com/resizer?imageKeys=1&width=200&height=150",
|
||||
|
||||
@@ -52,8 +52,9 @@ class TestScraperSync(unittest.TestCase):
|
||||
scraper.persistence.upsert_car = MagicMock(return_value={"action": "inserted", "images_upserted": 1})
|
||||
|
||||
scraper.collect_listing = MagicMock(return_value={"vehicle_urls": ["https://www.iaai.com/VehicleDetail/222~US"]})
|
||||
page = MagicMock()
|
||||
scraper._scrape_on_page = MagicMock(return_value={"db_record": make_db_record("222")})
|
||||
scraper._get_page = MagicMock(return_value=MagicMock())
|
||||
scraper._get_page = MagicMock(return_value=page)
|
||||
scraper.car_mapper.map_to_car_record = MagicMock(side_effect=AssertionError("should not be called"))
|
||||
|
||||
result = scraper.sync_listing()
|
||||
@@ -63,6 +64,7 @@ class TestScraperSync(unittest.TestCase):
|
||||
self.assertIn("trace_id", result)
|
||||
self.assertIn("elapsed_seconds", result)
|
||||
self.assertEqual(scraper.persistence.upsert_car.call_count, 1)
|
||||
page.close.assert_called_once()
|
||||
|
||||
def test_sync_listing_respects_limit(self) -> None:
|
||||
scraper = self._make_scraper()
|
||||
@@ -84,6 +86,21 @@ class TestScraperSync(unittest.TestCase):
|
||||
|
||||
self.assertEqual(scraper._scrape_on_page.call_count, 1)
|
||||
|
||||
def test_sync_listing_does_not_count_skipped_as_upserted(self) -> None:
|
||||
scraper = self._make_scraper()
|
||||
|
||||
scraper.persistence.create_tables = MagicMock()
|
||||
scraper.persistence.start_sync_run = MagicMock(return_value=4)
|
||||
scraper.persistence.finish_sync_run = MagicMock()
|
||||
scraper.persistence.upsert_car = MagicMock(return_value={"action": "skipped", "images_upserted": 0})
|
||||
scraper.collect_listing = MagicMock(return_value={"vehicle_urls": ["https://www.iaai.com/VehicleDetail/444~US"]})
|
||||
scraper._scrape_on_page = MagicMock(return_value={"db_record": make_db_record("444")})
|
||||
scraper._get_page = MagicMock(return_value=MagicMock())
|
||||
|
||||
result = scraper.sync_listing()
|
||||
|
||||
self.assertEqual(result["cars_upserted"], 0)
|
||||
|
||||
def test_close_resets_browser_state(self) -> None:
|
||||
scraper = self._make_scraper()
|
||||
scraper.context = MagicMock()
|
||||
|
||||
Reference in New Issue
Block a user