Files
iaai-fast/iaai_sync_service/sync_engine.py
2026-04-24 00:20:42 +03:00

1277 lines
42 KiB
Python

from __future__ import annotations
import concurrent.futures
import logging
import re
import secrets
import string
import time
from dataclasses import dataclass
from datetime import UTC, datetime
from typing import Any
from sqlalchemy import delete, insert, select, text, update
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.orm import Session
from iaai_sync_service.client import IAAIClient, ListingVehicle, build_resizer_images_from_keys
from iaai_sync_service.models import Car, Image
from iaai_sync_service.runtime_config import RuntimeConfig, RuntimeFilters
from iaai_sync_service.settings import Settings
logger = logging.getLogger(__name__)
ORIGIN_PREFIX = "iaai:"
ORIGIN_URL_BASE = "https://www.iaai.com/VehicleDetail"
DAMAGE_NEUTRAL_VALUES = {
"NORMAL WEAR & TEAR",
"NORMAL WEAR",
"NORMALWEAR&TEAR",
"NONE",
"NO DAMAGE",
"NO VISIBLE DAMAGE",
"MINOR DENT/SCRATCHES",
}
INACTIVE_STATUS_VALUES = {
"SOLD",
"SO",
"CLOSED",
"CN",
"DELIVERED",
"WITHDRAWN",
"WDR",
"COMPLETE",
"COMPLETED",
}
PARSER_ID_PREFIX = "car-"
PARSER_ID_CHARS = string.ascii_letters + string.digits
PARSER_ID_RANDOM_LEN = 22
PARSER_ID_PATTERN = re.compile(r"^car-[A-Za-z0-9]{22}$")
DB_IN_CLAUSE_CHUNK_SIZE = 2000
@dataclass
class SyncStats:
ids_fetched: int = 0
cars_upserted: int = 0
cars_failed: int = 0
images_upserted: int = 0
@dataclass(frozen=True)
class PreparedCar:
inventory_id: str
brand: str
model: str
year: int | None
price: int | None
currency: str
mileage: int
country: str
color: str
drive: str | None
gearbox: str | None
steering_wheel: str
body_type: str
engine_volume: int | None
selling_type: str
origin: str
origin_url: str
origin_id: str
is_damaged: bool
evaluation: str | None
rental: bool
images: list[dict[str, str | int]]
class RowSkipError(RuntimeError):
pass
def utc_now() -> datetime:
return datetime.now(UTC)
def run_sync_once(
*,
session: Session,
client: IAAIClient,
settings: Settings,
runtime_config: RuntimeConfig | None = None,
) -> tuple[SyncStats, list[str]]:
runtime = runtime_config or RuntimeConfig()
stats = SyncStats()
errors: list[str] = []
seen_at = utc_now()
rows_seen = 0
rows_skipped_condition = 0
rows_selected = 0
candidates: dict[str, ListingVehicle] = {}
logger.info(
"IAAI listing collection started condition_check_enabled=%s filters_enabled=%s",
runtime.condition_check_enabled is True,
runtime.filters.is_enabled(),
)
for vehicle in client.iter_listing_vehicles(runtime.filters):
rows_seen += 1
if runtime.condition_check_enabled and not passes_condition_check(vehicle):
rows_skipped_condition += 1
continue
candidates[vehicle.inventory_id] = vehicle
rows_selected = len(candidates)
logger.info(
"IAAI listing summary rows_seen=%s rows_filtered_condition=%s condition_check_enabled=%s rows_selected=%s",
rows_seen,
rows_skipped_condition,
runtime.condition_check_enabled is True,
rows_selected,
)
if not candidates:
return stats, errors
currency_labels = get_enum_labels(session, "currencyenum")
country_labels = get_enum_labels(session, "countryenum")
drive_labels = get_enum_labels(session, "driveenum")
gearbox_labels = get_enum_labels(session, "gearboxenum")
body_type_labels = get_enum_labels(session, "bodytypeenum")
body_type_default = resolve_body_type_default(session)
steering_left = resolve_steering_left(session)
selling_type_auction = resolve_selling_type_auction(session)
origin_code = resolve_origin_code(session)
origin_ids = {f"{ORIGIN_PREFIX}{inventory_id}" for inventory_id in candidates}
origin_urls = {f"{ORIGIN_URL_BASE}/{inventory_id}" for inventory_id in candidates}
existing_cars_by_id, existing_cars_by_url = preload_existing_cars(
session=session,
origin_ids=origin_ids,
origin_urls=origin_urls,
)
existing_image_signatures = preload_image_signatures(
session=session,
car_ids=list({car.id for car in existing_cars_by_id.values()}),
)
commit_batch_size = max(1, settings.db_commit_batch_size)
started_at = time.perf_counter()
prepared_rows: list[PreparedCar] = []
with concurrent.futures.ThreadPoolExecutor(max_workers=settings.fetch_concurrency) as executor:
future_to_vehicle = {
executor.submit(client.fetch_vehicle_detail_payload, inventory_id): vehicle
for inventory_id, vehicle in candidates.items()
}
for index, future in enumerate(concurrent.futures.as_completed(future_to_vehicle), start=1):
vehicle = future_to_vehicle[future]
try:
payload = future.result()
prepared = prepare_car(
listing_vehicle=vehicle,
detail_payload=payload,
currency_labels=currency_labels,
country_labels=country_labels,
drive_labels=drive_labels,
gearbox_labels=gearbox_labels,
body_type_labels=body_type_labels,
body_type_default=body_type_default,
steering_left=steering_left,
selling_type_auction=selling_type_auction,
origin_code=origin_code,
)
except Exception as exc: # noqa: BLE001
stats.cars_failed += 1
errors.append(f"inventory_id={vehicle.inventory_id}: detail parse failed: {exc}")
logger.exception("Detail parse failed inventory_id=%s: %s", vehicle.inventory_id, exc)
continue
if not passes_runtime_filters_prepared(prepared, runtime.filters):
continue
stats.ids_fetched += 1
prepared_rows.append(prepared)
if index % 100 == 0 or index == rows_selected:
elapsed = max(0.001, time.perf_counter() - started_at)
logger.info(
"IAAI sync progress %s/%s ids_fetched=%s cars_upserted=%s cars_failed=%s images_upserted=%s queued_for_db=%s throughput=%.2f details/s",
index,
rows_selected,
stats.ids_fetched,
stats.cars_upserted,
stats.cars_failed,
stats.images_upserted,
len(prepared_rows),
index / elapsed,
)
if prepared_rows:
logger.info(
"IAAI DB apply started rows=%s mode=%s commit_batch_size=%s",
len(prepared_rows),
"postgres_bulk" if is_postgresql_session(session) else "rowwise",
commit_batch_size,
)
if is_postgresql_session(session):
apply_prepared_rows_postgres(
session=session,
prepared_rows=prepared_rows,
seen_at=seen_at,
commit_batch_size=commit_batch_size,
existing_cars_by_id=existing_cars_by_id,
existing_cars_by_url=existing_cars_by_url,
existing_image_signatures=existing_image_signatures,
stats=stats,
errors=errors,
)
else:
apply_prepared_rows_rowwise(
session=session,
prepared_rows=prepared_rows,
seen_at=seen_at,
commit_batch_size=commit_batch_size,
existing_cars_by_id=existing_cars_by_id,
existing_cars_by_url=existing_cars_by_url,
existing_image_signatures=existing_image_signatures,
stats=stats,
errors=errors,
)
client.persist_session_state()
if stats.cars_failed > 0:
logger.warning("Sold reconcile skipped due to failures cars_failed=%s", stats.cars_failed)
return stats, errors
if runtime.filters.is_enabled():
logger.info("Sold reconcile skipped because runtime filters are enabled")
return stats, errors
if runtime.condition_check_enabled is True:
logger.info("Sold reconcile skipped because condition check is enabled")
return stats, errors
if stats.ids_fetched == 0:
logger.warning("Sold reconcile skipped because no rows passed filters")
return stats, errors
reconcile_sold_cars(session=session, seen_at=seen_at)
session.commit()
return stats, errors
def is_postgresql_session(session: Session) -> bool:
bind = session.get_bind()
return bind is not None and bind.dialect.name == "postgresql"
def apply_prepared_rows_rowwise(
*,
session: Session,
prepared_rows: list[PreparedCar],
seen_at: datetime,
commit_batch_size: int,
existing_cars_by_id: dict[str, Car],
existing_cars_by_url: dict[str, Car],
existing_image_signatures: dict[int, list[tuple[int, str, str]]],
stats: SyncStats,
errors: list[str],
allow_lookup: bool = False,
) -> None:
pending_db_changes = 0
total_rows = len(prepared_rows)
started_at = time.perf_counter()
for index, prepared in enumerate(prepared_rows, start=1):
try:
with session.begin_nested():
car = upsert_car(
session=session,
row=prepared,
seen_at=seen_at,
existing_entity=existing_cars_by_id.get(prepared.origin_id)
or existing_cars_by_url.get(prepared.origin_url),
allow_lookup=allow_lookup,
)
existing_cars_by_id[prepared.origin_id] = car
existing_cars_by_url[prepared.origin_url] = car
existing_signature = existing_image_signatures.get(car.id)
images_written, normalized_images = replace_images_if_changed_cached(
session=session,
car_id=car.id,
images=prepared.images,
existing_normalized=existing_signature,
)
existing_image_signatures[car.id] = normalized_images
car.is_hidden = not bool(normalized_images)
stats.cars_upserted += 1
stats.images_upserted += images_written
pending_db_changes += 1
if pending_db_changes >= commit_batch_size:
session.commit()
pending_db_changes = 0
except Exception as exc: # noqa: BLE001
stats.cars_failed += 1
errors.append(f"inventory_id={prepared.inventory_id}: upsert failed: {exc}")
logger.exception("Upsert failed inventory_id=%s: %s", prepared.inventory_id, exc)
if index % 100 == 0 or index == total_rows:
elapsed = max(0.001, time.perf_counter() - started_at)
logger.info(
"IAAI DB progress %s/%s cars_upserted=%s cars_failed=%s images_upserted=%s throughput=%.2f rows/s",
index,
total_rows,
stats.cars_upserted,
stats.cars_failed,
stats.images_upserted,
index / elapsed,
)
if pending_db_changes > 0:
session.commit()
def apply_prepared_rows_postgres(
*,
session: Session,
prepared_rows: list[PreparedCar],
seen_at: datetime,
commit_batch_size: int,
existing_cars_by_id: dict[str, Car],
existing_cars_by_url: dict[str, Car],
existing_image_signatures: dict[int, list[tuple[int, str, str]]],
stats: SyncStats,
errors: list[str],
) -> None:
total_rows = len(prepared_rows)
started_at = time.perf_counter()
for batch_start in range(0, total_rows, commit_batch_size):
batch = prepared_rows[batch_start : batch_start + commit_batch_size]
try:
_apply_prepared_batch_postgres(
session=session,
prepared_rows=batch,
seen_at=seen_at,
existing_cars_by_id=existing_cars_by_id,
existing_cars_by_url=existing_cars_by_url,
existing_image_signatures=existing_image_signatures,
stats=stats,
)
except Exception as exc: # noqa: BLE001
session.rollback()
logger.exception(
"PostgreSQL bulk apply failed batch_start=%s batch_size=%s. Falling back to row-wise mode: %s",
batch_start + 1,
len(batch),
exc,
)
apply_prepared_rows_rowwise(
session=session,
prepared_rows=batch,
seen_at=seen_at,
commit_batch_size=max(1, min(commit_batch_size, 50)),
existing_cars_by_id=existing_cars_by_id,
existing_cars_by_url=existing_cars_by_url,
existing_image_signatures=existing_image_signatures,
stats=stats,
errors=errors,
allow_lookup=True,
)
processed = batch_start + len(batch)
if processed % 100 == 0 or processed == total_rows:
elapsed = max(0.001, time.perf_counter() - started_at)
logger.info(
"IAAI DB progress %s/%s cars_upserted=%s cars_failed=%s images_upserted=%s throughput=%.2f rows/s",
processed,
total_rows,
stats.cars_upserted,
stats.cars_failed,
stats.images_upserted,
processed / elapsed,
)
def _apply_prepared_batch_postgres(
*,
session: Session,
prepared_rows: list[PreparedCar],
seen_at: datetime,
existing_cars_by_id: dict[str, Car],
existing_cars_by_url: dict[str, Car],
existing_image_signatures: dict[int, list[tuple[int, str, str]]],
stats: SyncStats,
) -> None:
if not prepared_rows:
return
reserved_slugs: set[str] = set()
upsert_values: list[dict[str, Any]] = []
normalized_images_by_origin: dict[str, list[tuple[int, str, str]]] = {}
for row in prepared_rows:
existing = existing_cars_by_id.get(row.origin_id) or existing_cars_by_url.get(row.origin_url)
parser_id = parse_text(existing.parser_id) if existing is not None else None
if not is_valid_parser_id(parser_id):
parser_id = generate_parser_id(session)
slug_needs_refresh = (
existing is None
or not parse_text(existing.slug)
or existing.brand != row.brand
or existing.model != row.model
or existing.year != row.year
)
if slug_needs_refresh:
slug = generate_unique_slug(
session=session,
brand=row.brand,
model=row.model,
year=row.year,
parser_id=parser_id,
)
slug = reserve_slug_in_batch(
session=session,
base_slug=slug,
reserved_slugs=reserved_slugs,
)
else:
slug = str(existing.slug)
reserved_slugs.add(slug)
normalized_images = normalize_images(row.images)
normalized_images_by_origin[row.origin_id] = normalized_images
upsert_values.append(
{
"parser_id": parser_id,
"brand": row.brand,
"model": row.model,
"year": row.year,
"price": row.price,
"currency": row.currency,
"mileage": row.mileage,
"country": row.country,
"is_sold": False,
"color": row.color,
"drive": row.drive,
"gearbox": row.gearbox,
"steering_wheel": row.steering_wheel,
"body_type": row.body_type,
"engine_volume": row.engine_volume,
"selling_type": row.selling_type,
"one_owner": False,
"new_car": False,
"is_hidden": not bool(normalized_images),
"origin": row.origin,
"origin_url": row.origin_url,
"origin_id": row.origin_id,
"is_damaged": row.is_damaged,
"evaluation": row.evaluation,
"non_smoking": True,
"rental": row.rental,
"repair_history": False,
"slug": slug,
"last_seen_at": seen_at,
}
)
statement = pg_insert(Car).values(upsert_values)
excluded = statement.excluded
update_values = {
"parser_id": excluded.parser_id,
"brand": excluded.brand,
"model": excluded.model,
"year": excluded.year,
"price": excluded.price,
"currency": excluded.currency,
"mileage": excluded.mileage,
"country": excluded.country,
"is_sold": excluded.is_sold,
"color": excluded.color,
"drive": excluded.drive,
"gearbox": excluded.gearbox,
"steering_wheel": excluded.steering_wheel,
"body_type": excluded.body_type,
"engine_volume": excluded.engine_volume,
"selling_type": excluded.selling_type,
"one_owner": excluded.one_owner,
"new_car": excluded.new_car,
"is_hidden": excluded.is_hidden,
"origin": excluded.origin,
"origin_url": excluded.origin_url,
"is_damaged": excluded.is_damaged,
"evaluation": excluded.evaluation,
"non_smoking": excluded.non_smoking,
"rental": excluded.rental,
"repair_history": excluded.repair_history,
"slug": excluded.slug,
"last_seen_at": excluded.last_seen_at,
}
result_rows = session.execute(
statement.on_conflict_do_update(
index_elements=[Car.origin_id],
set_=update_values,
).returning(Car.id, Car.origin_id)
).all()
car_id_by_origin = {origin_id: car_id for car_id, origin_id in result_rows}
changed_car_ids: list[int] = []
image_rows: list[dict[str, Any]] = []
local_images_upserted = 0
signature_updates: dict[int, list[tuple[int, str, str]]] = {}
for payload in upsert_values:
origin_id = str(payload["origin_id"])
car_id = car_id_by_origin.get(origin_id)
if car_id is None:
raise RuntimeError(f"Bulk upsert did not return car id for origin_id={origin_id}")
normalized_new = normalized_images_by_origin[origin_id]
existing_normalized = existing_image_signatures.get(car_id)
if normalized_new == existing_normalized:
continue
changed_car_ids.append(car_id)
signature_updates[car_id] = normalized_new
local_images_upserted += len(normalized_new)
for order_index, fullres_image, preview_image in normalized_new:
image_rows.append(
{
"car_id": car_id,
"order_index": order_index,
"fullres_image": fullres_image,
"preview_image": preview_image,
}
)
if changed_car_ids:
session.execute(delete(Image).where(Image.car_id.in_(changed_car_ids)))
if image_rows:
session.execute(insert(Image), image_rows)
session.commit()
existing_image_signatures.update(signature_updates)
stats.cars_upserted += len(prepared_rows)
stats.images_upserted += local_images_upserted
def reserve_slug_in_batch(
*,
session: Session,
base_slug: str,
reserved_slugs: set[str],
) -> str:
if base_slug not in reserved_slugs:
reserved_slugs.add(base_slug)
return base_slug
counter = 1
while True:
candidate = f"{base_slug}-{counter}"
counter += 1
if candidate in reserved_slugs:
continue
exists = session.execute(select(Car.id).where(Car.slug == candidate)).scalar_one_or_none()
if exists is None:
reserved_slugs.add(candidate)
return candidate
def passes_condition_check(row: ListingVehicle) -> bool:
if row.timed_auction_closed:
return False
status = (row.inventory_status or "").strip().upper()
if status in INACTIVE_STATUS_VALUES:
return False
return True
def passes_runtime_filters_prepared(row: PreparedCar, filters: RuntimeFilters) -> bool:
if not filters.is_enabled():
return True
brand_key = row.brand.casefold()
model_key = row.model.casefold()
body_key = row.body_type.casefold()
if filters.brands and brand_key not in filters.brands:
return False
if filters.models and model_key not in filters.models:
return False
if filters.years and (row.year is None or row.year not in filters.years):
return False
if filters.body_types and body_key not in filters.body_types:
return False
return True
def prepare_car(
*,
listing_vehicle: ListingVehicle,
detail_payload: dict[str, Any],
currency_labels: set[str],
country_labels: set[str],
drive_labels: set[str],
gearbox_labels: set[str],
body_type_labels: set[str],
body_type_default: str,
steering_left: str,
selling_type_auction: str,
origin_code: str,
) -> PreparedCar:
inventory_view = detail_payload.get("inventoryView")
if not isinstance(inventory_view, dict):
raise RowSkipError("detail payload missing inventoryView")
attributes = inventory_view.get("attributes")
if not isinstance(attributes, dict):
raise RowSkipError("detail payload missing inventoryView.attributes")
inventory_id = parse_text(attributes.get("Id")) or listing_vehicle.inventory_id
if not inventory_id:
raise RowSkipError("missing inventory id")
brand = limit_text(parse_text(attributes.get("Make")) or "unknown", 50)
model = build_model_name(
parse_text(attributes.get("Model")),
parse_text(attributes.get("Series")),
)
if not model:
model = "unknown"
model = limit_text(model, 50)
year = parse_year(attributes.get("Year"))
auction_info = detail_payload.get("auctionInformation")
auction_info = auction_info if isinstance(auction_info, dict) else {}
bidding_info = auction_info.get("biddingInformation")
bidding_info = bidding_info if isinstance(bidding_info, dict) else {}
prebid_info = auction_info.get("prebidInformation")
prebid_info = prebid_info if isinstance(prebid_info, dict) else {}
high_bid = first_positive_int(
prebid_info.get("decimalHighBidAmount"),
prebid_info.get("highBidAmount"),
bidding_info.get("highBidAmount"),
)
buy_now = first_positive_int(
bidding_info.get("buyNowAmount"),
prebid_info.get("buyNowPrice"),
bidding_info.get("buyNowPrice"),
)
price = high_bid if high_bid is not None else buy_now
currency = map_currency(
parse_text(attributes.get("Currency")) or listing_vehicle.currency,
currency_labels=currency_labels,
)
country = map_country(
inventory_id=inventory_id,
country_labels=country_labels,
)
mileage = parse_non_negative_int(attributes.get("ODOValue")) or 0
color = normalize_color(
parse_text(attributes.get("ExteriorColor"))
or parse_text(attributes.get("ColorDesc"))
or parse_text(attributes.get("colorDesc"))
)
drive = map_drive(parse_text(attributes.get("DriveLineTypeDesc")), drive_labels=drive_labels)
gearbox = map_gearbox(parse_text(attributes.get("Transmission")), gearbox_labels=gearbox_labels)
body_type = (
map_body_type(
parse_text(attributes.get("BodyStyleName")) or parse_text(attributes.get("VehicleClass")),
body_type_labels=body_type_labels,
)
or body_type_default
)
engine_volume = parse_engine_volume(
parse_text(attributes.get("EngineSize")) or parse_text(attributes.get("EngineInformation"))
)
evaluation = parse_text(attributes.get("VehicleGrade"))
is_damaged = derive_is_damaged(
primary_damage=parse_text(attributes.get("PrimaryDamageDesc")),
secondary_damage=parse_text(attributes.get("SecondaryDamageDesc")),
)
image_dimensions = inventory_view.get("imageDimensions")
image_dimensions = image_dimensions if isinstance(image_dimensions, dict) else {}
keys_container = image_dimensions.get("keys")
keys_container = keys_container if isinstance(keys_container, dict) else {}
image_keys = keys_container.get("$values")
image_keys = image_keys if isinstance(image_keys, list) else []
images = build_resizer_images_from_keys(image_keys)
origin_id = f"{ORIGIN_PREFIX}{inventory_id}"
origin_url = f"{ORIGIN_URL_BASE}/{inventory_id}"
return PreparedCar(
inventory_id=inventory_id,
brand=brand,
model=model,
year=year,
price=price,
currency=currency,
mileage=mileage,
country=country,
color=color,
drive=drive,
gearbox=gearbox,
steering_wheel=steering_left,
body_type=body_type,
engine_volume=engine_volume,
selling_type=selling_type_auction,
origin=origin_code,
origin_url=origin_url,
origin_id=origin_id,
is_damaged=is_damaged,
evaluation=evaluation,
rental=False,
images=images,
)
def preload_existing_cars(
*,
session: Session,
origin_ids: set[str],
origin_urls: set[str],
) -> tuple[dict[str, Car], dict[str, Car]]:
by_origin_id: dict[str, Car] = {}
by_origin_url: dict[str, Car] = {}
if origin_ids:
ids = sorted(origin_ids)
for chunk in chunked(ids, DB_IN_CLAUSE_CHUNK_SIZE):
rows = session.execute(select(Car).where(Car.origin_id.in_(chunk))).scalars().all()
for row in rows:
by_origin_id[row.origin_id] = row
by_origin_url[row.origin_url] = row
if origin_urls:
urls_to_lookup = sorted(url for url in origin_urls if url not in by_origin_url)
for chunk in chunked(urls_to_lookup, DB_IN_CLAUSE_CHUNK_SIZE):
rows = session.execute(select(Car).where(Car.origin_url.in_(chunk))).scalars().all()
for row in rows:
by_origin_url[row.origin_url] = row
by_origin_id.setdefault(row.origin_id, row)
return by_origin_id, by_origin_url
def preload_image_signatures(*, session: Session, car_ids: list[int]) -> dict[int, list[tuple[int, str, str]]]:
if not car_ids:
return {}
signatures: dict[int, list[tuple[int, str, str]]] = {}
for chunk in chunked(car_ids, DB_IN_CLAUSE_CHUNK_SIZE):
rows = session.execute(
select(Image)
.where(Image.car_id.in_(chunk))
.order_by(Image.car_id, Image.order_index, Image.fullres_image, Image.preview_image)
).scalars()
for row in rows:
signatures.setdefault(row.car_id, []).append(
(
row.order_index,
row.fullres_image.strip(),
row.preview_image.strip(),
)
)
return signatures
def chunked(values: list[Any], size: int) -> list[list[Any]]:
if size <= 0:
return [values]
return [values[i : i + size] for i in range(0, len(values), size)]
def upsert_car(
*,
session: Session,
row: PreparedCar,
seen_at: datetime,
existing_entity: Car | None = None,
allow_lookup: bool = True,
) -> Car:
entity = existing_entity
if entity is None and allow_lookup:
entity = session.execute(select(Car).where(Car.origin_id == row.origin_id)).scalar_one_or_none()
if entity is None and allow_lookup:
entity = session.execute(select(Car).where(Car.origin_url == row.origin_url)).scalar_one_or_none()
is_new_entity = False
if entity is None:
entity = Car(
parser_id=generate_parser_id(session),
brand=row.brand,
model=row.model,
origin_id=row.origin_id,
origin_url=row.origin_url,
)
is_new_entity = True
elif not is_valid_parser_id(entity.parser_id):
entity.parser_id = generate_parser_id(session)
slug_needs_refresh = (
is_new_entity
or not parse_text(entity.slug)
or entity.brand != row.brand
or entity.model != row.model
or entity.year != row.year
)
if slug_needs_refresh:
slug = generate_unique_slug(
session=session,
brand=row.brand,
model=row.model,
year=row.year,
parser_id=entity.parser_id,
)
else:
slug = entity.slug
entity.brand = row.brand
entity.model = row.model
entity.year = row.year
entity.price = row.price
entity.currency = row.currency
entity.mileage = row.mileage
entity.country = row.country
entity.is_sold = False
entity.color = row.color
entity.drive = row.drive
entity.gearbox = row.gearbox
entity.steering_wheel = row.steering_wheel
entity.body_type = row.body_type
entity.engine_volume = row.engine_volume
entity.selling_type = row.selling_type
entity.one_owner = False
entity.new_car = False
entity.origin = row.origin
entity.origin_url = row.origin_url
entity.origin_id = row.origin_id
entity.is_damaged = row.is_damaged
entity.evaluation = row.evaluation
entity.non_smoking = True
entity.rental = row.rental
entity.repair_history = False
entity.slug = slug
entity.last_seen_at = seen_at
if is_new_entity:
session.add(entity)
session.flush()
return entity
def reconcile_sold_cars(*, session: Session, seen_at: datetime) -> None:
session.execute(
update(Car)
.where(
Car.origin_id.like(f"{ORIGIN_PREFIX}%"),
Car.last_seen_at < seen_at,
Car.is_sold.is_(False),
)
.values(is_sold=True)
)
def replace_images_if_changed(
*,
session: Session,
car_id: int,
images: list[dict[str, str | int]],
) -> int:
written, _ = replace_images_if_changed_cached(
session=session,
car_id=car_id,
images=images,
existing_normalized=None,
)
return written
def replace_images_if_changed_cached(
*,
session: Session,
car_id: int,
images: list[dict[str, str | int]],
existing_normalized: list[tuple[int, str, str]] | None,
) -> tuple[int, list[tuple[int, str, str]]]:
normalized_new = normalize_images(images)
if existing_normalized is None:
existing_rows = session.execute(select(Image).where(Image.car_id == car_id)).scalars().all()
existing_normalized = sorted(
((row.order_index, row.fullres_image.strip(), row.preview_image.strip()) for row in existing_rows),
key=lambda item: (item[0], item[1], item[2]),
)
if normalized_new == existing_normalized:
return 0, existing_normalized
replace_images(session=session, car_id=car_id, images=images)
return len(normalized_new), normalized_new
def replace_images(
*,
session: Session,
car_id: int,
images: list[dict[str, str | int]],
) -> None:
session.execute(delete(Image).where(Image.car_id == car_id))
for image in images:
fullres_image = parse_text(image.get("fullres_image"))
if fullres_image is None:
continue
preview_image = parse_text(image.get("preview_image")) or fullres_image
order_index = parse_non_negative_int(image.get("order_index")) or 0
session.add(
Image(
car_id=car_id,
fullres_image=fullres_image,
preview_image=preview_image,
order_index=order_index,
)
)
def normalize_images(images: list[dict[str, str | int]]) -> list[tuple[int, str, str]]:
normalized: list[tuple[int, str, str]] = []
for image in images:
fullres_image = parse_text(image.get("fullres_image"))
if fullres_image is None:
continue
preview_image = parse_text(image.get("preview_image")) or fullres_image
order_index = parse_non_negative_int(image.get("order_index")) or 0
normalized.append((order_index, fullres_image.strip(), preview_image.strip()))
normalized.sort(key=lambda item: (item[0], item[1], item[2]))
return normalized
def build_model_name(model: str | None, series: str | None) -> str:
values = [model, series]
unique_parts: list[str] = []
seen: set[str] = set()
for value in values:
if not value:
continue
normalized = " ".join(value.split())
if not normalized:
continue
key = normalized.casefold()
if key in seen:
continue
seen.add(key)
unique_parts.append(normalized)
return " ".join(unique_parts).strip()
def parse_year(value: Any) -> int | None:
parsed = parse_non_negative_int(value)
if parsed is None or parsed <= 0:
return None
if parsed < 1900 or parsed > 2100:
return None
return parsed
def parse_non_negative_int(value: Any) -> int | None:
if value is None or isinstance(value, bool):
return None
if isinstance(value, int):
return value if value >= 0 else None
if isinstance(value, float):
if value < 0:
return None
return int(round(value))
if isinstance(value, str):
raw = value.strip()
if not raw:
return None
compact = raw.replace(",", "").replace(" ", "")
compact = compact.replace("$", "")
match = re.search(r"-?\d+(?:\.\d+)?", compact)
if match is None:
return None
try:
parsed = float(match.group(0))
except ValueError:
return None
if parsed < 0:
return None
return int(round(parsed))
return None
def first_positive_int(*values: Any) -> int | None:
for value in values:
parsed = parse_non_negative_int(value)
if parsed is not None and parsed > 0:
return parsed
return None
def parse_text(value: Any) -> str | None:
if isinstance(value, str):
text = value.strip()
return text if text else None
return None
def normalize_color(value: str | None) -> str:
if not value:
return "other"
normalized = value.strip().lower()
return normalized or "other"
def map_currency(value: str | None, *, currency_labels: set[str]) -> str:
if value is None:
return select_enum_label(currency_labels, ("USD", "NA"), fallback="NA") or "NA"
normalized = value.strip().upper()
if normalized == "CAD":
return select_enum_label(currency_labels, ("CAD", "USD", "NA"), fallback="NA") or "NA"
if normalized == "USD":
return select_enum_label(currency_labels, ("USD", "NA"), fallback="NA") or "NA"
return select_enum_label(currency_labels, ("NA", "USD", "CAD"), fallback="NA") or "NA"
def map_country(*, inventory_id: str, country_labels: set[str]) -> str:
upper_id = inventory_id.strip().upper()
if upper_id.endswith("~US"):
return select_enum_label(country_labels, ("US", "NA"), fallback="NA") or "NA"
if upper_id.endswith("~CA"):
return select_enum_label(country_labels, ("CA", "US", "NA"), fallback="NA") or "NA"
return select_enum_label(country_labels, ("NA", "US", "CA"), fallback="NA") or "NA"
def map_drive(value: str | None, *, drive_labels: set[str]) -> str | None:
if not value:
return None
normalized = value.strip().lower()
candidates: tuple[str, ...] | None = None
if "front" in normalized or normalized == "fwd":
candidates = ("FWD",)
elif "all wheel" in normalized or "4x4" in normalized or normalized == "awd":
candidates = ("4WD", "FOUR_WD")
elif "rear" in normalized or normalized == "rwd":
candidates = ("RWD",)
elif "2wd" in normalized or "two wheel" in normalized:
candidates = ("2WD", "TWO_WD")
elif "unknown" in normalized or normalized in {"na", "n/a"}:
candidates = ("NA",)
if candidates is None:
return None
return select_enum_label(drive_labels, candidates, fallback="NA")
def map_gearbox(value: str | None, *, gearbox_labels: set[str]) -> str | None:
if not value:
return None
normalized = value.strip().lower()
candidates: tuple[str, ...] | None = None
if "cvt" in normalized:
candidates = ("CVT",)
elif "manual" in normalized or normalized == "mt":
candidates = ("MT",)
elif "electric" in normalized or normalized == "ev":
candidates = ("EV",)
elif "auto" in normalized or normalized == "at":
candidates = ("AT",)
elif "unknown" in normalized or normalized in {"na", "n/a"}:
candidates = ("NA",)
if candidates is None:
return None
return select_enum_label(gearbox_labels, candidates, fallback="NA")
def map_body_type(value: str | None, *, body_type_labels: set[str]) -> str | None:
if not value:
return None
normalized = value.strip().lower()
candidates: tuple[str, ...] | None = None
if "sedan" in normalized:
candidates = ("SEDAN",)
elif "sport utility" in normalized or normalized == "suv":
candidates = ("SUV",)
elif "hatch" in normalized:
candidates = ("HATCHBACK",)
elif "wagon" in normalized or normalized == "station":
candidates = ("STATION_WAGON", "Station Wagon")
elif "coupe" in normalized:
candidates = ("COUPE",)
elif "pickup" in normalized or "crew" in normalized and "cab" in normalized:
candidates = ("PICKUP", "Pickup")
elif "convertible" in normalized or "roadster" in normalized or "cabrio" in normalized:
candidates = ("OPEN", "Open")
elif "van" in normalized:
candidates = ("MINIVAN",)
elif "truck" in normalized or "chassis" in normalized:
candidates = ("TRUCK", "Truck")
elif "rv" in normalized or "motorized" in normalized:
candidates = ("RV",)
elif normalized in {"other", "unknown"}:
candidates = ("OTHER", "Other")
if candidates is None:
return None
return select_enum_label(body_type_labels, candidates, fallback="OTHER")
def parse_engine_volume(value: str | None) -> int | None:
if not value:
return None
match = re.search(r"(\d+(?:\.\d+)?)\s*[lL]\b", value)
if not match:
return None
try:
liters = float(match.group(1))
except ValueError:
return None
cc = int(round(liters * 1000))
if cc <= 0 or cc > 10000:
return None
return cc
def derive_is_damaged(*, primary_damage: str | None, secondary_damage: str | None) -> bool:
for value in (primary_damage, secondary_damage):
if not value:
continue
normalized = value.replace(" ", "").strip().upper()
if not normalized:
continue
if normalized in {item.replace(" ", "") for item in DAMAGE_NEUTRAL_VALUES}:
continue
return True
return False
def resolve_origin_code(session: Session) -> str:
labels = get_enum_labels(session, "originenum")
if labels:
return select_enum_label(labels, ("IAAI", "iaai", "NA"), fallback="NA") or "NA"
return "IAAI"
def resolve_selling_type_auction(session: Session) -> str:
labels = get_enum_labels(session, "sellingtypeenum")
if labels:
return select_enum_label(labels, ("AUCTION", "auction", "NA"), fallback="NA") or "NA"
return "AUCTION"
def resolve_steering_left(session: Session) -> str:
labels = get_enum_labels(session, "steeringwheelenum")
if labels:
return select_enum_label(labels, ("LEFT", "left", "NA"), fallback="NA") or "NA"
return "LEFT"
def resolve_body_type_default(session: Session) -> str:
labels = get_enum_labels(session, "bodytypeenum")
if labels:
return select_enum_label(labels, ("OTHER", "Other", "NA"), fallback="NA") or "NA"
return "OTHER"
def get_enum_labels(session: Session, type_name: str) -> set[str]:
bind = session.get_bind()
if bind is None or bind.dialect.name != "postgresql":
return set()
rows = session.execute(
text(
"""
SELECT e.enumlabel
FROM pg_type t
JOIN pg_namespace n ON n.oid = t.typnamespace
JOIN pg_enum e ON e.enumtypid = t.oid
WHERE n.nspname = 'public'
AND t.typname = :type_name
"""
),
{"type_name": type_name},
).scalars()
return {value for value in rows}
def select_enum_label(
labels: set[str] | None,
candidates: tuple[str, ...],
*,
fallback: str | None = None,
) -> str | None:
if not labels:
return candidates[0] if candidates else fallback
for candidate in candidates:
if candidate in labels:
return candidate
if fallback and fallback in labels:
return fallback
return sorted(labels)[0] if labels else fallback
def limit_text(value: str, max_length: int) -> str:
if len(value) <= max_length:
return value
return value[:max_length].rstrip()
def is_valid_parser_id(value: str | None) -> bool:
if value is None:
return False
return bool(PARSER_ID_PATTERN.fullmatch(value))
def generate_parser_id(session: Session) -> str:
while True:
candidate = PARSER_ID_PREFIX + "".join(secrets.choice(PARSER_ID_CHARS) for _ in range(PARSER_ID_RANDOM_LEN))
exists = session.execute(select(Car.id).where(Car.parser_id == candidate)).scalar_one_or_none()
if exists is None:
return candidate
def slugify_text(value: str) -> str:
lowered = value.lower()
slug = re.sub(r"[^a-z0-9]+", "-", lowered).strip("-")
return slug or "car"
def generate_unique_slug(
*,
session: Session,
brand: str,
model: str,
year: int | None,
parser_id: str,
) -> str:
if year is None:
base_slug = slugify_text(f"{brand}-{model}")
else:
base_slug = slugify_text(f"{brand}-{model}-{year}")
slug = base_slug
counter = 1
while True:
existing = session.execute(
select(Car.id).where(
Car.slug == slug,
Car.parser_id != parser_id,
)
).scalar_one_or_none()
if existing is None:
return slug
slug = f"{base_slug}-{counter}"
counter += 1