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:`` while some newer helper code used ``mobilede:``. 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