5 Commits

Author SHA1 Message Date
qananasikq
54fea100dc add docker proxy bridge 2026-04-08 18:29:06 +03:00
qananasikq
447a88feed persist raw attributes and tighten sync assertions 2026-04-08 14:30:57 +03:00
qananasikq
b2d3902bcc stabilize daemon sync and docker startup 2026-04-08 14:30:44 +03:00
qananasikq
d1844ba625 document docker runtime and tune env defaults 2026-04-08 14:30:27 +03:00
ff2a1c7ac8 Обновить README.md 2026-04-08 12:38:15 +02:00
18 changed files with 564 additions and 59 deletions

View File

@@ -14,3 +14,4 @@ coverage.xml
*.egg-info/
dist/
build/
artifacts/

View File

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

View File

@@ -5,11 +5,17 @@ ENV PYTHONDONTWRITEBYTECODE=1 \
WORKDIR /app
# Xvfb for headed Chromium in container (IAAI blocks headless)
RUN apt-get update -qq && apt-get install -y --no-install-recommends xvfb \
&& rm -rf /var/lib/apt/lists/*
COPY requirements.txt ./
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"]

View File

@@ -1,6 +1,6 @@
# IAAI Scraper
Скрапер публичного листинга автомобилей с сайта IAAI.
Скрапер листинга автомобилей с сайта IAAI.
Собирает данные карточек через Playwright, парсит HTML и перехваченные JSON ответы,
нормализует и сохраняет в PostgreSQL или SQLite через SQLAlchemy.
@@ -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`.
## Защита от дублей

View File

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

View File

@@ -1,4 +1,31 @@
#!/bin/bash
set -e
# Start HTTP→SOCKS5 proxy bridge if SOCKS5 upstream is configured
if [ -n "$SOCKS5_PROXY_HOST" ]; then
echo "[entrypoint] Starting proxy bridge (HTTP :8899 → SOCKS5 $SOCKS5_PROXY_HOST:${SOCKS5_PROXY_PORT:-1002})..."
if [ -z "$IAAI_PROXY_SERVER" ]; then
export IAAI_PROXY_SERVER="http://127.0.0.1:8899"
elif [ "$IAAI_PROXY_SERVER" = "http://localhost:8899" ]; then
export IAAI_PROXY_SERVER="http://127.0.0.1:8899"
fi
echo "[entrypoint] Using browser proxy: $IAAI_PROXY_SERVER"
python -m iaai_scraper.proxy_bridge &
BRIDGE_PID=$!
sleep 1
if ! kill -0 $BRIDGE_PID 2>/dev/null; then
echo "[entrypoint] ERROR: iaai_scraper.proxy_bridge failed to start"
exit 1
fi
echo "[entrypoint] Proxy bridge started (PID $BRIDGE_PID)"
fi
# IAAI blocks headless Chromium on Linux; run headed via Xvfb virtual display
export DISPLAY=:99
export IAAI_HEADLESS=false
Xvfb :99 -screen 0 1920x1080x24 -nolisten tcp &
XVFB_PID=$!
sleep 0.5
echo "[entrypoint] Xvfb started (PID $XVFB_PID, DISPLAY=$DISPLAY)"
exec "$@"

View File

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

View File

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

View File

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

View File

@@ -0,0 +1,317 @@
from __future__ import annotations
import logging
import os
import select
import socket
import socketserver
import struct
import threading
from urllib.parse import urlsplit
BUFFER_SIZE = 65536
CRLF = b"\r\n"
DEFAULT_LISTEN_HOST = os.getenv("PROXY_BRIDGE_HOST", "0.0.0.0")
DEFAULT_LISTEN_PORT = int(os.getenv("PROXY_BRIDGE_PORT", "8899"))
SOCKS5_HOST = os.getenv("SOCKS5_PROXY_HOST", "")
SOCKS5_PORT = int(os.getenv("SOCKS5_PROXY_PORT", "1002"))
SOCKS5_USER = os.getenv("SOCKS5_PROXY_USER", "")
SOCKS5_PASS = os.getenv("SOCKS5_PROXY_PASS", "")
RELAY_IDLE_TIMEOUT_SECONDS = int(os.getenv("PROXY_BRIDGE_RELAY_IDLE_TIMEOUT_SECONDS", "60"))
MAX_WORKERS = int(os.getenv("PROXY_BRIDGE_MAX_WORKERS", "64"))
logger = logging.getLogger("proxy_bridge")
class ThreadingTCPServer(socketserver.ThreadingMixIn, socketserver.TCPServer):
allow_reuse_address = True
daemon_threads = True
def __init__(self, server_address, request_handler_class):
super().__init__(server_address, request_handler_class)
self._worker_semaphore = threading.BoundedSemaphore(MAX_WORKERS)
def process_request_thread(self, request, client_address):
with self._worker_semaphore:
super().process_request_thread(request, client_address)
def _recv_exact(sock: socket.socket, size: int) -> bytes:
data = b""
while len(data) < size:
chunk = sock.recv(size - len(data))
if not chunk:
raise ConnectionError("Unexpected EOF from SOCKS5 server")
data += chunk
return data
def _socks5_connect(host: str, port: int) -> socket.socket:
if not SOCKS5_HOST:
raise RuntimeError("SOCKS5_PROXY_HOST is not configured")
upstream = socket.create_connection((SOCKS5_HOST, SOCKS5_PORT), timeout=30)
upstream.settimeout(30)
methods = [0x00]
if SOCKS5_USER or SOCKS5_PASS:
methods = [0x02]
upstream.sendall(bytes([0x05, len(methods), *methods]))
version, method = _recv_exact(upstream, 2)
if version != 0x05 or method == 0xFF:
upstream.close()
raise ConnectionError("SOCKS5 authentication negotiation failed")
if method == 0x02:
username = SOCKS5_USER.encode("utf-8")
password = SOCKS5_PASS.encode("utf-8")
if len(username) > 255 or len(password) > 255:
upstream.close()
raise ValueError("SOCKS5 username/password too long")
upstream.sendall(bytes([0x01, len(username)]) + username + bytes([len(password)]) + password)
auth_version, auth_status = _recv_exact(upstream, 2)
if auth_version != 0x01 or auth_status != 0x00:
upstream.close()
raise ConnectionError("SOCKS5 username/password authentication failed")
try:
socket.inet_aton(host)
addr_type = 0x01
addr_payload = socket.inet_aton(host)
except OSError:
host_bytes = host.encode("idna")
if len(host_bytes) > 255:
upstream.close()
raise ValueError("Target host is too long for SOCKS5 domain format")
addr_type = 0x03
addr_payload = bytes([len(host_bytes)]) + host_bytes
request = bytes([0x05, 0x01, 0x00, addr_type]) + addr_payload + struct.pack("!H", port)
upstream.sendall(request)
response_head = _recv_exact(upstream, 4)
version, reply, _reserved, reply_addr_type = response_head
if version != 0x05 or reply != 0x00:
upstream.close()
raise ConnectionError(f"SOCKS5 connect failed with code {reply}")
if reply_addr_type == 0x01:
_recv_exact(upstream, 4)
elif reply_addr_type == 0x03:
domain_len = _recv_exact(upstream, 1)[0]
_recv_exact(upstream, domain_len)
elif reply_addr_type == 0x04:
_recv_exact(upstream, 16)
_recv_exact(upstream, 2)
upstream.settimeout(RELAY_IDLE_TIMEOUT_SECONDS)
return upstream
def _relay_bidirectional(left: socket.socket, right: socket.socket) -> None:
sockets = [left, right]
left.settimeout(RELAY_IDLE_TIMEOUT_SECONDS)
right.settimeout(RELAY_IDLE_TIMEOUT_SECONDS)
try:
while True:
readable, _, exceptional = select.select(sockets, [], sockets, RELAY_IDLE_TIMEOUT_SECONDS)
if exceptional:
break
if not readable:
logger.debug("Relay idle timeout reached; closing sockets")
return
for current in readable:
other = right if current is left else left
data = current.recv(BUFFER_SIZE)
if not data:
return
other.sendall(data)
finally:
for sock in sockets:
try:
sock.shutdown(socket.SHUT_RDWR)
except OSError:
pass
try:
sock.close()
except OSError:
pass
class ProxyHandler(socketserver.StreamRequestHandler):
def handle(self) -> None:
try:
request_line = self.rfile.readline(BUFFER_SIZE).decode("iso-8859-1").strip()
if not request_line:
return
method, target, version = request_line.split()
headers = self._read_headers()
logger.info("%s %s", method, target)
if method.upper() == "CONNECT":
host, port = self._parse_connect_target(target)
logger.info("CONNECT %s:%s", host, port)
upstream = _socks5_connect(host, port)
self.wfile.write(f"{version} 200 Connection Established".encode("ascii") + CRLF + CRLF)
self.wfile.flush()
_relay_bidirectional(self.connection, upstream)
return
host, port, path = self._parse_forward_target(target, headers)
upstream = _socks5_connect(host, port)
self._send_forward_request(upstream, method, path, version, headers)
body = self._read_request_body(headers)
if body:
upstream.sendall(body)
_relay_bidirectional(self.connection, upstream)
except Exception as exc:
logger.exception("Proxy bridge request failed: %s", exc)
try:
self.wfile.write(
b"HTTP/1.1 502 Bad Gateway" + CRLF
+ b"Connection: close" + CRLF
+ b"Content-Type: text/plain; charset=utf-8" + CRLF + CRLF
+ b"Bad Gateway"
)
self.wfile.flush()
except OSError:
pass
def _read_headers(self) -> list[tuple[str, str]]:
headers: list[tuple[str, str]] = []
while True:
line = self.rfile.readline(BUFFER_SIZE)
if line in {CRLF, b"\n", b""}:
break
decoded = line.decode("iso-8859-1")
if ":" not in decoded:
continue
name, value = decoded.split(":", 1)
headers.append((name.strip(), value.strip()))
return headers
@staticmethod
def _parse_connect_target(target: str) -> tuple[str, int]:
if target.startswith("["):
end = target.find("]")
if end == -1 or len(target) <= end + 2 or target[end + 1] != ":":
raise ValueError("Invalid CONNECT target")
host = target[1:end]
port_text = target[end + 2 :]
return host, int(port_text)
host, port_text = target.rsplit(":", 1)
return host, int(port_text)
@staticmethod
def _parse_forward_target(target: str, headers: list[tuple[str, str]]) -> tuple[str, int, str]:
if target.startswith("http://"):
parts = urlsplit(target)
port = parts.port or 80
path = parts.path or "/"
if parts.query:
path += f"?{parts.query}"
return parts.hostname or "", port, path
if target.startswith("https://"):
raise ValueError("HTTPS absolute-form request must use CONNECT")
host_header = next((value for name, value in headers if name.lower() == "host"), "")
if not host_header:
raise ValueError("Missing Host header")
if ":" in host_header:
host, port_text = host_header.rsplit(":", 1)
return host, int(port_text), target
return host_header, 80, target
def _send_forward_request(
self,
upstream: socket.socket,
method: str,
path: str,
version: str,
headers: list[tuple[str, str]],
) -> None:
filtered_headers: list[tuple[str, str]] = []
hop_by_hop = {
"proxy-connection",
"proxy-authorization",
"connection",
"keep-alive",
"te",
"trailer",
"transfer-encoding",
"upgrade",
}
for name, value in headers:
if name.lower() in hop_by_hop:
continue
filtered_headers.append((name, value))
request_head = [f"{method} {path} {version}\r\n"]
request_head.extend(f"{name}: {value}\r\n" for name, value in filtered_headers)
request_head.append("\r\n")
upstream.sendall("".join(request_head).encode("iso-8859-1"))
def _read_request_body(self, headers: list[tuple[str, str]]) -> bytes:
transfer_encoding = next((value for name, value in headers if name.lower() == "transfer-encoding"), "")
if "chunked" in transfer_encoding.lower():
return self._read_chunked_request_body()
content_length = next((value for name, value in headers if name.lower() == "content-length"), None)
if not content_length:
return b""
return self.rfile.read(int(content_length))
def _read_chunked_request_body(self) -> bytes:
chunks: list[bytes] = []
while True:
size_line = self.rfile.readline(BUFFER_SIZE)
if not size_line:
raise ConnectionError("Unexpected EOF in chunked request")
size_text = size_line.strip().split(b";", 1)[0]
chunk_size = int(size_text, 16)
chunks.append(size_line)
if chunk_size == 0:
while True:
trailer_line = self.rfile.readline(BUFFER_SIZE)
if not trailer_line:
raise ConnectionError("Unexpected EOF in chunked trailers")
chunks.append(trailer_line)
if trailer_line in {CRLF, b"\n"}:
return b"".join(chunks)
chunk_data = self.rfile.read(chunk_size)
if len(chunk_data) != chunk_size:
raise ConnectionError("Unexpected EOF in chunk body")
chunks.append(chunk_data)
chunk_end = self.rfile.read(2)
if chunk_end != CRLF:
raise ConnectionError("Invalid chunk terminator")
chunks.append(chunk_end)
def main() -> None:
if not SOCKS5_HOST:
raise SystemExit("SOCKS5_PROXY_HOST is required")
logging.basicConfig(
level=os.getenv("PROXY_BRIDGE_LOG_LEVEL", "INFO").upper(),
format="[%(asctime)s] [proxy_bridge] %(levelname)s: %(message)s",
)
with ThreadingTCPServer((DEFAULT_LISTEN_HOST, DEFAULT_LISTEN_PORT), ProxyHandler) as server:
logger.info(
"Listening on %s:%s -> socks5://%s:%s (max_workers=%s)",
DEFAULT_LISTEN_HOST,
DEFAULT_LISTEN_PORT,
SOCKS5_HOST,
SOCKS5_PORT,
MAX_WORKERS,
)
server.serve_forever()
if __name__ == "__main__":
main()

View File

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

View File

@@ -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()

View File

@@ -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")

View File

@@ -2,5 +2,6 @@ playwright>=1.53.0
python-dotenv>=1.0.1
pydantic>=2.8.2
SQLAlchemy>=2.0.32
PySocks>=1.7.1
pytest>=8.3.0
pytest-cov>=5.0.0

View File

@@ -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()

View File

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

View File

@@ -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()