Files
mobile.de/mobilede_scraper/storage/db.py
T
2026-08-19 11:18:00 +03:00

1188 lines
39 KiB
Python

import logging
import os
import re
import uuid
from contextlib import contextmanager
from datetime import datetime, timedelta, timezone
from typing import Any, Iterator
from sqlalchemy import case, create_engine, delete, func, or_, select, text, update
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.orm import Session, sessionmaker
from ..core.config import Settings
from .models import Base, Car, DetailJob, Image, SyncRun
from .schemas import CarRecord
logger = logging.getLogger("MOBILEDE_scraper.db")
MOBILEDE_SKIP_IMAGES_FOR_UPDATED = os.getenv("MOBILEDE_SKIP_IMAGES_FOR_UPDATED", "1").strip().lower() in {"1", "true", "yes", "on"}
CAR_DB_FIELDS = {
col.key for col in Car.__table__.columns
if col.key not in ("id",)
}
CAR_UPDATE_FIELDS = CAR_DB_FIELDS - {"first_seen_at"}
# Поля, изменение которых видно уже в Search и требует повторного Detail.
# Технические timestamps и sold-флаги сюда намеренно не входят.
MOBILEDE_MATERIAL_CHANGE_FIELDS = (
"brand",
"model",
"year",
"price",
"currency",
"mileage",
"country",
"color",
"drive",
"gearbox",
"steering_wheel",
"body_type",
"engine_volume",
"selling_type",
"one_owner",
"new_car",
"is_hidden",
"origin_url",
"is_damaged",
"evaluation",
"slug",
)
_IN_CHUNK_SIZE = 5000
MOBILEDE_ORIGIN_PREFIXES = ("mobile.de:", "mobilede:")
PARSER_ID_RE = re.compile(r"^car-[A-Za-z0-9]{22}$")
def _utc_datetime(value: datetime | None) -> datetime | None:
if value is None:
return None
if value.tzinfo is None:
return value.replace(tzinfo=timezone.utc)
return value.astimezone(timezone.utc)
DETAIL_VALUE_FALLBACKS: dict[str, set[object]] = {
"brand": {None, "", "UNKNOWN"},
"model": {None, "", "UNKNOWN"},
"year": {None},
"price": {None},
"mileage": {None, 0},
"country": {None, "", "NA"},
"color": {None, "", "other"},
"drive": {None, "NA"},
"gearbox": {None, "NA"},
"body_type": {None, "OTHER", "NA"},
"engine_volume": {None},
"evaluation": {None, ""},
}
def _origin_prefix_filter(column, prefixes: tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES):
"""Match all supported mobile.de origin_id prefixes.
Older records were stored as ``mobile.de:<id>`` while some newer helper code
used ``mobilede:<id>``. Cleanup and refresh queries must include both to avoid
leaving stale active cars in the DB.
"""
return or_(*[column.like(f"{prefix}%") for prefix in prefixes])
class PersistenceService:
def __init__(self, settings: Settings) -> None:
self.settings = settings
engine_kwargs = {
"echo": settings.database.echo,
"future": True,
}
if "postgresql" in settings.database.url:
engine_kwargs["pool_size"] = settings.database.pool_size
engine_kwargs["max_overflow"] = settings.database.max_overflow
engine_kwargs["pool_pre_ping"] = True
engine_kwargs["pool_recycle"] = settings.database.pool_recycle_seconds
engine_kwargs["pool_timeout"] = 30
# Таймауты запросов и блокировок.
engine_kwargs["connect_args"] = {
"options": "-c statement_timeout=120000 -c lock_timeout=30000"
}
self.engine = create_engine(settings.database.url, **engine_kwargs)
self.session_factory = sessionmaker(bind=self.engine, expire_on_commit=False, future=True)
def create_tables(self) -> None:
# Для SQLite можно create_all.
is_sqlite = self.settings.database.url.startswith("sqlite")
if not is_sqlite and not self.settings.database.auto_create_tables:
return
try:
Base.metadata.create_all(self.engine)
except Exception:
logger.debug("create_tables skipped (schema already exists)")
@contextmanager
def session_scope(self) -> Iterator[Session]:
session = self.session_factory()
try:
yield session
session.commit()
except Exception:
session.rollback()
raise
finally:
session.close()
def start_sync_run(self, lane: str) -> int:
with self.session_scope() as session:
now = datetime.now(timezone.utc)
stale_seconds = max(3600, int(os.getenv("MOBILEDE_SYNC_RUN_STALE_SECONDS", "90000")))
stale_before = now - timedelta(seconds=stale_seconds)
stale_runs = session.execute(
select(SyncRun).where(
SyncRun.status == "running",
SyncRun.started_at < stale_before,
)
).scalars().all()
for stale in stale_runs:
stale.status = "failed"
stale.finished_at = now
if not stale.error_summary:
stale.error_summary = "Recovered stale running sync run before starting a new run"
run = SyncRun(status="running", lane=lane, ids_fetched=0, cars_upserted=0, cars_failed=0, images_upserted=0)
session.add(run)
session.flush()
return int(run.id)
def finish_sync_run(self, run_id: int, *, status: str, ids_fetched: int, cars_upserted: int, cars_failed: int, images_upserted: int, error_summary: str | None = None) -> None:
with self.session_scope() as session:
run = session.get(SyncRun, run_id)
if run is None:
return
run.finished_at = datetime.now(timezone.utc)
run.status = status
run.ids_fetched = ids_fetched
run.cars_upserted = cars_upserted
run.cars_failed = cars_failed
run.images_upserted = images_upserted
run.error_summary = error_summary
def get_existing_origin_ids(self, origin_ids: list[str]) -> set[str]:
if not origin_ids:
return set()
with self.session_scope() as session:
rows = session.execute(select(Car.origin_id).where(Car.origin_id.in_(origin_ids))).all()
return {str(row[0]) for row in rows if row and row[0]}
@staticmethod
def _add_images(session: Session, car_id: int, images: list[dict[str, object]]) -> None:
if not images:
return
session.add_all([
Image(
fullres_image=str(img["fullres_image"]),
preview_image=str(img["preview_image"]),
order_index=int(img.get("order_index", 0)),
car_id=car_id,
)
for img in images
])
@staticmethod
def _car_payload(record: CarRecord) -> dict[str, object]:
payload = record.model_dump(mode="python")
if record.details_confirmed:
payload["details_fetched_at"] = datetime.now(timezone.utc)
return {key: value for key, value in payload.items() if key in CAR_DB_FIELDS}
@staticmethod
def _prepare_update_payload(
car: Car,
payload: dict[str, object],
record: CarRecord,
) -> dict[str, object]:
prepared = dict(payload)
if prepared.get("gallery_fetched_at") is None:
prepared["gallery_fetched_at"] = car.gallery_fetched_at
if not record.details_confirmed:
prepared["details_fetched_at"] = car.details_fetched_at
if record.preserve_existing_details and not record.details_confirmed:
for key, fallback_values in DETAIL_VALUE_FALLBACKS.items():
if prepared.get(key) in fallback_values:
prepared[key] = getattr(car, key)
return prepared
@staticmethod
def _apply_update_payload(car: Car, payload: dict[str, object]) -> None:
for key, value in payload.items():
if key in CAR_UPDATE_FIELDS:
if key == "parser_id" and PARSER_ID_RE.fullmatch(str(car.parser_id or "")):
continue
setattr(car, key, value)
@staticmethod
def _material_changed_fields(car: Car, payload: dict[str, object]) -> tuple[str, ...]:
return tuple(
field
for field in MOBILEDE_MATERIAL_CHANGE_FIELDS
if field in payload and getattr(car, field) != payload[field]
)
@classmethod
def _classify_search_update(
cls,
car: Car,
payload: dict[str, object],
record: CarRecord,
) -> tuple[tuple[str, ...], bool]:
if record.details_confirmed:
return (), False
changed_fields = cls._material_changed_fields(car, payload)
reappeared = bool(car.is_sold and payload.get("is_sold") is False)
return changed_fields, reappeared
@staticmethod
def _skip_image_sync(record: CarRecord) -> bool:
return bool(getattr(record, "skip_image_sync", False))
@classmethod
def _should_sync_images(
cls,
record: CarRecord,
images: list[dict[str, object]],
existing_urls: set[str],
) -> bool:
if not images or cls._skip_image_sync(record):
return False
if record.images_confirmed:
return True
return not existing_urls
@staticmethod
def _incoming_image_urls(images: list[dict[str, object]]) -> set[str]:
return {
str(image.get("fullres_image") or "")
for image in images
if image.get("fullres_image")
}
def _is_postgres(self) -> bool:
return self.engine.dialect.name == "postgresql"
def _load_existing_cars(
self,
session: Session,
origin_ids: list[str],
origin_urls: list[str],
) -> tuple[dict[str, Car], dict[str, Car]]:
existing_by_id: dict[str, Car] = {}
existing_by_url: dict[str, Car] = {}
if not origin_ids and not origin_urls:
return existing_by_id, existing_by_url
existing_cars: list[Car] = []
max_len = max(len(origin_ids), len(origin_urls), 1)
for i in range(0, max_len, _IN_CHUNK_SIZE):
id_chunk = origin_ids[i:i + _IN_CHUNK_SIZE]
url_chunk = origin_urls[i:i + _IN_CHUNK_SIZE]
conditions = []
if id_chunk:
conditions.append(Car.origin_id.in_(id_chunk))
if url_chunk:
conditions.append(Car.origin_url.in_(url_chunk))
if not conditions:
continue
rows = session.execute(
select(Car).where(or_(*conditions))
).scalars().all()
existing_cars.extend(rows)
for car in existing_cars:
if car.origin_id:
existing_by_id[car.origin_id] = car
if car.origin_url:
existing_by_url[car.origin_url] = car
return existing_by_id, existing_by_url
def _load_existing_image_urls(self, session: Session, car_ids: set[int]) -> dict[int, set[str]]:
existing_images_map: dict[int, set[str]] = {}
if not car_ids:
return existing_images_map
car_id_list = list(car_ids)
for i in range(0, len(car_id_list), _IN_CHUNK_SIZE):
chunk = car_id_list[i:i + _IN_CHUNK_SIZE]
img_rows = session.execute(
select(Image.car_id, Image.fullres_image).where(Image.car_id.in_(chunk))
).all()
for cid, furl in img_rows:
existing_images_map.setdefault(int(cid), set()).add(str(furl))
return existing_images_map
@staticmethod
def _postgres_upsert_set_map(insert_stmt) -> dict[str, object]:
return {
key: getattr(insert_stmt.excluded, key)
for key in CAR_UPDATE_FIELDS
}
def _replace_images_for_car(
self,
session: Session,
car_id: int,
images: list[dict[str, object]],
origin_id: str,
) -> int:
if not images:
return 0
nested = session.begin_nested()
try:
session.execute(delete(Image).where(Image.car_id == car_id))
self._add_images(session, car_id, images)
session.flush()
nested.commit()
return len(images)
except Exception:
nested.rollback()
logger.warning("Image replacement failed for car %s, keeping old images", origin_id, exc_info=True)
return 0
def _upsert_cars_batch_postgres(self, records: list[CarRecord]) -> dict[str, object]:
inserted = 0
updated = 0
changed = 0
unchanged = 0
reappeared = 0
images_total = 0
inserted_origin_ids: list[str] = []
changed_origin_ids: list[str] = []
reappeared_origin_ids: list[str] = []
with self.session_scope() as session:
origin_ids = [r.origin_id for r in records if r.origin_id]
origin_urls = [r.origin_url for r in records if r.origin_url]
existing_by_id, existing_by_url = self._load_existing_cars(session, origin_ids, origin_urls)
entries: list[dict[str, object]] = []
upsert_payloads: list[dict[str, object]] = []
for record in records:
payload = self._car_payload(record)
images = [image.model_dump(mode="python") for image in record.images]
car_by_id = existing_by_id.get(record.origin_id)
car_by_url = existing_by_url.get(record.origin_url)
existing_car = car_by_id or car_by_url
if existing_car is not None:
payload = self._prepare_update_payload(existing_car, payload, record)
changed_fields, was_reappeared = self._classify_search_update(existing_car, payload, record)
if was_reappeared:
reappeared += 1
reappeared_origin_ids.append(record.origin_id)
elif changed_fields:
changed += 1
changed_origin_ids.append(record.origin_id)
else:
unchanged += 1
if car_by_id is not None and PARSER_ID_RE.fullmatch(str(car_by_id.parser_id or "")):
payload["parser_id"] = car_by_id.parser_id
entry: dict[str, object] = {
"record": record,
"images": images,
"car_id": None,
"action": "inserted",
"skip_image_sync": self._skip_image_sync(record),
}
if car_by_url is not None and car_by_url.origin_id != record.origin_id and car_by_id is None:
self._apply_update_payload(car_by_url, payload)
entry["car_id"] = int(car_by_url.id)
entry["action"] = "updated"
updated += 1
else:
upsert_payloads.append(payload)
if car_by_id is not None:
entry["action"] = "updated"
updated += 1
else:
inserted += 1
inserted_origin_ids.append(record.origin_id)
entries.append(entry)
session.flush()
if upsert_payloads:
insert_stmt = pg_insert(Car).values(upsert_payloads)
upsert_stmt = insert_stmt.on_conflict_do_update(
index_elements=[Car.origin_id],
set_=self._postgres_upsert_set_map(insert_stmt),
).returning(Car.id, Car.origin_id)
rows = session.execute(upsert_stmt).all()
car_ids_by_origin_id = {str(origin_id): int(car_id) for car_id, origin_id in rows}
for entry in entries:
if entry["car_id"] is not None:
continue
record = entry["record"]
car_id = car_ids_by_origin_id.get(record.origin_id)
if car_id is None:
raise RuntimeError(f"PostgreSQL upsert did not return car_id for {record.origin_id}")
entry["car_id"] = car_id
insert_images_by_car_id: dict[int, list[dict[str, object]]] = {}
updated_entries = [
entry for entry in entries
if str(entry.get("action") or "updated") == "updated"
]
existing_images_map = self._load_existing_image_urls(
session,
{int(entry["car_id"]) for entry in updated_entries},
)
updated_entries_needing_compare: list[dict[str, object]] = []
for entry in entries:
car_id = int(entry["car_id"])
action = str(entry.get("action") or "updated")
images = entry["images"]
skip_image_sync = bool(entry.get("skip_image_sync"))
if action == "inserted":
insert_images_by_car_id[car_id] = images
continue
if skip_image_sync:
continue
existing_urls = existing_images_map.get(car_id, set())
if not self._should_sync_images(entry["record"], images, existing_urls):
continue
if MOBILEDE_SKIP_IMAGES_FOR_UPDATED and existing_urls and not entry["record"].images_confirmed:
continue
updated_entries_needing_compare.append(entry)
for car_id, images in insert_images_by_car_id.items():
self._add_images(session, car_id, images)
images_total += len(images)
if updated_entries_needing_compare:
images_by_car_id: dict[int, list[dict[str, object]]] = {}
replace_ids: list[int] = []
for entry in updated_entries_needing_compare:
car_id = int(entry["car_id"])
images = entry["images"]
new_image_urls = {
str(img.get("fullres_image", ""))
for img in images
if img.get("fullres_image")
}
old_image_urls = existing_images_map.get(car_id, set())
if new_image_urls != old_image_urls:
replace_ids.append(car_id)
images_by_car_id[car_id] = images
else:
images_total += len(old_image_urls)
if replace_ids:
for i in range(0, len(replace_ids), _IN_CHUNK_SIZE):
chunk = replace_ids[i:i + _IN_CHUNK_SIZE]
session.execute(delete(Image).where(Image.car_id.in_(chunk)))
for car_id in replace_ids:
images = images_by_car_id[car_id]
self._add_images(session, car_id, images)
images_total += len(images)
return {
"inserted": inserted,
"updated": updated,
"changed": changed,
"unchanged": unchanged,
"reappeared": reappeared,
"images_upserted": images_total,
"inserted_origin_ids": inserted_origin_ids,
"changed_origin_ids": changed_origin_ids,
"reappeared_origin_ids": reappeared_origin_ids,
}
def upsert_car(self, record: CarRecord):
# Вставка или обновление авто.
payload = self._car_payload(record)
images = [image.model_dump(mode="python") for image in record.images]
with self.session_scope() as session:
if self._is_postgres():
car_by_id = session.execute(
select(Car).where(Car.origin_id == record.origin_id)
).scalar_one_or_none()
if car_by_id is not None:
payload = self._prepare_update_payload(car_by_id, payload, record)
if car_by_id is not None and PARSER_ID_RE.fullmatch(str(car_by_id.parser_id or "")):
payload["parser_id"] = car_by_id.parser_id
car_by_url = session.execute(
select(Car).where(Car.origin_url == record.origin_url)
).scalar_one_or_none()
if car_by_id is None and car_by_url is not None:
payload = self._prepare_update_payload(car_by_url, payload, record)
if car_by_url is not None and car_by_url.origin_id != record.origin_id:
changed_fields, reappeared = self._classify_search_update(car_by_url, payload, record)
self._apply_update_payload(car_by_url, payload)
session.flush()
car_id = int(car_by_url.id)
action = "updated"
else:
existed = car_by_id is not None
changed_fields, reappeared = (
self._classify_search_update(car_by_id, payload, record)
if car_by_id is not None
else ((), False)
)
insert_stmt = pg_insert(Car).values(**payload)
upsert_stmt = insert_stmt.on_conflict_do_update(
index_elements=[Car.origin_id],
set_=self._postgres_upsert_set_map(insert_stmt),
).returning(Car.id)
car_id = int(session.execute(upsert_stmt).scalar_one())
action = "updated" if existed or car_by_url is not None else "inserted"
existing_urls = self._load_existing_image_urls(session, {car_id}).get(car_id, set())
images_upserted = 0
if self._should_sync_images(record, images, existing_urls):
if self._incoming_image_urls(images) == existing_urls:
images_upserted = len(existing_urls)
else:
images_upserted = self._replace_images_for_car(
session, car_id, images, record.origin_id,
)
return {
"car_id": car_id,
"images_upserted": images_upserted,
"action": action,
"material_changed": bool(changed_fields),
"changed_fields": list(changed_fields),
"reappeared": reappeared,
}
car = session.execute(
select(Car).where(or_(Car.origin_id == record.origin_id, Car.origin_url == record.origin_url))
).scalar_one_or_none()
action = "inserted"
if car is None:
car = Car(**payload)
session.add(car)
session.flush()
else:
action = "updated"
payload = self._prepare_update_payload(car, payload, record)
changed_fields, reappeared = self._classify_search_update(car, payload, record)
self._apply_update_payload(car, payload)
session.flush()
existing_urls = self._load_existing_image_urls(session, {int(car.id)}).get(int(car.id), set())
images_upserted = 0
if self._should_sync_images(record, images, existing_urls):
if self._incoming_image_urls(images) == existing_urls:
images_upserted = len(existing_urls)
else:
images_upserted = self._replace_images_for_car(session, int(car.id), images, record.origin_id)
return {
"car_id": int(car.id),
"images_upserted": images_upserted,
"action": action,
"material_changed": bool(changed_fields),
"changed_fields": list(changed_fields),
"reappeared": reappeared,
}
self._add_images(session, int(car.id), images)
session.flush()
return {
"car_id": int(car.id),
"images_upserted": len(images),
"action": action,
"material_changed": False,
"changed_fields": [],
"reappeared": False,
}
def upsert_cars_batch(self, records: list[CarRecord]) -> dict[str, object]:
"""Пакетный upsert нескольких автомобилей в одной транзакции.
Оптимизации:
- Дедупликация записей по origin_id перед вставкой.
- Chunked IN-queries для больших списков (обход лимита PG параметров).
- Пропуск перезаписи изображений, если набор URL не изменился.
- Один DELETE по car_id IN (...) вместо удаления по одному.
- Fallback на по-одному upsert если batch commit упал.
"""
# Дедупликация батча.
seen_ids: dict[str, int] = {}
unique_records: list[CarRecord] = []
for idx, r in enumerate(records):
key = r.origin_id or r.origin_url
if key in seen_ids:
logger.debug("Dedup: skipping duplicate record %s (idx %d vs %d)", key, idx, seen_ids[key])
continue
seen_ids[key] = idx
unique_records.append(r)
if len(unique_records) < len(records):
logger.info("Deduped batch: %d → %d records", len(records), len(unique_records))
records = unique_records
try:
return self._upsert_cars_batch_inner(records)
except Exception as exc:
logger.warning("Batch upsert failed (%s), falling back to individual upserts", exc)
return self._upsert_cars_individually(records)
def _upsert_cars_batch_inner(self, records: list[CarRecord]) -> dict[str, object]:
"""Внутренняя реализация batched upsert (одна транзакция)."""
if self._is_postgres():
return self._upsert_cars_batch_postgres(records)
inserted = 0
updated = 0
changed = 0
unchanged = 0
reappeared = 0
images_total = 0
inserted_origin_ids: list[str] = []
changed_origin_ids: list[str] = []
reappeared_origin_ids: list[str] = []
with self.session_scope() as session:
# Загружаем существующие записи.
origin_ids = [r.origin_id for r in records if r.origin_id]
origin_urls = [r.origin_url for r in records if r.origin_url]
existing_by_id, existing_by_url = self._load_existing_cars(session, origin_ids, origin_urls)
new_cars: list[tuple[Car, list[dict]]] = []
updated_cars: list[tuple[Car, list[dict], bool, CarRecord]] = []
for record in records:
payload = self._car_payload(record)
images = [image.model_dump(mode="python") for image in record.images]
car = existing_by_id.get(record.origin_id) or existing_by_url.get(record.origin_url)
if car is None:
car = Car(**payload)
session.add(car)
inserted += 1
inserted_origin_ids.append(record.origin_id)
new_cars.append((car, images))
else:
payload = self._prepare_update_payload(car, payload, record)
changed_fields, was_reappeared = self._classify_search_update(car, payload, record)
if was_reappeared:
reappeared += 1
reappeared_origin_ids.append(record.origin_id)
elif changed_fields:
changed += 1
changed_origin_ids.append(record.origin_id)
else:
unchanged += 1
self._apply_update_payload(car, payload)
updated += 1
updated_cars.append((car, images, self._skip_image_sync(record), record))
# Один flush.
session.flush()
# Добавляем картинки новым авто.
for car, images in new_cars:
self._add_images(session, int(car.id), images)
images_total += len(images)
update_cars_needing_images: list[tuple[Car, list[dict]]] = []
if updated_cars:
update_ids = [int(item[0].id) for item in updated_cars]
existing_images_map = self._load_existing_image_urls(session, set(update_ids))
for car, images, skip_image_sync, record in updated_cars:
old_image_urls = existing_images_map.get(int(car.id), set())
if skip_image_sync:
continue
if not self._should_sync_images(record, images, old_image_urls):
continue
if MOBILEDE_SKIP_IMAGES_FOR_UPDATED and old_image_urls and not record.images_confirmed:
continue
new_image_urls = {img.get("fullres_image", "") for img in images}
if new_image_urls != old_image_urls:
update_cars_needing_images.append((car, images))
else:
images_total += len(old_image_urls)
# Обновляем только изменённые картинки.
if update_cars_needing_images:
update_ids = [int(car.id) for car, _ in update_cars_needing_images]
for i in range(0, len(update_ids), _IN_CHUNK_SIZE):
chunk = update_ids[i:i + _IN_CHUNK_SIZE]
session.execute(delete(Image).where(Image.car_id.in_(chunk)))
for car, images in update_cars_needing_images:
self._add_images(session, int(car.id), images)
images_total += len(images)
return {
"inserted": inserted,
"updated": updated,
"changed": changed,
"unchanged": unchanged,
"reappeared": reappeared,
"images_upserted": images_total,
"inserted_origin_ids": inserted_origin_ids,
"changed_origin_ids": changed_origin_ids,
"reappeared_origin_ids": reappeared_origin_ids,
}
def _upsert_cars_individually(self, records: list[CarRecord]) -> dict[str, object]:
# Запасной поштучный upsert.
inserted = 0
updated = 0
changed = 0
unchanged = 0
reappeared = 0
images_total = 0
inserted_origin_ids: list[str] = []
changed_origin_ids: list[str] = []
reappeared_origin_ids: list[str] = []
failed_origin_ids: list[str] = []
for record in records:
try:
result = self.upsert_car(record)
action = result.get("action", "inserted")
if action == "inserted":
inserted += 1
inserted_origin_ids.append(record.origin_id)
else:
updated += 1
if bool(result.get("reappeared")):
reappeared += 1
reappeared_origin_ids.append(record.origin_id)
elif bool(result.get("material_changed")):
changed += 1
changed_origin_ids.append(record.origin_id)
else:
unchanged += 1
images_total += int(result.get("images_upserted", 0))
except Exception as exc:
logger.error("Individual upsert failed for %s: %s", record.origin_id, exc)
failed_origin_ids.append(record.origin_id)
return {
"inserted": inserted,
"updated": updated,
"changed": changed,
"unchanged": unchanged,
"reappeared": reappeared,
"images_upserted": images_total,
"inserted_origin_ids": inserted_origin_ids,
"changed_origin_ids": changed_origin_ids,
"reappeared_origin_ids": reappeared_origin_ids,
"failed": len(failed_origin_ids),
"failed_origin_ids": failed_origin_ids,
}
def mark_sold_not_seen_since(
self,
since_ts: datetime,
*,
prefix: str | tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES,
safety_ratio: float = 0.8,
brands: tuple[str, ...] = (),
) -> int:
"""Mark active cars as sold when they were not seen during a full refresh cycle.
Any unsold mobile.de car with ``last_seen_at < since_ts`` is considered absent
from the latest completed refresh cycle and can be marked as sold.
A safety guard prevents anomalous bulk updates.
"""
prefixes = (prefix,) if isinstance(prefix, str) else tuple(prefix)
scope_conditions = [
_origin_prefix_filter(Car.origin_id, prefixes),
Car.is_sold == False, # noqa: E712
]
if brands:
scope_conditions.append(func.lower(Car.brand).in_([brand.strip().casefold() for brand in brands if brand.strip()]))
with self.session_scope() as session:
total_active = int(
session.execute(
select(text("count(*)")).select_from(Car).where(
*scope_conditions,
)
).scalar()
or 0
)
if total_active <= 0:
return 0
would_mark = int(
session.execute(
select(text("count(*)")).select_from(Car).where(
*scope_conditions,
Car.last_seen_at < since_ts,
)
).scalar()
or 0
)
if would_mark <= 0:
return 0
if total_active > 100 and would_mark > int(total_active * max(0.0, min(1.0, safety_ratio))):
logger.error(
"mark_sold(last_seen) safety abort: would mark %d/%d (%.0f%%) as sold",
would_mark,
total_active,
(would_mark / max(1, total_active)) * 100,
)
return 0
result = session.execute(
update(Car)
.where(*scope_conditions)
.where(Car.last_seen_at < since_ts)
.values(is_sold=True, sold_at=since_ts)
)
count = int(result.rowcount or 0)
if count:
logger.info(
"Marked %d cars as sold by last_seen cutoff=%s brands=%s",
count,
since_ts.isoformat(),
list(brands) if brands else "all",
)
return count
def get_active_cars_batch_for_sold_probe(
self,
prefix: str | tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES,
*,
limit: int = 200,
newest_first: bool = True,
) -> list[tuple[int, str, datetime]]:
prefixes = (prefix,) if isinstance(prefix, str) else tuple(prefix)
order_by = Car.last_seen_at.desc() if newest_first else Car.last_seen_at.asc()
with self.session_scope() as session:
result = session.execute(
select(Car.id, Car.origin_url, Car.last_seen_at).where(
_origin_prefix_filter(Car.origin_id, prefixes),
Car.is_sold == False, # noqa: E712
).order_by(order_by).limit(limit)
)
return [
(int(row[0]), str(row[1]), row[2])
for row in result
if row and row[0] and row[1] and row[2]
]
def get_active_cars_batch_for_image_enrich(
self,
prefix: str | tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES,
*,
limit: int | None = None,
max_existing_images: int = 1,
stale_before: datetime | None = None,
) -> list[tuple[int, str, str, int]]:
if limit is not None and int(limit) <= 0:
return []
prefixes = (prefix,) if isinstance(prefix, str) else tuple(prefix)
with self.session_scope() as session:
image_count = func.count(Image.id)
conditions = [
_origin_prefix_filter(Car.origin_id, prefixes),
Car.is_sold == False, # noqa: E712
]
if stale_before is not None:
conditions.append(or_(Car.gallery_fetched_at.is_(None), Car.gallery_fetched_at <= stale_before))
else:
conditions.append(Car.gallery_fetched_at.is_(None))
stmt = (
select(Car.id, Car.origin_id, Car.origin_url, image_count.label("image_count"))
.outerjoin(Image, Image.car_id == Car.id)
.where(*conditions)
.group_by(Car.id, Car.origin_id, Car.origin_url, Car.gallery_fetched_at)
.order_by(
case((Car.gallery_fetched_at.is_(None), 0), else_=1).asc(),
case((Car.gallery_fetched_at.is_(None), Car.first_seen_at), else_=None).desc(),
Car.gallery_fetched_at.asc(),
Car.id.desc(),
)
)
if limit is not None:
stmt = stmt.limit(int(limit))
result = session.execute(stmt)
return [
(int(row[0]), str(row[1]), str(row[2]), int(row[3] or 0))
for row in result
if row and row[0] and row[1] and row[2]
]
def replace_car_gallery(self, origin_id: str, images: list[object]) -> dict[str, object]:
"""Replace only gallery rows and mark gallery completion; never update car data fields."""
image_payloads = [
image.model_dump(mode="python") if hasattr(image, "model_dump") else dict(image)
for image in images
]
listing_id = self._detail_listing_id(origin_id)
candidate_origin_ids = {str(origin_id)}
candidate_origin_ids.update(f"{prefix}{listing_id}" for prefix in MOBILEDE_ORIGIN_PREFIXES)
with self.session_scope() as session:
car = session.execute(
select(Car).where(Car.origin_id.in_(candidate_origin_ids))
).scalar_one_or_none()
if car is None:
raise LookupError(f"Car not found for gallery {origin_id}")
old_urls = self._load_existing_image_urls(session, {int(car.id)}).get(int(car.id), set())
new_urls = self._incoming_image_urls(image_payloads)
# Успешный Detail без изображений считаем завершённым, но не удаляем
# preview из Search: пустая галерея может быть временной особенностью payload.
changed = bool(image_payloads) and old_urls != new_urls
if changed:
session.execute(delete(Image).where(Image.car_id == int(car.id)))
self._add_images(session, int(car.id), image_payloads)
car.gallery_fetched_at = datetime.now(timezone.utc)
return {
"car_id": int(car.id),
"images_upserted": len(image_payloads) if image_payloads else len(old_urls),
"gallery_changed": changed,
}
def get_active_cars_batch_for_detail(
self,
prefix: str | tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES,
*,
limit: int | None = None,
stale_before: datetime | None = None,
) -> list[tuple[int, str, str]]:
"""Return fair Detail candidates without using gallery size as a proxy."""
if limit is not None and int(limit) <= 0:
return []
prefixes = (prefix,) if isinstance(prefix, str) else tuple(prefix)
conditions = [
_origin_prefix_filter(Car.origin_id, prefixes),
Car.is_sold == False, # noqa: E712
]
if stale_before is None:
conditions.append(Car.details_fetched_at.is_(None))
else:
conditions.append(or_(Car.details_fetched_at.is_(None), Car.details_fetched_at <= stale_before))
stmt = (
select(Car.id, Car.origin_id, Car.origin_url)
.where(*conditions)
.order_by(
case((Car.details_fetched_at.is_(None), 0), else_=1).asc(),
case((Car.details_fetched_at.is_(None), Car.first_seen_at), else_=None).asc(),
Car.details_fetched_at.asc(),
Car.id.asc(),
)
)
if limit is not None:
stmt = stmt.limit(int(limit))
with self.session_scope() as session:
rows = session.execute(stmt).all()
return [
(int(row[0]), str(row[1]), str(row[2]))
for row in rows
if row and row[0] and row[1] and row[2]
]
@staticmethod
def _detail_listing_id(origin_id: str) -> str:
return str(origin_id).rsplit(":", 1)[-1].strip()
def enqueue_detail_jobs(
self,
origin_ids: list[str],
*,
priority: int = 3,
reset: bool = False,
) -> dict[str, int]:
"""Create or reactivate durable jobs for known cars."""
unique_inputs = list(dict.fromkeys(str(value).strip() for value in origin_ids if str(value).strip()))
listing_ids = [self._detail_listing_id(value) for value in unique_inputs]
listing_ids = [value for value in listing_ids if value]
if not listing_ids:
return {"ready": 0, "deferred": 0, "missing": 0}
candidate_origin_ids = set(unique_inputs)
for listing_id in listing_ids:
candidate_origin_ids.update(f"{prefix}{listing_id}" for prefix in MOBILEDE_ORIGIN_PREFIXES)
now = datetime.now(timezone.utc)
ready = 0
deferred = 0
with self.session_scope() as session:
cars = session.execute(
select(Car).where(Car.origin_id.in_(candidate_origin_ids))
).scalars().all()
cars_by_listing = {self._detail_listing_id(car.origin_id): car for car in cars}
existing_jobs = session.execute(
select(DetailJob).where(DetailJob.listing_id.in_(listing_ids))
).scalars().all()
jobs_by_listing = {job.listing_id: job for job in existing_jobs}
for listing_id in listing_ids:
car = cars_by_listing.get(listing_id)
if car is None:
continue
job = jobs_by_listing.get(listing_id)
if job is None:
job = DetailJob(
car_id=int(car.id),
listing_id=listing_id,
status="pending",
priority=max(0, min(9, int(priority))),
available_at=now,
)
session.add(job)
jobs_by_listing[listing_id] = job
ready += 1
continue
job.car_id = int(car.id)
job.priority = max(int(job.priority or 0), max(0, min(9, int(priority))))
lease_expires_at = _utc_datetime(job.lease_expires_at)
available_at = _utc_datetime(job.available_at)
lease_active = job.status == "leased" and lease_expires_at is not None and lease_expires_at > now
cooldown_active = (
job.status in {"unavailable", "retry"}
and available_at is not None
and available_at > now
)
if lease_active or (cooldown_active and not reset):
deferred += 1
continue
if job.status == "completed" and car.gallery_fetched_at is not None and not reset:
deferred += 1
continue
job.status = "pending"
job.available_at = now
job.lease_owner = None
job.lease_expires_at = None
if reset:
job.error_kind = None
job.last_error = None
ready += 1
return {"ready": ready, "deferred": deferred, "missing": len(listing_ids) - len(cars_by_listing)}
def reserve_detail_jobs(
self,
*,
limit: int,
lease_seconds: int = 900,
listing_ids: list[str] | None = None,
) -> list[dict[str, object]]:
"""Atomically lease ready jobs; PostgreSQL workers use SKIP LOCKED."""
if int(limit) <= 0:
return []
now = datetime.now(timezone.utc)
lease_until = now + timedelta(seconds=max(60, int(lease_seconds)))
with self.session_scope() as session:
session.execute(
update(DetailJob)
.where(DetailJob.status == "leased")
.where(DetailJob.lease_expires_at <= now)
.values(status="pending", lease_owner=None, lease_expires_at=None, updated_at=now)
)
conditions = [
DetailJob.status.in_(("pending", "retry")),
DetailJob.available_at <= now,
]
if listing_ids is not None:
conditions.append(DetailJob.listing_id.in_(listing_ids))
stmt = (
select(DetailJob)
.where(*conditions)
.order_by(
DetailJob.priority.desc(),
DetailJob.available_at.asc(),
DetailJob.created_at.asc(),
DetailJob.id.asc(),
)
.limit(int(limit))
)
if self._is_postgres():
stmt = stmt.with_for_update(skip_locked=True)
jobs = session.execute(stmt).scalars().all()
leases: list[dict[str, object]] = []
for job in jobs:
owner = uuid.uuid4().hex
job.status = "leased"
job.lease_owner = owner
job.lease_expires_at = lease_until
job.attempts = int(job.attempts or 0) + 1
job.updated_at = now
leases.append({
"job_id": int(job.id),
"listing_id": str(job.listing_id),
"lease_owner": owner,
"attempts": int(job.attempts),
"priority": int(job.priority),
})
return leases
def complete_detail_job(self, listing_id: str, *, lease_owner: str | None = None) -> bool:
now = datetime.now(timezone.utc)
conditions = [DetailJob.listing_id == self._detail_listing_id(listing_id)]
if lease_owner:
conditions.append(DetailJob.lease_owner == lease_owner)
with self.session_scope() as session:
result = session.execute(
update(DetailJob)
.where(*conditions)
.values(
status="completed",
available_at=now,
lease_owner=None,
lease_expires_at=None,
last_http_status=200,
error_kind=None,
last_error=None,
updated_at=now,
)
)
return bool(result.rowcount)
def fail_detail_job(
self,
listing_id: str,
*,
lease_owner: str | None = None,
status: str = "retry",
cooldown_seconds: int = 60,
http_status: int | None = None,
error_kind: str | None = None,
last_error: str | None = None,
) -> bool:
now = datetime.now(timezone.utc)
conditions = [DetailJob.listing_id == self._detail_listing_id(listing_id)]
if lease_owner:
conditions.append(DetailJob.lease_owner == lease_owner)
with self.session_scope() as session:
result = session.execute(
update(DetailJob)
.where(*conditions)
.values(
status=status,
available_at=now + timedelta(seconds=max(0, int(cooldown_seconds))),
lease_owner=None,
lease_expires_at=None,
last_http_status=http_status,
error_kind=error_kind,
last_error=(last_error or "")[:4000] or None,
updated_at=now,
)
)
return bool(result.rowcount)
def release_detail_job(self, listing_id: str, *, lease_owner: str | None = None) -> bool:
return self.fail_detail_job(
listing_id,
lease_owner=lease_owner,
status="pending",
cooldown_seconds=0,
)
def mark_cars_sold_by_ids(self, car_ids: list[int], *, sold_at: datetime | None = None) -> int:
if not car_ids:
return 0
ts = sold_at or datetime.now(timezone.utc)
with self.session_scope() as session:
result = session.execute(
update(Car)
.where(Car.id.in_(car_ids))
.where(Car.is_sold == False) # noqa: E712
.values(is_sold=True, sold_at=ts)
)
count = int(result.rowcount or 0)
if count:
logger.info("Marked %d cars as sold by explicit id probe", count)
return count