Add gallery fetch status

This commit is contained in:
qananasikq
2026-08-20 10:43:00 +03:00
parent 7068b7e8fa
commit 1891f84291
9 changed files with 54 additions and 492 deletions
+12 -243
View File
@@ -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:
-38
View File
@@ -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",
-27
View File
@@ -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)