From 2cebb783160aa944d6497e98182906f25cfd935c Mon Sep 17 00:00:00 2001 From: qananasikq Date: Mon, 11 May 2026 12:46:07 +0300 Subject: [PATCH] refresh api and tests --- mobilede_scraper/api/app.py | 7 +- mobilede_scraper/api/deps.py | 2 +- mobilede_scraper/api/routes/__init__.py | 8 +- mobilede_scraper/api/routes/cars.py | 10 +- mobilede_scraper/api/routes/health.py | 4 +- mobilede_scraper/api/routes/tasks.py | 34 ++- tests/test_db.py | 30 +++ tests/test_mappers.py | 296 +++++++++++++++++++++++- tests/test_mobilede_scraper.py | 71 +++++- tests/test_self_heal.py | 2 +- tests/test_worker_runtime_tasks.py | 183 +++++++++++++++ uv.lock | 2 +- 12 files changed, 621 insertions(+), 28 deletions(-) diff --git a/mobilede_scraper/api/app.py b/mobilede_scraper/api/app.py index ac2215c..1a9e86d 100644 --- a/mobilede_scraper/api/app.py +++ b/mobilede_scraper/api/app.py @@ -1,4 +1,4 @@ -# Создание FastAPI-приложения и настройка его жизненного цикла. +# FastAPI-приложение и его жизненный цикл. from contextlib import asynccontextmanager @@ -11,8 +11,7 @@ from .routes import cars, health, tasks @asynccontextmanager async def lifespan(app: FastAPI): - # Жизненный цикл. - # Таблицы создаются через Alembic-миграции (сервис migrate). + # Таблицы создаются миграциями Alembic. yield @@ -36,5 +35,5 @@ def create_app(settings: Settings | None = None) -> FastAPI: return app -# Экземпляр приложения для запуска через uvicorn. +# Точка входа для uvicorn. app = create_app() diff --git a/mobilede_scraper/api/deps.py b/mobilede_scraper/api/deps.py index 86b4bc3..e9abf43 100644 --- a/mobilede_scraper/api/deps.py +++ b/mobilede_scraper/api/deps.py @@ -1,4 +1,4 @@ -# Вспомогательные зависимости для FastAPI-роутов. +# Зависимости FastAPI. from fastapi import Request diff --git a/mobilede_scraper/api/routes/__init__.py b/mobilede_scraper/api/routes/__init__.py index a9cf553..49c9739 100644 --- a/mobilede_scraper/api/routes/__init__.py +++ b/mobilede_scraper/api/routes/__init__.py @@ -1,7 +1,3 @@ -"""API route package. - -Routers are defined in dedicated modules (`cars`, `health`, `tasks`). This -package intentionally does not create routes to avoid duplicate task endpoints -and accidental imports of legacy parser channels during FastAPI startup. -""" +# Пакет роутов API. +# Здесь только импорт модулей роутов без побочных эффектов. diff --git a/mobilede_scraper/api/routes/cars.py b/mobilede_scraper/api/routes/cars.py index 7195d48..749e281 100644 --- a/mobilede_scraper/api/routes/cars.py +++ b/mobilede_scraper/api/routes/cars.py @@ -1,4 +1,4 @@ -# Роуты для просмотра автомобилей и агрегированной статистики. +# Роуты для авто и статистики. from fastapi import APIRouter, Depends, HTTPException, Query from sqlalchemy import func, select @@ -22,7 +22,7 @@ def list_cars( is_sold: bool | None = None, persistence: PersistenceService = Depends(get_persistence), ): - # Список автомобилей с пагинацией и фильтрами + # Список авто с фильтрами. with persistence.session_scope() as session: query = select(Car) @@ -59,7 +59,7 @@ def get_car( car_id: int, persistence: PersistenceService = Depends(get_persistence), ): - # Детальная информация об автомобиле с изображениями + # Карточка авто с фото. with persistence.session_scope() as session: car = session.get(Car, car_id) if car is None: @@ -72,7 +72,7 @@ def get_car_by_origin( origin_id: str, persistence: PersistenceService = Depends(get_persistence), ): - # Поиск автомобиля по origin_id + # Поиск по origin_id. with persistence.session_scope() as session: car = session.execute( select(Car).where(Car.origin_id == origin_id) @@ -84,7 +84,7 @@ def get_car_by_origin( @router.get("/stats") def get_stats(persistence: PersistenceService = Depends(get_persistence)): - # Общая статистика по БД + # Общая статистика. with persistence.session_scope() as session: total_cars = session.execute(select(func.count(Car.id))).scalar() or 0 total_images = session.execute(select(func.count(Image.id))).scalar() or 0 diff --git a/mobilede_scraper/api/routes/health.py b/mobilede_scraper/api/routes/health.py index d10c942..ae52e91 100644 --- a/mobilede_scraper/api/routes/health.py +++ b/mobilede_scraper/api/routes/health.py @@ -1,4 +1,4 @@ -# Роут проверки доступности сервиса и соединения с БД. +# Проверка API и БД. import logging @@ -14,7 +14,7 @@ logger = logging.getLogger("MOBILEDE_scraper.api.health") @router.get("/health") def health_check(persistence: PersistenceService = Depends(get_persistence)): - # Проверка API и БД + # Проверяем доступность БД. db_ok = False try: with persistence.session_scope() as session: diff --git a/mobilede_scraper/api/routes/tasks.py b/mobilede_scraper/api/routes/tasks.py index 0b1d02f..f9a134b 100644 --- a/mobilede_scraper/api/routes/tasks.py +++ b/mobilede_scraper/api/routes/tasks.py @@ -1,4 +1,4 @@ -# Роуты запуска задач синхронизации и просмотра истории sync-runs. +# Роуты запуска задач и истории sync-run. import json @@ -11,8 +11,14 @@ from ...core.config import Settings from ..deps import get_persistence from ...storage.db import PersistenceService from ...storage.models import SyncRun -from ...worker.celery_app import MOBILEDE_SYNC_QUEUE, celery_app -from ...worker.tasks import mobilede_sync_detail_task, mobilede_sync_runtime_segments_task, mobilede_sync_search_task +from ...worker.celery_app import celery_app +from ...worker.constants import MOBILEDE_SYNC_QUEUE +from ...worker.tasks import ( + mobilede_enrich_images_batch_task, + mobilede_sync_detail_task, + mobilede_sync_runtime_segments_task, + mobilede_sync_search_task, +) router = APIRouter() @@ -38,6 +44,13 @@ class MobileDeSyncDetailRequest(BaseModel): lane: str = "mobile_de_cars" +class MobileDeEnrichImagesRequest(BaseModel): + limit: int = 50 + lane: str = "mobile_de_cars" + max_existing_images: int = 1 + delay_seconds: float = 1.0 + + class MobileDeRuntimeSegmentsRequest(BaseModel): lane: str = "mobile_de_cars" delay_seconds: float = 0.7 @@ -72,6 +85,19 @@ def start_mobilede_sync_detail(body: MobileDeSyncDetailRequest): } +@router.post("/mobilede/tasks/enrich-images") +def start_mobilede_enrich_images(body: MobileDeEnrichImagesRequest): + result = mobilede_enrich_images_batch_task.apply_async( + kwargs=body.model_dump(), + queue=MOBILEDE_SYNC_QUEUE, + ) + return { + "task_id": result.id, + "status": "queued", + "queue": MOBILEDE_SYNC_QUEUE, + } + + @router.post("/mobilede/tasks/sync-runtime-segments") def start_mobilede_runtime_segments(body: MobileDeRuntimeSegmentsRequest): result = mobilede_sync_runtime_segments_task.apply_async( @@ -141,7 +167,7 @@ def list_sync_runs( per_page: int = Query(20, ge=1, le=100), persistence: PersistenceService = Depends(get_persistence), ): - # История запусков синхронизации + # История запусков. with persistence.session_scope() as session: total = session.execute(select(func.count(SyncRun.id))).scalar() or 0 offset = (page - 1) * per_page diff --git a/tests/test_db.py b/tests/test_db.py index b3f8cb9..b344043 100644 --- a/tests/test_db.py +++ b/tests/test_db.py @@ -182,6 +182,36 @@ 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: + low = self._record("mobile.de:low") + low.images = [ + ImageRecord( + fullres_image="https://img.classistatic.de/api/v1/mo-prod/images/low?rule=mo-640.jpg", + preview_image="https://img.classistatic.de/api/v1/mo-prod/images/low?rule=mo-200.jpg", + order_index=0, + ) + ] + rich = self._record("mobile.de:rich") + rich.images = [ + ImageRecord( + fullres_image=f"https://img.classistatic.de/api/v1/mo-prod/images/rich-{idx}?rule=mo-640.jpg", + preview_image=f"https://img.classistatic.de/api/v1/mo-prod/images/rich-{idx}?rule=mo-200.jpg", + order_index=idx, + ) + for idx in range(3) + ] + sold = self._record("mobile.de:sold") + sold.images = [] + sold.is_sold = True + + self.persistence.upsert_car(low) + self.persistence.upsert_car(rich) + self.persistence.upsert_car(sold) + + 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)]) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_mappers.py b/tests/test_mappers.py index 5c0652d..1b1f578 100644 --- a/tests/test_mappers.py +++ b/tests/test_mappers.py @@ -27,8 +27,51 @@ class TestMobileDeMapper(unittest.TestCase): ) self.assertEqual(len(record.images), 2) - self.assertEqual(record.images[0].fullres_image, urls[0]) - self.assertEqual(record.images[1].fullres_image, "https://img.classistatic.de/api/v1/mo-prod/images/2") + self.assertEqual(record.images[0].fullres_image, f"{urls[0]}?rule=mo-640.jpg") + self.assertEqual(record.images[1].fullres_image, "https://img.classistatic.de/api/v1/mo-prod/images/2?rule=mo-640.jpg") + + def test_extracts_nested_gallery_images_without_detail_fetch(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="124", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=124", + title="Honda Civic", + raw={ + "image": "https://img.classistatic.de/api/v1/mo-prod/images/main", + "images": [{"ref": "img.classistatic.de/api/v1/mo-prod/images/ref-1"}], + "mediaGallery": [ + {"uri": "//img.classistatic.de/api/v1/mo-prod/images/gallery-1"}, + {"picture": {"src": "https://img.classistatic.de/api/v1/mo-prod/images/gallery-2"}}, + ], + "trackingUrl": "https://example.test/not-an-image", + }, + ) + ) + + self.assertEqual( + [image.fullres_image for image in record.images], + [ + "https://img.classistatic.de/api/v1/mo-prod/images/main?rule=mo-640.jpg", + "https://img.classistatic.de/api/v1/mo-prod/images/ref-1?rule=mo-640.jpg", + "https://img.classistatic.de/api/v1/mo-prod/images/gallery-1?rule=mo-640.jpg", + "https://img.classistatic.de/api/v1/mo-prod/images/gallery-2?rule=mo-640.jpg", + ], + ) + + def test_keeps_existing_mobilede_image_rule(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="125", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=125", + title="BMW 320", + raw={"images": [{"uri": "img.classistatic.de/api/v1/mo-prod/images/abc?rule=mo-1024.jpg"}]}, + ) + ) + + self.assertEqual( + [image.fullres_image for image in record.images], + ["https://img.classistatic.de/api/v1/mo-prod/images/abc?rule=mo-1024.jpg"], + ) def test_origin_id_uses_canonical_prefix(self) -> None: record = self.mapper.listing_to_car_record( @@ -84,6 +127,255 @@ class TestMobileDeMapper(unittest.TestCase): self.assertEqual(record.year, 2019) self.assertEqual(record.mileage, 50100) + def test_listing_uses_mobilede_attr_fallbacks(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="1001", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1001", + title="BMW X1 xDrive", + subtitle="Sport", + raw={ + "attr": { + "yc": "2016", + "cn": "DE", + "ecol": "Серебряный", + "c": "OffRoad", + "cc": "1 998 ccm", + } + }, + ) + ) + + self.assertEqual(record.year, 2016) + self.assertEqual(record.country, "DE") + self.assertEqual(record.color, "silver") + self.assertEqual(record.body_type, "SUV") + self.assertEqual(record.engine_volume, 1998) + self.assertEqual(record.drive, "4WD") + + def test_listing_body_and_color_normalization_is_broader(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="1002", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1002", + title="BMW 320 Touring", + subtitle="EstateCar", + raw={"attr": {"ecol": "Коричневый"}}, + ) + ) + + self.assertEqual(record.body_type, "STATION_WAGON") + self.assertEqual(record.color, "brown") + + def test_listing_maps_available_search_payload_fields(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="1003", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1003", + title="BMW 120 Frontantrieb", + raw={ + "price": {"grs": {"amount": "12250"}}, + "contact": {"country": "IT"}, + "priceRating": {"rating": "GOOD_PRICE", "ratingLabel": "Good price"}, + "hasDamage": True, + "isConditionNew": True, + "cubicCapacity": "1 995 ccm", + "color": "Black Metallic", + "attr": { + "yc": "2020", + "ml": "55 000 km", + "c": "SmallCar", + "tr": "Manual", + "pvo": "1", + }, + }, + ) + ) + + self.assertEqual(record.year, 2020) + self.assertEqual(record.price, 12250) + self.assertEqual(record.mileage, 55000) + self.assertEqual(record.country, "IT") + self.assertEqual(record.color, "black") + self.assertEqual(record.drive, "FWD") + self.assertEqual(record.gearbox, "MT") + self.assertEqual(record.body_type, "HATCHBACK") + self.assertEqual(record.engine_volume, 1995) + self.assertTrue(record.one_owner) + self.assertTrue(record.new_car) + self.assertTrue(record.is_damaged) + self.assertTrue(record.repair_history) + self.assertEqual(record.evaluation, "GOOD_PRICE") + + def test_drive_mapping_uses_specific_markers_only(self) -> None: + allrad = self.mapper.listing_to_car_record( + MobileDeListing( + id="1004", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1004", + title="BMW X3 Allrad", + ) + ) + unrelated = self.mapper.listing_to_car_record( + MobileDeListing( + id="1005", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1005", + title="BMW 320 All inclusive", + ) + ) + + self.assertEqual(allrad.drive, "4WD") + self.assertIsNone(unrelated.drive) + + def test_search_mapping_uses_country_codes_and_body_fallbacks(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="1006", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1006", + title="BMW X3", + subtitle="M40 d Facelift 1 Hand Scheckheft BMW", + raw={ + "category": "Иное", + "attr": { + "cn": "NL", + "c": "OtherCar", + "cc": "2 993 ccm", + }, + }, + ) + ) + + self.assertEqual(record.country, "NL") + self.assertEqual(record.body_type, "SUV") + self.assertEqual(record.engine_volume, 2993) + + def test_search_mapping_infers_explicit_decimal_engine_volume(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="1007", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1007", + title="Honda Civic", + subtitle="1.6 Automatik", + raw={"attr": {"cn": "AT", "c": "SmallCar"}}, + ) + ) + + self.assertEqual(record.country, "AT") + self.assertEqual(record.body_type, "HATCHBACK") + self.assertEqual(record.engine_volume, 1600) + + def test_search_mapping_does_not_treat_plain_year_text_as_engine_volume(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="1009", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1009", + title="Honda Civic", + subtitle="1983", + raw={"attr": {"cn": "SK", "c": "SmallCar"}}, + ) + ) + + self.assertIsNone(record.engine_volume) + + def test_search_mapping_extracts_ccm_without_concatenating_all_title_digits(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="1010", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1010", + title="BMW 320", + subtitle="E 46 320i 2.200 ccm 170 PS 2005", + raw={"attr": {"cn": "DE", "c": "Cabrio"}}, + ) + ) + + self.assertEqual(record.engine_volume, 2200) + + def test_search_mapping_reads_nested_safe_field_aliases(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="1008", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1008", + title="BMW 218", + raw={ + "vehicle": { + "specs": { + "exteriorColor": "Silver", + "bodyType": "Cabrio", + "cubicCapacity": "1998", + "wheelDrive": "Allrad", + "transmission": "Automatic", + } + } + }, + ) + ) + + self.assertEqual(record.color, "silver") + self.assertEqual(record.body_type, "OPEN") + self.assertEqual(record.engine_volume, 1998) + self.assertEqual(record.drive, "4WD") + self.assertEqual(record.gearbox, "AT") + + def test_listing_mapping_supports_attr_aliases_and_liter_engine_text(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="1011", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1011", + title="BMW 320 Touring", + subtitle="2.0 l xDrive", + raw={ + "attr": { + "exteriorColor": "Skyscraper Grey Metallic", + "wheelDrive": "Antrieb vorne", + "engineDisplacement": "2.0 l", + }, + "driveType": "xDrive", + }, + ) + ) + + self.assertEqual(record.color, "gray") + self.assertEqual(record.drive, "FWD") + self.assertEqual(record.engine_volume, 2000) + + def test_detail_mapping_uses_aliases_for_color_drive_and_engine(self) -> None: + record = self.mapper.detail_to_car_record( + "1012", + { + "shortTitle": "BMW X3", + "subTitle": "xDrive30d", + "make": {"localized": "BMW"}, + "model": {"localized": "X3"}, + "price": {"grs": {"amount": "31990", "currency": "EUR"}}, + "contact": {"countryCode": "DE"}, + "attributes": [ + {"tag": "exteriorColor", "value": "Brooklyn Grey"}, + {"tag": "driveType", "value": "All wheel drive"}, + {"tag": "engineDisplacement", "value": "2993 cc"}, + {"tag": "mileage", "value": "145 000 km"}, + {"tag": "year", "value": "2019"}, + {"tag": "bodyType", "value": "Geländewagen"}, + ], + }, + ) + + self.assertEqual(record.color, "gray") + self.assertEqual(record.drive, "4WD") + self.assertEqual(record.engine_volume, 2993) + self.assertEqual(record.body_type, "SUV") + self.assertEqual(record.year, 2019) + + def test_drive_prefers_explicit_attr_marker_over_title_noise(self) -> None: + record = self.mapper.listing_to_car_record( + MobileDeListing( + id="1013", + url="https://suchen.mobile.de/fahrzeuge/details.html?id=1013", + title="BMW 320 xDrive", + raw={"attr": {"wheelDrive": "Front-wheel drive"}}, + ) + ) + + self.assertEqual(record.drive, "FWD") + if __name__ == "__main__": unittest.main() diff --git a/tests/test_mobilede_scraper.py b/tests/test_mobilede_scraper.py index eb17f70..f4e433c 100644 --- a/tests/test_mobilede_scraper.py +++ b/tests/test_mobilede_scraper.py @@ -2,16 +2,18 @@ from __future__ import annotations import unittest from datetime import datetime, timezone +from unittest.mock import patch from mobilede_scraper.mobile_de.models import MobileDeListing, MobileDeSearchPage from mobilede_scraper.mobile_de.scraper import MobileDeScraper class _FakePersistence: - def __init__(self) -> None: + def __init__(self, existing_ids: set[str] | None = None) -> None: self.upsert_calls: list[list[str]] = [] self.upsert_records: list[object] = [] self.finish_payload: dict[str, object] | None = None + self.existing_ids = existing_ids or set() def create_tables(self) -> None: return None @@ -20,7 +22,7 @@ class _FakePersistence: return 1 def get_existing_origin_ids(self, origin_ids: list[str]) -> set[str]: - return set() + return {origin_id for origin_id in origin_ids if origin_id in self.existing_ids} def upsert_cars_batch(self, records): self.upsert_calls.append([record.origin_id for record in records]) @@ -38,6 +40,7 @@ class _FakePersistence: class _FakeClient: def __init__(self, pages: list[MobileDeSearchPage]) -> None: self.pages = pages + self.detail_calls: list[str] = [] def build_make_model_param(self, make_id: str, model_id: str | None) -> str: return make_id if not model_id else f"{make_id};{model_id}" @@ -46,6 +49,26 @@ class _FakeClient: del kwargs yield from self.pages + def fetch_detail(self, listing_id: str) -> dict[str, object]: + self.detail_calls.append(str(listing_id)) + return { + "shortTitle": "BMW X1", + "subTitle": "xDrive20i", + "make": {"localized": "BMW"}, + "model": {"localized": "X1"}, + "price": {"grs": {"amount": "19999", "currency": "EUR"}}, + "contact": {"countryCode": "DE"}, + "attributes": [ + {"tag": "firstRegistration", "value": "05/2016"}, + {"tag": "mileage", "value": "71 500 km"}, + {"tag": "category", "value": "OffRoad"}, + {"tag": "color", "value": "Серый"}, + {"tag": "wheelDrive", "value": "xDrive"}, + {"tag": "cubicCapacity", "value": "1 998 ccm"}, + ], + "images": ["https://img.example.test/1.jpg"], + } + class TestMobileDeScraperStreamingSync(unittest.TestCase): def test_sync_search_upserts_each_page_during_run(self) -> None: @@ -110,6 +133,50 @@ class TestMobileDeScraperStreamingSync(unittest.TestCase): self.assertTrue(all(record.is_sold is False for record in persistence.upsert_records)) self.assertTrue(all(record.sold_at is None for record in persistence.upsert_records)) + def test_sync_search_selectively_enriches_new_records_with_missing_fields(self) -> None: + pages = [ + MobileDeSearchPage( + url="https://example.test/page-1", + page_number=1, + total_results=2, + listings=[ + MobileDeListing( + id="1", + url="https://example.test/1", + title="BMW X1", + subtitle="xDrive20i", + raw={"attr": {"fr": "", "ml": "", "c": "", "ecol": "", "cc": ""}}, + ), + MobileDeListing( + id="2", + url="https://example.test/2", + title="BMW 320", + raw={"attr": {"fr": "09/2016", "ml": "115 000 km", "c": "EstateCar", "ecol": "Серый", "cc": "1 998 ccm"}}, + ), + ], + ), + ] + client = _FakeClient(pages) + persistence = _FakePersistence(existing_ids={"mobile.de:2"}) + scraper = MobileDeScraper(client=client, persistence=persistence) + + with ( + patch("mobilede_scraper.mobile_de.scraper.MOBILEDE_SELECTIVE_DETAIL_ENRICH_ENABLED", True), + patch("mobilede_scraper.mobile_de.scraper.MOBILEDE_SELECTIVE_DETAIL_ENRICH_MAX_PER_RUN", 1), + patch("mobilede_scraper.mobile_de.scraper.MOBILEDE_SELECTIVE_DETAIL_ENRICH_MAX_PER_PAGE", 1), + ): + result = scraper.sync_search(max_pages=1) + + self.assertEqual(client.detail_calls, ["1"]) + self.assertEqual(result["detail_enriched"], 1) + self.assertEqual(result["detail_enrich_failed"], 0) + record_by_id = {record.origin_id: record for record in persistence.upsert_records} + self.assertEqual(record_by_id["mobile.de:1"].drive, "4WD") + self.assertEqual(record_by_id["mobile.de:1"].engine_volume, 1998) + self.assertEqual(record_by_id["mobile.de:1"].body_type, "SUV") + self.assertEqual(record_by_id["mobile.de:1"].year, 2016) + self.assertEqual(record_by_id["mobile.de:2"].engine_volume, 1998) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_self_heal.py b/tests/test_self_heal.py index d0d19ca..3c073ed 100644 --- a/tests/test_self_heal.py +++ b/tests/test_self_heal.py @@ -65,7 +65,7 @@ class TestSelfHeal(unittest.TestCase): @patch("builtins.open") def test_kill_worker_process_sends_term_and_kill(self, open_mock, kill_mock, _sleep_mock) -> None: open_mock.return_value.__enter__.return_value.read.return_value = "123" - # SIGTERM -> process alive check (pid,0) -> SIGKILL + # SIGTERM -> проверка, что процесс жив (pid,0) -> SIGKILL kill_mock.side_effect = [None, None, None] self_heal._kill_worker_process() diff --git a/tests/test_worker_runtime_tasks.py b/tests/test_worker_runtime_tasks.py index 87bc454..42765ff 100644 --- a/tests/test_worker_runtime_tasks.py +++ b/tests/test_worker_runtime_tasks.py @@ -1,6 +1,7 @@ from __future__ import annotations from datetime import datetime, timezone +from types import SimpleNamespace import unittest from unittest.mock import MagicMock, patch @@ -12,6 +13,7 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): def __init__(self) -> None: self.store: dict[str, str] = {} self.sets: dict[str, set[str]] = {} + self.lists: dict[str, list[str]] = {} def get(self, key: str): return self.store.get(key) @@ -34,6 +36,26 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): bucket.add(member) return 1 if len(bucket) > before else 0 + def scard(self, key: str): + return len(self.sets.get(key, set())) + + def llen(self, key: str): + return len(self.lists.get(key, [])) + + def delete(self, *keys: str): + removed = 0 + for key in keys: + if key in self.store: + del self.store[key] + removed += 1 + if key in self.sets: + del self.sets[key] + removed += 1 + if key in self.lists: + del self.lists[key] + removed += 1 + return removed + def expire(self, key: str, ttl: int): # noqa: ARG002 return True @@ -203,6 +225,58 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): redis_client.set.assert_any_call(tasks.MOBILEDE_BOOTSTRAP_DONE_KEY, "1") redis_client.set.assert_any_call(tasks.MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY, "1", ex=30 * 24 * 60 * 60) + def test_force_full_scan_only_new_stays_full_pass_for_continuous_hourly_cycle(self) -> None: + redis_client = self._FakeRedis() + redis_client.store[tasks.MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY] = "10" + redis_client.store[tasks.MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY] = "10" + + effective_only_new = tasks._mobilede_force_full_scan_only_new( + True, + redis_client=redis_client, + continuous=True, + ) + + self.assertFalse(effective_only_new) + + def test_post_bootstrap_refresh_disabled_for_continuous_hourly_cycle(self) -> None: + redis_client = self._FakeRedis() + redis_client.store[tasks.MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY] = "10" + redis_client.store[tasks.MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY] = "10" + + post_bootstrap_refresh = tasks._mobilede_post_bootstrap_full_refresh( + redis_client, + False, + continuous=True, + ) + + self.assertFalse(post_bootstrap_refresh) + + def test_stalled_bootstrap_recovery_ignores_stale_global_progress_key(self) -> None: + redis_client = self._FakeRedis() + redis_client.store[tasks.MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY] = "13" + redis_client.store[tasks.MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY] = "261" + redis_client.store[tasks.GLOBAL_PROGRESS_TS_KEY] = "1" + redis_client.sets[tasks.MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY] = {"seg-a", "seg-b"} + + with patch("mobilede_scraper.worker.tasks.time.time", return_value=10_000): + recovered = tasks._mobilede_try_recover_stalled_bootstrap_queue(redis_client) + + self.assertTrue(recovered) + self.assertNotIn(tasks.MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY, redis_client.sets) + + def test_stalled_bootstrap_recovery_keeps_recent_progress(self) -> None: + redis_client = self._FakeRedis() + redis_client.store[tasks.MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY] = "13" + redis_client.store[tasks.MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY] = "261" + redis_client.store[tasks.GLOBAL_PROGRESS_TS_KEY] = "9950" + redis_client.sets[tasks.MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY] = {"seg-a", "seg-b"} + + with patch("mobilede_scraper.worker.tasks.time.time", return_value=10_000): + recovered = tasks._mobilede_try_recover_stalled_bootstrap_queue(redis_client) + + self.assertFalse(recovered) + self.assertIn(tasks.MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY, redis_client.sets) + def test_bootstrap_completion_forces_refresh_sold_finalize_when_progress_lags(self) -> None: redis_client = self._FakeRedis() started_at = datetime(2026, 5, 6, 10, 0, tzinfo=timezone.utc) @@ -244,6 +318,58 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertTrue(should_skip) + def test_queue_mobilede_bootstrap_recovery_can_force_requeue_existing_pending(self) -> None: + redis_client = self._FakeRedis() + redis_client.store[tasks.MOBILEDE_BOOTSTRAP_RECOVERY_PENDING_KEY] = "1" + + with patch.object(tasks.mobilede_sync_runtime_segments_task, "apply_async") as apply_async: + queued = tasks._queue_mobilede_bootstrap_recovery( + redis_client, + lane="mobile_de_cars", + delay_seconds=0.7, + use_cursor=True, + continuous=True, + segment_label="runtime_segments", + reason="waiting_for_active_bootstrap_tasks", + countdown=5, + force=True, + ) + + self.assertTrue(queued) + apply_async.assert_called_once() + self.assertEqual(redis_client.store[tasks.MOBILEDE_BOOTSTRAP_RECOVERY_PENDING_KEY], "1") + + def test_runtime_segments_task_requeues_pending_bootstrap_recovery_while_tasks_are_still_active(self) -> None: + redis_client = self._FakeRedis() + redis_client.store[tasks.MOBILEDE_BOOTSTRAP_RECOVERY_PENDING_KEY] = "1" + + runtime_config = SimpleNamespace(sync=SimpleNamespace(only_new=False)) + + with ( + patch.object(tasks, "_get_redis", return_value=redis_client), + patch.object(tasks, "_mobilede_try_recover_stalled_bootstrap_queue", return_value=False), + patch.object(tasks, "_get_cached_mobilede_runtime_segments", return_value=[{"make_id": "3500"}]), + patch.object(tasks, "Settings", return_value=SimpleNamespace(runtime_config_file="runtime_config.json")), + patch.object(tasks.RuntimeConfig, "from_file", return_value=runtime_config), + patch.object(tasks, "_mobilede_force_full_scan_only_new", return_value=False), + patch.object(tasks, "_mobilede_bootstrap_done", return_value=False), + patch.object(tasks, "_mobilede_post_bootstrap_full_refresh", return_value=False), + patch.object(tasks, "_mobilede_bootstrap_progress", return_value=(13, 261, 248)), + patch.object(tasks, "_mobilede_bootstrap_progress_snapshot", return_value=(13, 261, 248, 2)), + patch.object(tasks, "_has_recent_global_progress", return_value=False), + patch.object(tasks, "_queue_mobilede_bootstrap_recovery", return_value=True) as queue_recovery, + ): + result = tasks.mobilede_sync_runtime_segments_task.run( + lane="mobile_de_cars", + delay_seconds=0.7, + use_cursor=True, + continuous=True, + full_pass_repeat=False, + ) + + self.assertEqual(result["status"], "bootstrap_active") + queue_recovery.assert_called_once() + def test_exhausted_non_cursor_segment_window_is_skipped(self) -> None: exhausted = tasks._mobilede_segment_window_is_exhausted( {"make_id": "3500", "max_pages": 50}, @@ -336,6 +462,32 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertEqual(added, 0) redis_client.sadd.assert_not_called() + def test_late_overflow_children_are_not_queued_during_preplanned_bootstrap(self) -> None: + redis_client = MagicMock() + settings = MagicMock() + segment = { + "label": "Cars | ms=3500 | price=15001-20000", + "search_url": "https://www.mobile.de/ru/search.html?isSearchRequest=true&s=Car&vc=Car&ms=3500&p=15001:20000", + "make_id": "3500", + } + + with patch.object(tasks, "os") as os_mock: + os_mock.getenv.return_value = "true" + queued = tasks._queue_mobilede_overflow_child_segments( + redis_client, + settings, + parent_segment=segment, + lane="mobile_de_cars", + delay_seconds=0, + use_cursor=True, + only_new=False, + bootstrap_run=True, + refresh_cycle_id="cycle-a", + ) + + self.assertEqual(queued, 0) + redis_client.get.assert_not_called() + def test_pre_split_mileage_requires_explicit_flag(self) -> None: original_split = tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE try: @@ -390,6 +542,37 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertTrue(all("fr=" not in str(segment.get("search_url")) for segment in segments)) self.assertTrue(all("ml=" not in str(segment.get("search_url")) for segment in segments)) + def test_build_runtime_segments_splits_multi_make_filtered_search_url(self) -> None: + settings = MagicMock() + settings.listing.filtered_search_urls = [ + "https://www.mobile.de/ru/search.html?isSearchRequest=true&s=Car&vc=Car&ms=3500&ms=11000" + ] + + segments = tasks._build_mobilede_runtime_segments(settings) + + self.assertGreater(len(segments), 2) + self.assertEqual({str(segment.get("make_id")) for segment in segments}, {"3500", "11000"}) + self.assertTrue(all("pageNumber=1" in str(segment.get("search_url")) for segment in segments)) + self.assertTrue(all("ms=3500&ms=11000" not in str(segment.get("search_url")) for segment in segments)) + + def test_full_link_coverage_can_pre_split_every_segment_by_mileage(self) -> None: + settings = MagicMock() + settings.listing.filtered_search_urls = [ + "https://www.mobile.de/ru/search.html?isSearchRequest=true&s=Car&vc=Car&ms=111&ms=222" + ] + + original_split = tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE + try: + tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = True + segments = tasks._build_mobilede_runtime_segments(settings) + finally: + tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = original_split + + self.assertGreater(len(segments), 300) + self.assertTrue(all("ml=" in str(segment.get("search_url")) for segment in segments)) + self.assertEqual({str(segment.get("make_id")) for segment in segments}, {"111", "222"}) + self.assertTrue(all("ms=111&ms=222" not in str(segment.get("search_url")) for segment in segments)) + if __name__ == "__main__": unittest.main() diff --git a/uv.lock b/uv.lock index e2820a9..2b94cf9 100644 --- a/uv.lock +++ b/uv.lock @@ -342,7 +342,7 @@ wheels = [ ] [[package]] -name = "iaai-scraper" +name = "mobile-de-scraper" version = "0.1.0" source = { editable = "." } dependencies = [