Add detail job storage

This commit is contained in:
qananasikq
2026-08-19 11:18:00 +03:00
parent f3329f76a5
commit 33b85ccecd
5 changed files with 832 additions and 265 deletions
+478 -255
View File
@@ -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
+28
View File
@@ -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)
+2
View File
@@ -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)