From 1891f84291e7ab45cafbfa3f53194caff4b2ce4b Mon Sep 17 00:00:00 2001 From: qananasikq Date: Thu, 20 Aug 2026 10:43:00 +0300 Subject: [PATCH] Add gallery fetch status --- ...py => 004_add_car_seen_sold_timestamps.py} | 6 +- alembic/versions/004_add_ingestion_tables.py | 17 -- ...ed_at.py => 005_add_details_fetched_at.py} | 4 +- alembic/versions/006_runtime_storage.py | 35 +++ alembic/versions/007_runtime_storage.py | 73 ----- mobilede_scraper/storage/db.py | 255 +----------------- mobilede_scraper/storage/enums.py | 38 --- mobilede_scraper/storage/models.py | 27 -- tests/test_db.py | 91 +------ 9 files changed, 54 insertions(+), 492 deletions(-) rename alembic/versions/{005_add_car_seen_sold_timestamps.py => 004_add_car_seen_sold_timestamps.py} (86%) delete mode 100644 alembic/versions/004_add_ingestion_tables.py rename alembic/versions/{006_add_details_fetched_at.py => 005_add_details_fetched_at.py} (85%) create mode 100644 alembic/versions/006_runtime_storage.py delete mode 100644 alembic/versions/007_runtime_storage.py diff --git a/alembic/versions/005_add_car_seen_sold_timestamps.py b/alembic/versions/004_add_car_seen_sold_timestamps.py similarity index 86% rename from alembic/versions/005_add_car_seen_sold_timestamps.py rename to alembic/versions/004_add_car_seen_sold_timestamps.py index 993e1d1..b9ad4f3 100644 --- a/alembic/versions/005_add_car_seen_sold_timestamps.py +++ b/alembic/versions/004_add_car_seen_sold_timestamps.py @@ -3,8 +3,8 @@ from typing import Sequence, Union import sqlalchemy as sa from alembic import op -revision: str = "005_add_car_seen_sold_timestamps" -down_revision: Union[str, None] = "004_add_ingestion_tables" +revision: str = "004_add_car_seen_sold_timestamps" +down_revision: Union[str, None] = "003_add_composite_indexes" branch_labels: Union[str, Sequence[str], None] = None depends_on: Union[str, Sequence[str], None] = None @@ -28,4 +28,4 @@ def downgrade() -> None: op.drop_index("ix_MOBILEDE_cars_sold_at", table_name="MOBILEDE_cars") op.drop_index("ix_MOBILEDE_cars_first_seen_at", table_name="MOBILEDE_cars") op.drop_column("MOBILEDE_cars", "sold_at") - op.drop_column("MOBILEDE_cars", "first_seen_at") + op.drop_column("MOBILEDE_cars", "first_seen_at") \ No newline at end of file diff --git a/alembic/versions/004_add_ingestion_tables.py b/alembic/versions/004_add_ingestion_tables.py deleted file mode 100644 index a204e00..0000000 --- a/alembic/versions/004_add_ingestion_tables.py +++ /dev/null @@ -1,17 +0,0 @@ -# Заглушка миграции (таблицы ingestion пока не требуются) -from typing import Sequence, Union - -from alembic import op - -revision: str = "004_add_ingestion_tables" -down_revision: Union[str, None] = "003_add_composite_indexes" -branch_labels: Union[str, Sequence[str], None] = None -depends_on: Union[str, Sequence[str], None] = None - - -def upgrade() -> None: - pass - - -def downgrade() -> None: - pass diff --git a/alembic/versions/006_add_details_fetched_at.py b/alembic/versions/005_add_details_fetched_at.py similarity index 85% rename from alembic/versions/006_add_details_fetched_at.py rename to alembic/versions/005_add_details_fetched_at.py index 3e9dc4a..aee15a4 100644 --- a/alembic/versions/006_add_details_fetched_at.py +++ b/alembic/versions/005_add_details_fetched_at.py @@ -3,8 +3,8 @@ from typing import Sequence, Union import sqlalchemy as sa from alembic import op -revision: str = "006_add_details_fetched_at" -down_revision: Union[str, None] = "005_add_car_seen_sold_timestamps" +revision: str = "005_add_details_fetched_at" +down_revision: Union[str, None] = "004_add_car_seen_sold_timestamps" branch_labels: Union[str, Sequence[str], None] = None depends_on: Union[str, Sequence[str], None] = None diff --git a/alembic/versions/006_runtime_storage.py b/alembic/versions/006_runtime_storage.py new file mode 100644 index 0000000..b702837 --- /dev/null +++ b/alembic/versions/006_runtime_storage.py @@ -0,0 +1,35 @@ +from typing import Sequence, Union + +import sqlalchemy as sa +from alembic import op + +revision: str = "006_runtime_storage" +down_revision: Union[str, None] = "005_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.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.execute('DROP INDEX IF EXISTS "ix_MOBILEDE_cars_active_last_seen_mobile_de"') \ No newline at end of file diff --git a/alembic/versions/007_runtime_storage.py b/alembic/versions/007_runtime_storage.py deleted file mode 100644 index a3e22e4..0000000 --- a/alembic/versions/007_runtime_storage.py +++ /dev/null @@ -1,73 +0,0 @@ -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 8fe4b16..9bf000f 100644 --- a/mobilede_scraper/storage/db.py +++ b/mobilede_scraper/storage/db.py @@ -1,17 +1,16 @@ import logging import os import re -import uuid from contextlib import contextmanager from datetime import datetime, timedelta, timezone -from typing import Any, Iterator +from typing import 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 .models import Base, Car, Image, SyncRun from .schemas import CarRecord logger = logging.getLogger("MOBILEDE_scraper.db") @@ -53,14 +52,6 @@ MOBILEDE_MATERIAL_CHANGE_FIELDS = ( _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"}, @@ -866,7 +857,6 @@ class PersistenceService: 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: @@ -933,242 +923,21 @@ class PersistenceService: "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 + def is_gallery_fetched(self, origin_id: str) -> bool: + 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: - 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, - ) + return bool(session.execute( + select(Car.id) + .where(Car.origin_id.in_(candidate_origin_ids)) + .where(Car.gallery_fetched_at.is_not(None)) + .limit(1) + ).scalar_one_or_none()) def mark_cars_sold_by_ids(self, car_ids: list[int], *, sold_at: datetime | None = None) -> int: if not car_ids: diff --git a/mobilede_scraper/storage/enums.py b/mobilede_scraper/storage/enums.py index 4357d2b..90bed7a 100644 --- a/mobilede_scraper/storage/enums.py +++ b/mobilede_scraper/storage/enums.py @@ -16,44 +16,6 @@ BODY_TYPE_ENUM_VALUES = ( "OTHER", "NA", ) -COUNTRY_ENUM_VALUES = ( - "AT", - "BE", - "BG", - "CA", - "CH", - "CY", - "CZ", - "DE", - "DK", - "EE", - "ES", - "FI", - "FR", - "GB", - "GR", - "HR", - "HU", - "IE", - "IT", - "JP", - "KR", - "LI", - "LT", - "LU", - "LV", - "MT", - "NA", - "NL", - "NO", - "PL", - "PT", - "RO", - "SE", - "SI", - "SK", - "US", -) ORIGIN_ENUM_VALUES = ( "MOBILEDE", "MOBILE_DE", diff --git a/mobilede_scraper/storage/models.py b/mobilede_scraper/storage/models.py index f0137d7..4f8bde5 100644 --- a/mobilede_scraper/storage/models.py +++ b/mobilede_scraper/storage/models.py @@ -71,33 +71,6 @@ 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/tests/test_db.py b/tests/test_db.py index 256d949..6b3d7d1 100644 --- a/tests/test_db.py +++ b/tests/test_db.py @@ -10,7 +10,7 @@ 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, DetailJob, Image, SyncRun +from mobilede_scraper.storage.models import Car, Image, SyncRun from mobilede_scraper.storage.schemas import CarRecord, ImageRecord @@ -353,91 +353,6 @@ class TestPersistenceServiceIntegration(unittest.TestCase): 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") @@ -594,7 +509,7 @@ class TestPersistenceServiceIntegration(unittest.TestCase): 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) + candidates = self.persistence.get_active_cars_batch_for_image_enrich(limit=10) self.assertEqual([(origin_id, image_count) for _id, origin_id, _url, image_count in candidates], [("mobile.de:low", 1)]) @@ -624,7 +539,6 @@ class TestPersistenceServiceIntegration(unittest.TestCase): candidates = self.persistence.get_active_cars_batch_for_image_enrich( limit=10, - max_existing_images=1, stale_before=now - timedelta(days=7), ) @@ -635,7 +549,6 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.assertEqual( self.persistence.get_active_cars_batch_for_image_enrich( limit=0, - max_existing_images=1, stale_before=now, ), [],