diff --git a/alembic/versions/007_runtime_storage.py b/alembic/versions/007_runtime_storage.py new file mode 100644 index 0000000..a3e22e4 --- /dev/null +++ b/alembic/versions/007_runtime_storage.py @@ -0,0 +1,73 @@ +from typing import Sequence, Union + +import sqlalchemy as sa +from alembic import op + +revision: str = "007_runtime_storage" +down_revision: Union[str, None] = "006_add_details_fetched_at" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.execute( + 'CREATE INDEX IF NOT EXISTS "ix_MOBILEDE_cars_active_last_seen_mobile_de" ' + 'ON "MOBILEDE_cars" (last_seen_at) ' + "WHERE is_sold = FALSE AND origin_id LIKE 'mobile.de:%'" + ) + op.create_table( + "MOBILEDE_detail_jobs", + sa.Column("id", sa.BigInteger(), autoincrement=True, nullable=False), + sa.Column("car_id", sa.BigInteger(), nullable=False), + sa.Column("listing_id", sa.String(length=128), nullable=False), + sa.Column("status", sa.String(length=32), server_default="pending", nullable=False), + sa.Column("priority", sa.Integer(), server_default="3", nullable=False), + sa.Column("attempts", sa.Integer(), server_default="0", nullable=False), + sa.Column("available_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False), + sa.Column("lease_owner", sa.String(length=64), nullable=True), + sa.Column("lease_expires_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("last_http_status", sa.Integer(), nullable=True), + sa.Column("error_kind", sa.String(length=64), nullable=True), + sa.Column("last_error", sa.Text(), nullable=True), + sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False), + sa.ForeignKeyConstraint(["car_id"], ["MOBILEDE_cars.id"], ondelete="CASCADE"), + sa.PrimaryKeyConstraint("id"), + sa.UniqueConstraint("car_id"), + sa.UniqueConstraint("listing_id"), + ) + op.create_index("ix_MOBILEDE_detail_jobs_listing_id", "MOBILEDE_detail_jobs", ["listing_id"]) + op.create_index("ix_MOBILEDE_detail_jobs_status", "MOBILEDE_detail_jobs", ["status"]) + op.create_index("ix_MOBILEDE_detail_jobs_available_at", "MOBILEDE_detail_jobs", ["available_at"]) + op.create_index("ix_MOBILEDE_detail_jobs_lease_expires_at", "MOBILEDE_detail_jobs", ["lease_expires_at"]) + op.create_index( + "ix_MOBILEDE_detail_jobs_ready", + "MOBILEDE_detail_jobs", + ["status", "available_at", "priority", "created_at"], + ) + op.create_index("ix_MOBILEDE_detail_jobs_lease", "MOBILEDE_detail_jobs", ["status", "lease_expires_at"]) + op.add_column("MOBILEDE_cars", sa.Column("gallery_fetched_at", sa.DateTime(timezone=True), nullable=True)) + op.create_index("ix_MOBILEDE_cars_gallery_fetched_at", "MOBILEDE_cars", ["gallery_fetched_at"]) + op.execute( + ''' + UPDATE "MOBILEDE_cars" AS car + SET gallery_fetched_at = car.details_fetched_at + WHERE car.details_fetched_at IS NOT NULL + AND EXISTS ( + SELECT 1 FROM "MOBILEDE_images" AS image WHERE image.car_id = car.id + ) + ''' + ) + + +def downgrade() -> None: + op.drop_index("ix_MOBILEDE_cars_gallery_fetched_at", table_name="MOBILEDE_cars") + op.drop_column("MOBILEDE_cars", "gallery_fetched_at") + op.drop_index("ix_MOBILEDE_detail_jobs_lease", table_name="MOBILEDE_detail_jobs") + op.drop_index("ix_MOBILEDE_detail_jobs_ready", table_name="MOBILEDE_detail_jobs") + op.drop_index("ix_MOBILEDE_detail_jobs_lease_expires_at", table_name="MOBILEDE_detail_jobs") + op.drop_index("ix_MOBILEDE_detail_jobs_available_at", table_name="MOBILEDE_detail_jobs") + op.drop_index("ix_MOBILEDE_detail_jobs_status", table_name="MOBILEDE_detail_jobs") + op.drop_index("ix_MOBILEDE_detail_jobs_listing_id", table_name="MOBILEDE_detail_jobs") + op.drop_table("MOBILEDE_detail_jobs") + op.execute('DROP INDEX IF EXISTS "ix_MOBILEDE_cars_active_last_seen_mobile_de"') \ No newline at end of file diff --git a/mobilede_scraper/storage/db.py b/mobilede_scraper/storage/db.py index 86bd9d7..8fe4b16 100644 --- a/mobilede_scraper/storage/db.py +++ b/mobilede_scraper/storage/db.py @@ -1,8 +1,9 @@ import logging import os import re +import uuid from contextlib import contextmanager -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from typing import Any, Iterator from sqlalchemy import case, create_engine, delete, func, or_, select, text, update @@ -10,7 +11,7 @@ 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, Image, SyncRun +from .models import Base, Car, DetailJob, Image, SyncRun from .schemas import CarRecord logger = logging.getLogger("MOBILEDE_scraper.db") @@ -23,10 +24,43 @@ CAR_DB_FIELDS = { } 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 -CAR_TABLE_NAME = Car.__tablename__ 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"}, @@ -99,7 +133,14 @@ class PersistenceService: def start_sync_run(self, lane: str) -> int: with self.session_scope() as session: now = datetime.now(timezone.utc) - stale_runs = session.execute(select(SyncRun).where(SyncRun.status == "running")).scalars().all() + 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 @@ -124,13 +165,6 @@ class PersistenceService: run.images_upserted = images_upserted run.error_summary = error_summary - def get_existing_origin_urls(self, origin_urls: list[str]) -> set[str]: - if not origin_urls: - return set() - with self.session_scope() as session: - rows = session.execute(select(Car.origin_url).where(Car.origin_url.in_(origin_urls))).all() - return {str(row[0]) for row in rows if row and row[0]} - def get_existing_origin_ids(self, origin_ids: list[str]) -> set[str]: if not origin_ids: return set() @@ -138,36 +172,6 @@ class PersistenceService: 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]} - def get_existing_urls_and_ids( - self, origin_urls: list[str], origin_ids: list[str], - ) -> tuple[set[str], set[str]]: - # Загрузка URL и origin_id. - if not origin_urls and not origin_ids: - return set(), set() - urls: set[str] = set() - ids: set[str] = set() - with self.session_scope() as session: - max_len = max(len(origin_urls), len(origin_ids), 1) - for i in range(0, max_len, _IN_CHUNK_SIZE): - url_chunk = origin_urls[i:i + _IN_CHUNK_SIZE] - id_chunk = origin_ids[i:i + _IN_CHUNK_SIZE] - conditions = [] - if url_chunk: - conditions.append(Car.origin_url.in_(url_chunk)) - if id_chunk: - conditions.append(Car.origin_id.in_(id_chunk)) - if not conditions: - continue - rows = session.execute( - select(Car.origin_url, Car.origin_id).where(or_(*conditions)) - ).all() - for r in rows: - if r[0]: - urls.add(str(r[0])) - if r[1]: - ids.add(str(r[1])) - return urls, ids - @staticmethod def _add_images(session: Session, car_id: int, images: list[dict[str, object]]) -> None: if not images: @@ -196,6 +200,8 @@ class PersistenceService: 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: @@ -212,6 +218,27 @@ class PersistenceService: 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)) @@ -318,10 +345,16 @@ class PersistenceService: 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, int]: + 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] @@ -338,6 +371,15 @@ class PersistenceService: 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] = { @@ -360,6 +402,7 @@ class PersistenceService: updated += 1 else: inserted += 1 + inserted_origin_ids.append(record.origin_id) entries.append(entry) session.flush() @@ -441,7 +484,17 @@ class PersistenceService: self._add_images(session, car_id, images) images_total += len(images) - return {"inserted": inserted, "updated": updated, "images_upserted": images_total} + 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): # Вставка или обновление авто. @@ -462,12 +515,18 @@ class PersistenceService: 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], @@ -485,7 +544,14 @@ class PersistenceService: 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} + 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)) @@ -498,6 +564,7 @@ class PersistenceService: 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()) @@ -507,11 +574,25 @@ class PersistenceService: 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} + 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} - def upsert_cars_batch(self, records: list[CarRecord]) -> dict[str, int]: + 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 нескольких автомобилей в одной транзакции. Оптимизации: @@ -542,14 +623,20 @@ class PersistenceService: 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, int]: + 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: # Загружаем существующие записи. @@ -570,9 +657,19 @@ class PersistenceService: 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)) @@ -613,159 +710,64 @@ class PersistenceService: self._add_images(session, int(car.id), images) images_total += len(images) - return {"inserted": inserted, "updated": updated, "images_upserted": images_total} + 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, int]: + 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) - return {"inserted": inserted, "updated": updated, "images_upserted": images_total} - - def mark_sold_not_in_listing(self, active_origin_ids: set[str], lane: str = "MOBILEDE") -> int: - """Помечает авто как проданные, если их нет в активном листинге (по origin_id).""" - if not active_origin_ids: - return 0 - with self.session_scope() as session: - stmt = ( - update(Car) - .where(Car.origin_id.notin_(active_origin_ids)) - .where(Car.is_sold == False) # noqa: E712 - .where(_origin_prefix_filter(Car.origin_id)) - .values(is_sold=True, sold_at=datetime.now(timezone.utc)) - ) - result = session.execute(stmt) - count = result.rowcount or 0 - if count: - logger.info("Marked %d cars as sold (no longer in listing)", count) - return count - - def mark_sold_not_in_listing_by_urls(self, active_origin_urls: set[str], lane: str = "MOBILEDE") -> int: - """Помечает авто как проданные, если их URL нет в активном листинге. - - Для PostgreSQL использует временную таблицу + LEFT JOIN вместо NOT IN, - что кардинально быстрее при больших объёмах (100K+ URLs). - Встроенная защита: нормализация URL (тильда/дефис) + safety-check на аномальный процент. - """ - if not active_origin_urls: - return 0 - - # Нормализация URL. - def _norm(url: str) -> str: - return url.replace("~", "-") - - normalized_urls = {_norm(u) for u in active_origin_urls} - - is_postgres = "postgresql" in self.settings.database.url - - with self.session_scope() as session: - if is_postgres: - # Считаем кандидатов на sold. - total_active = session.execute( - text(f""" - SELECT count(*) - FROM {CAR_TABLE_NAME} - WHERE is_sold = FALSE - AND (origin_id LIKE 'mobile.de:%%' OR origin_id LIKE 'mobilede:%%') - """) - ).scalar() or 0 - - if total_active == 0: - return 0 - - # Временная таблица URL. - session.execute(text("CREATE TEMP TABLE IF NOT EXISTS _active_urls (url TEXT NOT NULL) ON COMMIT DROP")) - session.execute(text("TRUNCATE _active_urls")) - - # Вставляем URL чанками. - url_list = list(normalized_urls) - for i in range(0, len(url_list), _IN_CHUNK_SIZE): - chunk = url_list[i:i + _IN_CHUNK_SIZE] - values = ",".join(f"(:{f'u{j}'})" for j in range(len(chunk))) - params = {f"u{j}": url for j, url in enumerate(chunk)} - session.execute(text(f"INSERT INTO _active_urls (url) VALUES {values}"), params) - - # Индекс для JOIN. - session.execute(text("CREATE INDEX IF NOT EXISTS _ix_active_urls ON _active_urls (url)")) - - # Считаем будущие sold. - would_mark = session.execute(text(""" - SELECT count(*) - FROM {car_table} c - LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url - WHERE a.url IS NULL - AND c.is_sold = FALSE - AND (c.origin_id LIKE 'mobile.de:%%' OR c.origin_id LIKE 'mobilede:%%') - """.format(car_table=CAR_TABLE_NAME))).scalar() or 0 - - # Защита от аномалии. - if total_active > 100 and would_mark > total_active * 0.8: - logger.error( - "mark_sold safety abort: would mark %d/%d (%.0f%%) as sold — likely URL format mismatch", - would_mark, total_active, would_mark / total_active * 100, - ) - return 0 - - # Массовая пометка sold. - result = session.execute(text(""" - UPDATE {car_table} - SET is_sold = TRUE, sold_at = NOW() - FROM ( - SELECT c.id - FROM {car_table} c - LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url - WHERE a.url IS NULL - AND c.is_sold = FALSE - AND (c.origin_id LIKE 'mobile.de:%%' OR c.origin_id LIKE 'mobilede:%%') - ) sub - WHERE {car_table}.id = sub.id - """.format(car_table=CAR_TABLE_NAME))) - count = result.rowcount or 0 - else: - # Упрощённый путь для SQLite. - stmt = ( - update(Car) - .where(Car.is_sold == False) # noqa: E712 - .where(_origin_prefix_filter(Car.origin_id)) - .values(is_sold=True, sold_at=datetime.now(timezone.utc)) - ) - # Загружаем active URL. - all_active = session.execute( - select(Car.id, Car.origin_url).where( - Car.is_sold == False, _origin_prefix_filter(Car.origin_id) # noqa: E712 - ) - ).all() - mark_ids = [row[0] for row in all_active if _norm(row[1]) not in normalized_urls] - - if not mark_ids: - return 0 - total_active = len(all_active) - if total_active > 100 and len(mark_ids) > total_active * 0.8: - logger.error( - "mark_sold safety abort: would mark %d/%d (%.0f%%) as sold — likely URL format mismatch", - len(mark_ids), total_active, len(mark_ids) / total_active * 100, - ) - return 0 - - for i in range(0, len(mark_ids), _IN_CHUNK_SIZE): - chunk = mark_ids[i:i + _IN_CHUNK_SIZE] - session.execute(update(Car).where(Car.id.in_(chunk)).values(is_sold=True, sold_at=datetime.now(timezone.utc))) - count = len(mark_ids) - - if count: - logger.info("Marked %d cars as sold by URL (no longer in listing)", count) - return count + 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, @@ -773,6 +775,7 @@ class PersistenceService: *, 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. @@ -781,12 +784,17 @@ class PersistenceService: 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( - _origin_prefix_filter(Car.origin_id, prefixes), - Car.is_sold == False, # noqa: E712 + *scope_conditions, ) ).scalar() or 0 @@ -797,8 +805,7 @@ class PersistenceService: would_mark = int( session.execute( select(text("count(*)")).select_from(Car).where( - _origin_prefix_filter(Car.origin_id, prefixes), - Car.is_sold == False, # noqa: E712 + *scope_conditions, Car.last_seen_at < since_ts, ) ).scalar() @@ -818,72 +825,19 @@ class PersistenceService: result = session.execute( update(Car) - .where(_origin_prefix_filter(Car.origin_id, prefixes)) - .where(Car.is_sold == False) # noqa: E712 + .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", count, since_ts.isoformat()) - return count - - def get_all_origin_ids_for_lane(self, prefix: str | tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES) -> set[str]: - """Возвращает все известные origin_id для заданного префикса. - - Использует yield_per для потоковой загрузки при большом количестве записей. - """ - prefixes = (prefix,) if isinstance(prefix, str) else tuple(prefix) - with self.session_scope() as session: - result = session.execute( - select(Car.origin_id).where(_origin_prefix_filter(Car.origin_id, prefixes)).execution_options(yield_per=10000) - ) - return {str(row[0]) for row in result if row and row[0]} - - def get_all_active_origin_urls_for_lane(self, prefix: str | tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES) -> set[str]: - """Возвращает origin_url всех активных (не проданных) авто для заданного lane-префикса.""" - prefixes = (prefix,) if isinstance(prefix, str) else tuple(prefix) - with self.session_scope() as session: - result = session.execute( - select(Car.origin_url).where( - _origin_prefix_filter(Car.origin_id, prefixes), - Car.is_sold == False, # noqa: E712 - ).execution_options(yield_per=10000) - ) - return {str(row[0]) for row in result if row and row[0]} - - def count_active_cars_for_lane(self, prefix: str | tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES) -> int: - """Возвращает количество активных (не проданных) авто для заданного lane-префикса.""" - from sqlalchemy import func as sa_func - prefixes = (prefix,) if isinstance(prefix, str) else tuple(prefix) - with self.session_scope() as session: - result = session.execute( - select(sa_func.count()).select_from(Car).where( - _origin_prefix_filter(Car.origin_id, prefixes), - Car.is_sold == False, # noqa: E712 + 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 int(result.scalar() or 0) - - def get_active_origin_urls_batch_for_refresh( - self, - prefix: str | tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES, - offset: int = 0, - limit: int = 500, - ) -> list[str]: - """Возвращает батч origin_url активных авто для rolling refresh. - - Сортировка по last_seen_at ASC — давно не обновлённые идут первыми. - """ - prefixes = (prefix,) if isinstance(prefix, str) else tuple(prefix) - with self.session_scope() as session: - result = session.execute( - select(Car.origin_url).where( - _origin_prefix_filter(Car.origin_id, prefixes), - Car.is_sold == False, # noqa: E712 - ).order_by(Car.last_seen_at.asc()).offset(offset).limit(limit) - ) - return [str(row[0]) for row in result if row and row[0]] + return count def get_active_cars_batch_for_sold_probe( self, @@ -925,17 +879,19 @@ class PersistenceService: Car.is_sold == False, # noqa: E712 ] if stale_before is not None: - conditions.append(or_(Car.details_fetched_at.is_(None), Car.details_fetched_at <= stale_before)) + 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.details_fetched_at) - .having(image_count <= max(0, int(max_existing_images))) + .group_by(Car.id, Car.origin_id, Car.origin_url, Car.gallery_fetched_at) .order_by( - case((Car.details_fetched_at.is_(None), 0), else_=1).asc(), - Car.details_fetched_at.asc(), - Car.id.asc(), + 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: @@ -947,6 +903,273 @@ class PersistenceService: 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 diff --git a/mobilede_scraper/storage/models.py b/mobilede_scraper/storage/models.py index a7e2b4d..f0137d7 100644 --- a/mobilede_scraper/storage/models.py +++ b/mobilede_scraper/storage/models.py @@ -57,6 +57,7 @@ class Car(Base): last_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), index=True) sold_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True) details_fetched_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True) + gallery_fetched_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True) images: Mapped[list["Image"]] = relationship("Image", back_populates="car", cascade="all, delete-orphan") @@ -70,6 +71,33 @@ class Image(Base): car: Mapped[Car] = relationship("Car", back_populates="images") +class DetailJob(Base): + __tablename__ = "MOBILEDE_detail_jobs" + __table_args__ = ( + Index("ix_MOBILEDE_detail_jobs_ready", "status", "available_at", "priority", "created_at"), + Index("ix_MOBILEDE_detail_jobs_lease", "status", "lease_expires_at"), + ) + id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) + car_id: Mapped[int] = mapped_column( + BigInteger().with_variant(Integer, "sqlite"), + ForeignKey("MOBILEDE_cars.id", ondelete="CASCADE"), + nullable=False, + unique=True, + ) + listing_id: Mapped[str] = mapped_column(String(128), nullable=False, unique=True, index=True) + status: Mapped[str] = mapped_column(String(32), nullable=False, default="pending", index=True) + priority: Mapped[int] = mapped_column(Integer, nullable=False, default=3) + attempts: Mapped[int] = mapped_column(Integer, nullable=False, default=0) + available_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), index=True) + lease_owner: Mapped[str | None] = mapped_column(String(64), nullable=True) + lease_expires_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True) + last_http_status: Mapped[int | None] = mapped_column(Integer, nullable=True) + error_kind: Mapped[str | None] = mapped_column(String(64), nullable=True) + last_error: Mapped[str | None] = mapped_column(Text, nullable=True) + created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now()) + updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), onupdate=func.now()) + + class SyncRun(Base): __tablename__ = "MOBILEDE_sync_runs" id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) diff --git a/mobilede_scraper/storage/schemas.py b/mobilede_scraper/storage/schemas.py index 472d518..85e73c8 100644 --- a/mobilede_scraper/storage/schemas.py +++ b/mobilede_scraper/storage/schemas.py @@ -41,6 +41,7 @@ class CarRecord(BaseModel): last_seen_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc)) sold_at: datetime | None = None details_fetched_at: datetime | None = None + gallery_fetched_at: datetime | None = None skip_image_sync: bool = False details_confirmed: bool = False images_confirmed: bool = False @@ -93,5 +94,6 @@ class CarRead(BaseModel): last_seen_at: datetime | None = None sold_at: datetime | None = None details_fetched_at: datetime | None = None + gallery_fetched_at: datetime | None = None images: list[ImageRead] = Field(default_factory=list) diff --git a/tests/test_db.py b/tests/test_db.py index 14ff32e..256d949 100644 --- a/tests/test_db.py +++ b/tests/test_db.py @@ -4,12 +4,13 @@ import tempfile import unittest from datetime import datetime, timedelta, timezone from pathlib import Path +from unittest.mock import patch from sqlalchemy import select from mobilede_scraper.core.config import Settings from mobilede_scraper.storage.db import PersistenceService -from mobilede_scraper.storage.models import Car, Image, SyncRun +from mobilede_scraper.storage.models import Car, DetailJob, Image, SyncRun from mobilede_scraper.storage.schemas import CarRecord, ImageRecord @@ -52,6 +53,15 @@ class TestPersistenceServiceIntegration(unittest.TestCase): ], ) + @classmethod + def _search_record(cls, origin_id: str, *, price: int = 1000) -> CarRecord: + record = cls._record(origin_id, price=price) + record.details_confirmed = False + record.images_confirmed = False + record.preserve_existing_details = True + record.images = [] + return record + @staticmethod def _as_utc(value: datetime | None) -> datetime | None: if value is None: @@ -134,6 +144,59 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.assertEqual(result["images_upserted"], 1) self.assertEqual(len(images), 1) + def test_batch_upsert_returns_only_inserted_origin_ids(self) -> None: + existing = self._record("mobile.de:existing") + self.persistence.upsert_car(existing) + + updated = self._record("mobile.de:existing", price=2000) + inserted = self._record("mobile.de:new", price=3000) + result = self.persistence.upsert_cars_batch([updated, inserted]) + + self.assertEqual(result["inserted"], 1) + self.assertEqual(result["updated"], 1) + self.assertEqual(result["inserted_origin_ids"], ["mobile.de:new"]) + + def test_batch_search_upsert_classifies_changed_unchanged_and_reappeared(self) -> None: + unchanged = self._search_record("mobile.de:unchanged", price=1000) + changed = self._search_record("mobile.de:changed", price=1000) + reappeared = self._search_record("mobile.de:reappeared", price=1000) + self.persistence.upsert_cars_batch([unchanged, changed, reappeared]) + + with self.persistence.session_scope() as session: + car = session.execute(select(Car).where(Car.origin_id == reappeared.origin_id)).scalar_one() + car.is_sold = True + car.sold_at = datetime.now(timezone.utc) + + unchanged_again = self._search_record("mobile.de:unchanged", price=1000) + unchanged_again.last_seen_at = datetime.now(timezone.utc) + timedelta(minutes=1) + changed_again = self._search_record("mobile.de:changed", price=1500) + reappeared_again = self._search_record("mobile.de:reappeared", price=1000) + result = self.persistence.upsert_cars_batch([unchanged_again, changed_again, reappeared_again]) + + self.assertEqual(result["updated"], 3) + self.assertEqual(result["changed"], 1) + self.assertEqual(result["unchanged"], 1) + self.assertEqual(result["reappeared"], 1) + self.assertEqual(result["changed_origin_ids"], ["mobile.de:changed"]) + self.assertEqual(result["reappeared_origin_ids"], ["mobile.de:reappeared"]) + + def test_weak_search_fallback_does_not_create_false_material_change(self) -> None: + detail = self._record("mobile.de:fallback", price=1000) + detail.drive = "4WD" + detail.body_type = "SUV" + detail.color = "black" + self.persistence.upsert_car(detail) + + search = self._search_record("mobile.de:fallback", price=1000) + search.drive = None + search.body_type = "OTHER" + search.color = "other" + result = self.persistence.upsert_cars_batch([search]) + + self.assertEqual(result["changed"], 0) + self.assertEqual(result["unchanged"], 1) + self.assertEqual(result["changed_origin_ids"], []) + def test_search_update_preserves_confirmed_details_and_gallery(self) -> None: detail = self._record("mobile.de:confirmed", price=1000) detail.drive = "4WD" @@ -217,7 +280,7 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.assertEqual(car.body_type, "COUPE") self.assertEqual([image.fullres_image for image in images], ["https://example.test/refreshed.jpg"]) - def test_detail_timestamp_is_set_and_preserved_by_search_update(self) -> None: + def test_search_update_preserves_detail_and_gallery_completion(self) -> None: search = self._record("mobile.de:detail-ts") search.details_confirmed = False search.images_confirmed = False @@ -236,6 +299,15 @@ class TestPersistenceServiceIntegration(unittest.TestCase): car = session.execute(select(Car).where(Car.origin_id == search.origin_id)).scalar_one() fetched_at = self._as_utc(car.details_fetched_at) + gallery = [ + ImageRecord( + fullres_image="https://example.test/gallery.jpg", + preview_image="https://example.test/gallery-preview.jpg", + order_index=0, + ) + ] + self.persistence.replace_car_gallery(search.origin_id, gallery) + self.assertIsNotNone(fetched_at) self.assertGreaterEqual(fetched_at, before) @@ -249,8 +321,124 @@ class TestPersistenceServiceIntegration(unittest.TestCase): car = session.execute(select(Car).where(Car.origin_id == search.origin_id)).scalar_one() self.assertEqual(self._as_utc(car.details_fetched_at), fetched_at) + self.assertIsNotNone(car.gallery_fetched_at) - def test_start_sync_run_marks_stale_running_runs_as_failed(self) -> None: + def test_replace_car_gallery_changes_only_images_and_completion_timestamp(self) -> None: + search = self._search_record("mobile.de:gallery-only", price=1234) + search.brand = "BMW" + search.model = "M3" + self.persistence.upsert_car(search) + + gallery = [ + ImageRecord( + fullres_image=f"https://example.test/gallery-{index}.jpg", + preview_image=f"https://example.test/gallery-{index}-preview.jpg", + order_index=index, + ) + for index in range(3) + ] + result = self.persistence.replace_car_gallery(search.origin_id, gallery) + + with self.persistence.session_scope() as session: + car = session.execute(select(Car).where(Car.origin_id == search.origin_id)).scalar_one() + images = session.execute( + select(Image).where(Image.car_id == car.id).order_by(Image.order_index.asc()) + ).scalars().all() + + self.assertEqual(result["images_upserted"], 3) + self.assertEqual(car.brand, "BMW") + self.assertEqual(car.model, "M3") + self.assertEqual(car.price, 1234) + self.assertIsNone(car.details_fetched_at) + self.assertIsNotNone(car.gallery_fetched_at) + self.assertEqual(len(images), 3) + + def test_detail_candidates_ignore_gallery_size_and_prefer_oldest_unfetched(self) -> None: + old = self._search_record("mobile.de:old-detail") + old.first_seen_at = datetime(2026, 1, 1, tzinfo=timezone.utc) + old.images = [ + ImageRecord( + fullres_image=f"https://example.test/{index}.jpg", + preview_image=f"https://example.test/{index}-preview.jpg", + order_index=index, + ) + for index in range(4) + ] + old.images_confirmed = True + newer = self._search_record("mobile.de:new-detail") + newer.first_seen_at = datetime(2026, 2, 1, tzinfo=timezone.utc) + self.persistence.upsert_car(newer) + self.persistence.upsert_car(old) + + candidates = self.persistence.get_active_cars_batch_for_detail(limit=2) + + self.assertEqual([origin_id for _car_id, origin_id, _url in candidates], [ + "mobile.de:old-detail", + "mobile.de:new-detail", + ]) + + def test_detail_jobs_enforce_cooldown_reset_and_exclusive_leases(self) -> None: + self.persistence.upsert_car(self._search_record("mobile.de:job-1")) + self.persistence.upsert_car(self._search_record("mobile.de:job-2")) + + result = self.persistence.enqueue_detail_jobs(["mobile.de:job-1", "mobile.de:job-2"], priority=7) + self.assertEqual(result, {"ready": 2, "deferred": 0, "missing": 0}) + + first = self.persistence.reserve_detail_jobs(limit=1) + second = self.persistence.reserve_detail_jobs(limit=2) + self.assertEqual(len(first), 1) + self.assertEqual(len(second), 1) + self.assertNotEqual(first[0]["listing_id"], second[0]["listing_id"]) + + listing_id = str(first[0]["listing_id"]) + lease_owner = str(first[0]["lease_owner"]) + self.assertTrue(self.persistence.fail_detail_job( + listing_id, + lease_owner=lease_owner, + status="unavailable", + cooldown_seconds=3600, + http_status=404, + error_kind="unavailable", + last_error="not found", + )) + self.assertEqual( + self.persistence.enqueue_detail_jobs([f"mobile.de:{listing_id}"]), + {"ready": 0, "deferred": 1, "missing": 0}, + ) + self.assertEqual(self.persistence.reserve_detail_jobs(limit=1, listing_ids=[listing_id]), []) + + self.assertEqual( + self.persistence.enqueue_detail_jobs([f"mobile.de:{listing_id}"], priority=9, reset=True), + {"ready": 1, "deferred": 0, "missing": 0}, + ) + retry = self.persistence.reserve_detail_jobs(limit=1, listing_ids=[listing_id]) + self.assertEqual(len(retry), 1) + self.assertTrue(self.persistence.complete_detail_job(listing_id, lease_owner=str(retry[0]["lease_owner"]))) + + with self.persistence.session_scope() as session: + job = session.execute(select(DetailJob).where(DetailJob.listing_id == listing_id)).scalar_one() + self.assertEqual(job.status, "completed") + self.assertEqual(job.priority, 9) + self.assertEqual(job.last_http_status, 200) + self.assertIsNone(job.lease_owner) + self.assertIsNone(job.error_kind) + + def test_expired_detail_lease_is_recovered(self) -> None: + self.persistence.upsert_car(self._search_record("mobile.de:expired")) + self.persistence.enqueue_detail_jobs(["mobile.de:expired"]) + lease = self.persistence.reserve_detail_jobs(limit=1)[0] + + with self.persistence.session_scope() as session: + job = session.execute(select(DetailJob).where(DetailJob.listing_id == "expired")).scalar_one() + job.lease_expires_at = datetime.now(timezone.utc) - timedelta(minutes=1) + + recovered = self.persistence.reserve_detail_jobs(limit=1, listing_ids=["expired"]) + + self.assertEqual(len(recovered), 1) + self.assertNotEqual(recovered[0]["lease_owner"], lease["lease_owner"]) + self.assertEqual(recovered[0]["attempts"], 2) + + def test_start_sync_run_keeps_parallel_active_runs_running(self) -> None: first_run_id = self.persistence.start_sync_run("lane-a") second_run_id = self.persistence.start_sync_run("lane-b") @@ -260,9 +448,36 @@ class TestPersistenceServiceIntegration(unittest.TestCase): first = session.get(SyncRun, first_run_id) second = session.get(SyncRun, second_run_id) + self.assertEqual(first.status, "running") + self.assertIsNone(first.finished_at) + self.assertEqual(second.status, "running") + + def test_start_sync_run_marks_only_old_running_runs_as_failed(self) -> None: + first_run_id = self.persistence.start_sync_run("lane-a") + with self.persistence.session_scope() as session: + first = session.get(SyncRun, first_run_id) + first.started_at = datetime(2020, 1, 1, tzinfo=timezone.utc) + + self.persistence.start_sync_run("lane-b") + + with self.persistence.session_scope() as session: + first = session.get(SyncRun, first_run_id) self.assertEqual(first.status, "failed") self.assertIsNotNone(first.finished_at) - self.assertEqual(second.status, "running") + + def test_individual_upsert_fallback_reports_failed_records(self) -> None: + first = self._record("mobile.de:ok") + second = self._record("mobile.de:failed") + with patch.object( + self.persistence, + "upsert_car", + side_effect=[{"action": "inserted", "images_upserted": 0}, RuntimeError("db failure")], + ): + result = self.persistence._upsert_cars_individually([first, second]) + + self.assertEqual(result["inserted"], 1) + self.assertEqual(result["failed"], 1) + self.assertEqual(result["failed_origin_ids"], ["mobile.de:failed"]) def test_upsert_falls_back_to_origin_url_to_prevent_duplicates(self) -> None: first = self._record("OLD-ID") @@ -331,8 +546,28 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.assertEqual(sold_map["mobile.de:old"], (True, seen_at)) self.assertEqual(sold_map["mobile.de:current"], (False, None)) - def test_get_active_cars_batch_for_image_enrich_selects_low_image_active_cars(self) -> None: + def test_mark_sold_not_seen_since_limits_candidates_to_scope_brands(self) -> None: + seen_at = datetime(2026, 5, 5, tzinfo=timezone.utc) + bmw = self._record("mobile.de:bmw") + bmw.brand = "BMW" + bmw.last_seen_at = seen_at - timedelta(days=1) + toyota = self._record("mobile.de:toyota") + toyota.brand = "Toyota" + toyota.last_seen_at = seen_at - timedelta(days=1) + self.persistence.upsert_car(bmw) + self.persistence.upsert_car(toyota) + + marked = self.persistence.mark_sold_not_seen_since(seen_at, brands=("BMW",)) + + self.assertEqual(marked, 1) + with self.persistence.session_scope() as session: + cars = {car.origin_id: car for car in session.execute(select(Car)).scalars().all()} + self.assertTrue(cars[bmw.origin_id].is_sold) + self.assertFalse(cars[toyota.origin_id].is_sold) + + def test_get_active_cars_batch_for_image_enrich_selects_unfetched_galleries(self) -> None: low = self._record("mobile.de:low") + low.details_confirmed = False low.images = [ ImageRecord( fullres_image="https://img.classistatic.de/api/v1/mo-prod/images/low?rule=mo-640.jpg", @@ -357,6 +592,8 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.persistence.upsert_car(rich) self.persistence.upsert_car(sold) + self.persistence.replace_car_gallery(rich.origin_id, rich.images) + candidates = self.persistence.get_active_cars_batch_for_image_enrich(limit=10, max_existing_images=1) self.assertEqual([(origin_id, image_count) for _id, origin_id, _url, image_count in candidates], [("mobile.de:low", 1)]) @@ -364,10 +601,12 @@ class TestPersistenceServiceIntegration(unittest.TestCase): def test_image_enrich_selector_prefers_unfetched_then_oldest_stale(self) -> None: unfetched = self._record("mobile.de:unfetched") unfetched.details_confirmed = False + unfetched_new = self._record("mobile.de:unfetched-new") + unfetched_new.details_confirmed = False stale_old = self._record("mobile.de:stale-old") stale_new = self._record("mobile.de:stale-new") fresh = self._record("mobile.de:fresh") - for record in (unfetched, stale_old, stale_new, fresh): + for record in (unfetched, unfetched_new, stale_old, stale_new, fresh): record.images = [] self.persistence.upsert_car(record) @@ -377,9 +616,11 @@ class TestPersistenceServiceIntegration(unittest.TestCase): car.origin_id: car for car in session.execute(select(Car)).scalars().all() } - cars[stale_old.origin_id].details_fetched_at = now - timedelta(days=10) - cars[stale_new.origin_id].details_fetched_at = now - timedelta(days=8) - cars[fresh.origin_id].details_fetched_at = now - timedelta(days=1) + cars[unfetched.origin_id].first_seen_at = now - timedelta(hours=2) + cars[unfetched_new.origin_id].first_seen_at = now - timedelta(hours=1) + cars[stale_old.origin_id].gallery_fetched_at = now - timedelta(days=10) + cars[stale_new.origin_id].gallery_fetched_at = now - timedelta(days=8) + cars[fresh.origin_id].gallery_fetched_at = now - timedelta(days=1) candidates = self.persistence.get_active_cars_batch_for_image_enrich( limit=10, @@ -389,7 +630,7 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.assertEqual( [origin_id for _id, origin_id, _url, _count in candidates], - [unfetched.origin_id, stale_old.origin_id, stale_new.origin_id], + [unfetched_new.origin_id, unfetched.origin_id, stale_old.origin_id, stale_new.origin_id], ) self.assertEqual( self.persistence.get_active_cars_batch_for_image_enrich(