From ac2753146ea0b51688c24c9cfe52d23bf8994ae4 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Tue, 21 Apr 2026 23:50:31 +0300 Subject: [PATCH] add tests --- tests/conftest.py | 7 ++ tests/test_db.py | 127 ++++++++++++++++++++++ tests/test_mapper.py | 198 +++++++++++++++++++++++++++++++++++ tests/test_openlane.py | 129 +++++++++++++++++++++++ tests/test_scraper_images.py | 68 ++++++++++++ tests/test_utils.py | 24 +++++ tests/test_worker_tasks.py | 189 +++++++++++++++++++++++++++++++++ 7 files changed, 742 insertions(+) create mode 100644 tests/conftest.py create mode 100644 tests/test_db.py create mode 100644 tests/test_mapper.py create mode 100644 tests/test_openlane.py create mode 100644 tests/test_scraper_images.py create mode 100644 tests/test_utils.py create mode 100644 tests/test_worker_tasks.py diff --git a/tests/conftest.py b/tests/conftest.py new file mode 100644 index 0000000..b33f680 --- /dev/null +++ b/tests/conftest.py @@ -0,0 +1,7 @@ +from __future__ import annotations + +import logging + + +def pytest_configure() -> None: + logging.disable(logging.CRITICAL) diff --git a/tests/test_db.py b/tests/test_db.py new file mode 100644 index 0000000..4a910e5 --- /dev/null +++ b/tests/test_db.py @@ -0,0 +1,127 @@ +from __future__ import annotations + +import tempfile +import unittest +from pathlib import Path + +from sqlalchemy import select + +from openlane_scraper.core.config import Settings +from openlane_scraper.storage.db import PersistenceService +from openlane_scraper.storage.models import Car, Image, SyncRun +from openlane_scraper.storage.schemas import CarRecord, ImageRecord + + +class TestPersistenceServiceIntegration(unittest.TestCase): + def setUp(self) -> None: + self.tmp_dir = tempfile.TemporaryDirectory() + db_path = Path(self.tmp_dir.name) / "test.sqlite" + + self.settings = Settings() + self.settings.database.url = f"sqlite:///{db_path.as_posix()}" + self.settings.database.echo = False + + self.persistence = PersistenceService(self.settings) + self.persistence.create_tables() + + def tearDown(self) -> None: + self.persistence.engine.dispose() + self.tmp_dir.cleanup() + + @staticmethod + def _record(origin_id: str, *, price: int = 1000) -> CarRecord: + return CarRecord( + parser_id=f"openlane:{origin_id}", + brand="Toyota", + model="Camry", + year=2014, + price=price, + origin_url=f"https://app.openlane.com/vehicles/{origin_id}", + origin_id=f"openlane:{origin_id}", + slug=f"toyota-camry-{origin_id}", + images=[ + ImageRecord( + fullres_image="https://images.openlane.com/full/1.jpg", + preview_image="https://images.openlane.com/thumb/1.jpg", + order_index=0, + ) + ], + ) + + def test_insert_update_and_skip_flow(self) -> None: + first = self._record("777", price=1000) + inserted = self.persistence.upsert_car(first) + self.assertEqual(inserted["action"], "inserted") + self.assertEqual(inserted["images_upserted"], 1) + + same = self._record("777", price=1000) + updated_same = self.persistence.upsert_car(same) + self.assertEqual(updated_same["action"], "updated") + self.assertEqual(updated_same["images_upserted"], 1) + + changed = self._record("777", price=1500) + updated = self.persistence.upsert_car(changed) + self.assertEqual(updated["action"], "updated") + self.assertEqual(updated["images_upserted"], 1) + + with self.persistence.session_scope() as session: + cars = session.execute(select(Car)).scalars().all() + images = session.execute(select(Image)).scalars().all() + + self.assertEqual(len(cars), 1) + self.assertEqual(cars[0].price, 1500) + self.assertEqual(len(images), 1) + + def test_update_replaces_old_images(self) -> None: + first = self._record("888") + self.persistence.upsert_car(first) + + second = self._record("888") + second.images = [ + ImageRecord( + fullres_image="https://images.openlane.com/full/2.jpg", + preview_image="https://images.openlane.com/thumb/2.jpg", + order_index=0, + ) + ] + self.persistence.upsert_car(second) + + with self.persistence.session_scope() as session: + images = session.execute(select(Image)).scalars().all() + + self.assertEqual(len(images), 1) + self.assertIn("full/2", images[0].fullres_image) + + def test_start_sync_run_marks_stale_running_runs_as_failed(self) -> None: + first_run_id = self.persistence.start_sync_run("lane-a") + second_run_id = self.persistence.start_sync_run("lane-b") + + self.assertNotEqual(first_run_id, second_run_id) + + with self.persistence.session_scope() as session: + first = session.get(SyncRun, first_run_id) + second = session.get(SyncRun, second_run_id) + + self.assertEqual(first.status, "failed") + self.assertIsNotNone(first.finished_at) + self.assertEqual(second.status, "running") + + def test_upsert_falls_back_to_origin_url_to_prevent_duplicates(self) -> None: + first = self._record("OLD-ID") + first.origin_url = "https://app.openlane.com/vehicles/45089484" + self.persistence.upsert_car(first) + + second = self._record("NEW-ID") + second.origin_url = "https://app.openlane.com/vehicles/45089484" + result = self.persistence.upsert_car(second) + + self.assertEqual(result["action"], "updated") + with self.persistence.session_scope() as session: + cars = session.execute(select(Car)).scalars().all() + + self.assertEqual(len(cars), 1) + self.assertEqual(cars[0].origin_id, "openlane:NEW-ID") + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_mapper.py b/tests/test_mapper.py new file mode 100644 index 0000000..990363b --- /dev/null +++ b/tests/test_mapper.py @@ -0,0 +1,198 @@ +from __future__ import annotations + +import unittest + +from openlane_scraper.openlane.mapper import ( + map_openlane_record, + map_openlane_records, + _parse_mileage, + _parse_engine_volume_cc, + _generate_slug, +) + + +class TestMapOpenLaneRecord(unittest.TestCase): + def _sample_record(self, **overrides) -> dict: + base = { + "id": "12345", + "make": "Toyota", + "model": "Camry", + "year": 2020, + "mileage": 45000, + "price": 18500, + "color": "White", + "body_type": "sedan", + "transmission": "automatic", + "drivetrain": "fwd", + "vin": "1HGBH41JXMN109186", + "url": "https://app.openlane.com/vehicles/12345", + "images": [ + {"full": "https://img.openlane.com/1_full.jpg", "thumbnail": "https://img.openlane.com/1_thumb.jpg"}, + {"full": "https://img.openlane.com/2_full.jpg", "thumbnail": "https://img.openlane.com/2_thumb.jpg"}, + ], + } + base.update(overrides) + return base + + def test_basic_mapping(self) -> None: + record = self._sample_record() + result = map_openlane_record(record) + self.assertIsNotNone(result) + self.assertEqual(result.brand, "TOYOTA") + self.assertEqual(result.model, "CAMRY") + self.assertEqual(result.year, 2020) + self.assertEqual(result.price, 18500) + self.assertEqual(result.mileage, 45000) + self.assertEqual(result.color, "white") + self.assertEqual(result.body_type, "SEDAN") + self.assertEqual(result.gearbox, "AT") + self.assertEqual(result.drive, "FWD") + self.assertEqual(result.origin, "OPENLANE") + self.assertEqual(result.origin_id, "openlane:12345") + self.assertEqual(result.origin_url, "https://app.openlane.com/vehicles/12345") + self.assertEqual(result.evaluation, "1HGBH41JXMN109186") + self.assertEqual(len(result.images), 2) + self.assertEqual(result.images[0].fullres_image, "https://img.openlane.com/1_full.jpg") + self.assertEqual(result.images[0].preview_image, "https://img.openlane.com/1_thumb.jpg") + + def test_missing_id_returns_none(self) -> None: + record = {"make": "Toyota", "model": "Camry"} + self.assertIsNone(map_openlane_record(record)) + + def test_missing_brand_returns_none(self) -> None: + record = {"id": "123", "model": "Camry"} + self.assertIsNone(map_openlane_record(record)) + + def test_missing_model_returns_none(self) -> None: + record = {"id": "123", "make": "Toyota"} + self.assertIsNone(map_openlane_record(record)) + + def test_vin_as_id_fallback(self) -> None: + record = self._sample_record() + del record["id"] + record["vin"] = "WBAPH5C52BA012345" + result = map_openlane_record(record) + self.assertIsNotNone(result) + self.assertEqual(result.origin_id, "openlane:WBAPH5C52BA012345") + + def test_existing_openlane_prefix_is_normalized(self) -> None: + record = self._sample_record(id="openlane:777") + result = map_openlane_record(record) + self.assertIsNotNone(result) + self.assertEqual(result.origin_id, "openlane:777") + + def test_whitespace_id_is_normalized(self) -> None: + record = self._sample_record(id=" 888 ") + result = map_openlane_record(record) + self.assertIsNotNone(result) + self.assertEqual(result.origin_id, "openlane:888") + + def test_image_url_strings(self) -> None: + record = self._sample_record(images=["https://a.com/1.jpg", "https://a.com/2.jpg"]) + result = map_openlane_record(record) + self.assertEqual(len(result.images), 2) + self.assertEqual(result.images[0].fullres_image, "https://a.com/1.jpg") + + def test_image_deduplication(self) -> None: + record = self._sample_record(images=[ + {"full": "https://a.com/same.jpg"}, + {"full": "https://a.com/same.jpg"}, + {"full": "https://a.com/other.jpg"}, + ]) + result = map_openlane_record(record) + self.assertEqual(len(result.images), 2) + + def test_vehicle_images_api_format(self) -> None: + record = self._sample_record(images=[ + { + "url": "https://a.com/original.jpg", + "large_resolution_url": "https://a.com/large.jpg", + "low_resolution_url": "https://a.com/medium.jpg", + } + ]) + result = map_openlane_record(record) + self.assertEqual(len(result.images), 1) + self.assertEqual(result.images[0].fullres_image, "https://a.com/original.jpg") + self.assertEqual(result.images[0].preview_image, "https://a.com/medium.jpg") + + def test_default_url_generation(self) -> None: + record = self._sample_record() + del record["url"] + result = map_openlane_record(record) + self.assertEqual(result.origin_url, "https://app.openlane.com/vehicles/12345") + + def test_nested_vehicle_fields(self) -> None: + record = { + "id": "999", + "vehicle": { + "make": "BMW", + "model": "X5", + "year": 2022, + "mileage": 10000, + }, + } + result = map_openlane_record(record) + self.assertIsNotNone(result) + self.assertEqual(result.brand, "BMW") + self.assertEqual(result.model, "X5") + self.assertEqual(result.year, 2022) + + def test_map_records_filters_invalid(self) -> None: + records = [ + self._sample_record(id="1"), + {"no_id": True}, + self._sample_record(id="2"), + ] + results = map_openlane_records(records) + self.assertEqual(len(results), 2) + + def test_sale_type_mapping(self) -> None: + record = self._sample_record(sale_type="buy_now") + result = map_openlane_record(record) + self.assertEqual(result.selling_type, "STOCK") + + record2 = self._sample_record(sale_type="tender") + result2 = map_openlane_record(record2) + self.assertEqual(result2.selling_type, "TENDER") + + +class TestParseMileage(unittest.TestCase): + def test_integer(self) -> None: + self.assertEqual(_parse_mileage(45000), 45000) + + def test_string_with_commas(self) -> None: + self.assertEqual(_parse_mileage("45,000 miles"), 45000) + + def test_none_returns_zero(self) -> None: + self.assertEqual(_parse_mileage(None), 0) + + def test_float(self) -> None: + self.assertEqual(_parse_mileage(45000.5), 45000) + + +class TestParseEngineVolume(unittest.TestCase): + def test_liters_float(self) -> None: + self.assertEqual(_parse_engine_volume_cc(2.5), 2500) + + def test_cc_int(self) -> None: + self.assertEqual(_parse_engine_volume_cc(2500), 2500) + + def test_string_liters(self) -> None: + self.assertEqual(_parse_engine_volume_cc("2.5L"), 2500) + + def test_none(self) -> None: + self.assertIsNone(_parse_engine_volume_cc(None)) + + +class TestGenerateSlug(unittest.TestCase): + def test_basic(self) -> None: + slug = _generate_slug(2020, "TOYOTA", "CAMRY", "openlane:123") + self.assertEqual(slug, "2020-toyota-camry-openlane-123") + + def test_no_year(self) -> None: + slug = _generate_slug(None, "BMW", "X5", "openlane:456") + self.assertEqual(slug, "bmw-x5-openlane-456") + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_openlane.py b/tests/test_openlane.py new file mode 100644 index 0000000..47690c8 --- /dev/null +++ b/tests/test_openlane.py @@ -0,0 +1,129 @@ +from __future__ import annotations + +import json +import tempfile +import unittest +from pathlib import Path + +from openlane_scraper.core.config import OpenLaneConfig +from openlane_scraper.openlane.checkpoint import OpenLaneCheckpoint, OpenLaneCheckpointStore +from openlane_scraper.openlane.client import OpenLaneClient +from openlane_scraper.openlane.writer import OpenLaneResultWriter + + +class TestOpenLaneCheckpointStore(unittest.TestCase): + def test_save_and_load_checkpoint(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + path = Path(tmp_dir) / "checkpoint.json" + store = OpenLaneCheckpointStore(path) + checkpoint = OpenLaneCheckpoint( + max_pages=512, + completed_pages=[1, 2, 3], + failed_pages=[7], + total_records=99, + last_saved_at="2026-04-20T00:00:00+00:00", + ) + + store.save(checkpoint) + loaded = store.load(max_pages=512) + + self.assertEqual(loaded.max_pages, 512) + self.assertEqual(loaded.completed_pages, [1, 2, 3]) + self.assertEqual(loaded.failed_pages, [7]) + self.assertEqual(loaded.total_records, 99) + self.assertEqual(loaded.last_saved_at, "2026-04-20T00:00:00+00:00") + + +class TestOpenLaneWriter(unittest.TestCase): + def test_append_and_finalize_outputs(self) -> None: + with tempfile.TemporaryDirectory() as tmp_dir: + jsonl_path = Path(tmp_dir) / "openlane.jsonl" + aggregated_path = Path(tmp_dir) / "openlane_aggregated.json" + writer = OpenLaneResultWriter(jsonl_path, aggregated_path) + + written = writer.append_page(1, [{"id": 1}, {"id": 2}]) + self.assertEqual(written, 2) + + summary = writer.finalize(max_pages=10, completed_pages=[1], failed_pages=[2]) + + self.assertEqual(summary["total_records"], 2) + self.assertEqual(summary["total_pages_completed"], 1) + self.assertEqual(summary["total_pages_failed"], 1) + self.assertTrue(aggregated_path.exists()) + + payload = json.loads(aggregated_path.read_text(encoding="utf-8")) + self.assertEqual(payload["items"][0]["page"], 1) + self.assertEqual(payload["items"][0]["record"]["id"], 1) + + +class TestOpenLaneConfig(unittest.TestCase): + def test_retry_schedule_defaults(self) -> None: + config = OpenLaneConfig() + self.assertGreaterEqual(len(config.retry_schedule_seconds), 1) + self.assertEqual(tuple(float(x) for x in config.retry_schedule_seconds), config.retry_schedule_seconds) + + def test_storage_state_defaults(self) -> None: + config = OpenLaneConfig() + self.assertTrue(config.storage_state_file) + self.assertIn("artifacts/openlane", config.storage_state_file) + self.assertIsInstance(config.persist_storage_state, bool) + + +class TestOpenLaneClient(unittest.TestCase): + def test_extract_records_from_common_shapes(self) -> None: + self.assertEqual( + OpenLaneClient._extract_records({"results": [{"id": 1}, {"id": 2}]}), + [{"id": 1}, {"id": 2}], + ) + self.assertEqual( + OpenLaneClient._extract_records({"data": {"items": [{"id": 3}]}}), + [{"id": 3}], + ) + self.assertEqual( + OpenLaneClient._extract_records({"vehicles": [{"id": 4}]}), + [{"id": 4}], + ) + self.assertEqual(OpenLaneClient._extract_records({"data": {"foo": "bar"}}), []) + + def test_extract_vehicle_images_from_common_shapes(self) -> None: + payload = { + "vehicle_images": [ + {"url": "https://img/1.png", "low_resolution_url": "https://img/1_small.png"}, + ] + } + self.assertEqual( + OpenLaneClient._extract_vehicle_images(payload), + [{"url": "https://img/1.png", "low_resolution_url": "https://img/1_small.png"}], + ) + + payload_nested = { + "data": { + "images": [ + {"url": "https://img/2.png"}, + ] + } + } + self.assertEqual( + OpenLaneClient._extract_vehicle_images(payload_nested), + [{"url": "https://img/2.png"}], + ) + + self.assertEqual(OpenLaneClient._extract_vehicle_images({"data": {"foo": 1}}), []) + + def test_extract_vehicle_images_batch_helper_normalization(self) -> None: + # Косвенно проверяем нормализацию входа под батчевый метод + ids = [" 123 ", "123", 456, "", None] + normalized = [] + seen = set() + for v in ids: + s = str(v).strip() + if not s or s in seen: + continue + seen.add(s) + normalized.append(s) + + self.assertEqual(normalized, ["123", "456", "None"]) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_scraper_images.py b/tests/test_scraper_images.py new file mode 100644 index 0000000..7ed8cb7 --- /dev/null +++ b/tests/test_scraper_images.py @@ -0,0 +1,68 @@ +from __future__ import annotations + +import unittest + +from openlane_scraper.scraper import OpenLaneScraper +from openlane_scraper.storage.schemas import ImageRecord + + +class TestScraperImageHelpers(unittest.TestCase): + def test_build_image_records_supports_image_large_and_image(self) -> None: + payload = [ + { + "image": "https://img.cdn/vehicle-medium-1.jpg", + "image_large": "https://img.cdn/vehicle-large-1.jpg", + }, + { + "image": "https://img.cdn/vehicle-medium-2.jpg", + "large_resolution_url": "https://img.cdn/vehicle-large-2.jpg", + "low_resolution_url": "https://img.cdn/vehicle-low-2.jpg", + }, + ] + + records = OpenLaneScraper._build_image_records(payload) + + self.assertEqual(len(records), 2) + self.assertEqual(records[0].fullres_image, "https://img.cdn/vehicle-large-1.jpg") + self.assertEqual(records[0].preview_image, "https://img.cdn/vehicle-medium-1.jpg") + self.assertEqual(records[1].fullres_image, "https://img.cdn/vehicle-large-2.jpg") + self.assertEqual(records[1].preview_image, "https://img.cdn/vehicle-low-2.jpg") + + def test_merge_image_records_keeps_unique_and_reindexes(self) -> None: + existing = [ + ImageRecord( + fullres_image="https://img.cdn/1.jpg", + preview_image="https://img.cdn/1_small.jpg", + order_index=10, + ), + ImageRecord( + fullres_image="https://img.cdn/2.jpg", + preview_image="https://img.cdn/2_small.jpg", + order_index=11, + ), + ] + fetched = [ + ImageRecord( + fullres_image="https://img.cdn/2.jpg", + preview_image="https://img.cdn/2_small_new.jpg", + order_index=0, + ), + ImageRecord( + fullres_image="https://img.cdn/3.jpg", + preview_image="https://img.cdn/3_small.jpg", + order_index=1, + ), + ] + + merged = OpenLaneScraper._merge_image_records(existing, fetched) + + self.assertEqual([img.fullres_image for img in merged], [ + "https://img.cdn/2.jpg", + "https://img.cdn/3.jpg", + "https://img.cdn/1.jpg", + ]) + self.assertEqual([img.order_index for img in merged], [0, 1, 2]) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_utils.py b/tests/test_utils.py new file mode 100644 index 0000000..9e4559b --- /dev/null +++ b/tests/test_utils.py @@ -0,0 +1,24 @@ +from __future__ import annotations + +import unittest + +from openlane_scraper.core.utils import deep_find_key + + +class TestDeepFindKey(unittest.TestCase): + def test_finds_nested_and_respects_depth(self) -> None: + payload = { + "root": { + "target": "a", + "nested": [{"target": "b"}, {"x": 1}], + } + } + self.assertEqual(deep_find_key(payload, {"target"}), ["a", "b"]) + + deep_payload = {"l1": {"l2": {"l3": {"target": "value"}}}} + self.assertEqual(deep_find_key(deep_payload, {"target"}, max_depth=2), []) + self.assertEqual(deep_find_key(deep_payload, {"target"}, max_depth=8), ["value"]) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_worker_tasks.py b/tests/test_worker_tasks.py new file mode 100644 index 0000000..3c00a5f --- /dev/null +++ b/tests/test_worker_tasks.py @@ -0,0 +1,189 @@ +from __future__ import annotations + +import unittest +from unittest.mock import MagicMock, patch + +from openlane_scraper.worker import tasks + + +class TestWorkerTaskLockHelpers(unittest.TestCase): + def test_lock_acquire_refresh_release(self) -> None: + redis_client = MagicMock() + + # Acquire. + redis_client.set.return_value = True + self.assertTrue(tasks._acquire_lock(redis_client, "lock:key", "owner-token", 120)) + redis_client.set.assert_called_once_with("lock:key", "owner-token", nx=True, ex=120) + + # Refresh. + redis_client.eval.return_value = 1 + self.assertTrue(tasks._refresh_lock_if_owner(redis_client, "lock:key", "owner-token", 120)) + + # Release. + redis_client.eval.reset_mock() + tasks._release_lock_if_owner(redis_client, "lock:key", "owner-token") + args = redis_client.eval.call_args[0] + self.assertEqual(args[2], "lock:key") + self.assertEqual(args[3], "owner-token") + + def test_sync_listing_task_skips_when_lock_not_acquired(self) -> None: + with patch.object(tasks, "_get_persistence") as get_persistence, \ + patch.object(tasks, "_get_redis") as get_redis, \ + patch.object(tasks, "_acquire_lock", return_value=False): + persistence = MagicMock() + get_persistence.return_value = persistence + get_redis.return_value = MagicMock() + + tasks.sync_listing_task.push_request(id="task-123") + try: + result = tasks.sync_listing_task.run() + finally: + tasks.sync_listing_task.pop_request() + + persistence.create_tables.assert_called_once() + self.assertEqual(result["status"], "skipped") + self.assertEqual(result["reason"], "sync_already_running") + + def test_sync_listing_task_releases_owned_lock(self) -> None: + with patch.object(tasks, "_get_persistence") as get_persistence, \ + patch.object(tasks, "_get_redis") as get_redis, \ + patch.object(tasks, "_acquire_lock", return_value=True), \ + patch.object(tasks, "_is_full_scan_done", return_value=True), \ + patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \ + patch.object(tasks, "_release_lock_if_owner") as release_lock, \ + patch.object(tasks.sync_listing_task, "update_state"), \ + patch("openlane_scraper.worker.tasks.time.sleep"), \ + patch.object(tasks, "_run_browser_job", return_value={ + "run_id": 7, + "status": "success", + "cars_upserted": 2, + "cars_failed": 0, + "images_upserted": 4, + "skipped_existing": 1, + "elapsed_seconds": 1.25, + }): + persistence = MagicMock() + get_persistence.return_value = persistence + redis_client = MagicMock() + redis_client.get.return_value = None + get_redis.return_value = redis_client + stop_event = MagicMock() + heartbeat_thread = MagicMock() + start_heartbeat.return_value = (stop_event, heartbeat_thread) + + tasks.sync_listing_task.push_request(id="task-123") + try: + result = tasks.sync_listing_task.run() + finally: + tasks.sync_listing_task.pop_request() + + self.assertEqual(result["status"], "success") + stop_event.set.assert_called_once() + heartbeat_thread.join.assert_called_once() + release_lock.assert_called_once() + + def test_clear_orphan_sync_listing_lock(self) -> None: + # Нет запущенных задач → удаляет. + redis_client = MagicMock() + redis_client.get.return_value = "owner-token" + redis_client.ttl.return_value = 120 + celery_app = MagicMock() + inspector = MagicMock() + inspector.active.return_value = {"worker@node": []} + inspector.reserved.return_value = {"worker@node": []} + inspector.scheduled.return_value = {"worker@node": []} + celery_app.control.inspect.return_value = inspector + + self.assertTrue(tasks._clear_orphan_sync_listing_lock(redis_client, celery_app)) + redis_client.delete.assert_called_once_with(tasks.SYNC_LISTING_LOCK_KEY) + + # Задача активна → не удаляет. + redis_client2 = MagicMock() + redis_client2.get.return_value = "owner-token" + inspector2 = MagicMock() + inspector2.active.return_value = {"worker@node": [{"name": tasks.SYNC_LISTING_TASK_NAME}]} + inspector2.reserved.return_value = {"worker@node": []} + inspector2.scheduled.return_value = {"worker@node": []} + celery_app2 = MagicMock() + celery_app2.control.inspect.return_value = inspector2 + + self.assertFalse(tasks._clear_orphan_sync_listing_lock(redis_client2, celery_app2)) + redis_client2.delete.assert_not_called() + + def test_sync_listing_task_recovers_orphan_lock_and_runs(self) -> None: + with patch.object(tasks, "_get_persistence") as get_persistence, \ + patch.object(tasks, "_get_redis") as get_redis, \ + patch.object(tasks, "_acquire_lock", side_effect=[False, True]) as acquire_lock, \ + patch.object(tasks, "_clear_orphan_sync_listing_lock", return_value=True) as clear_orphan, \ + patch.object(tasks, "_is_full_scan_done", return_value=True), \ + patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \ + patch.object(tasks, "_release_lock_if_owner") as release_lock, \ + patch.object(tasks.sync_listing_task, "update_state"), \ + patch("openlane_scraper.worker.tasks.time.sleep"), \ + patch.object(tasks, "_run_browser_job", return_value={ + "run_id": 9, + "status": "success", + "cars_upserted": 3, + "cars_failed": 0, + "images_upserted": 5, + "skipped_existing": 0, + "elapsed_seconds": 2.0, + }): + persistence = MagicMock() + get_persistence.return_value = persistence + redis_client = MagicMock() + redis_client.get.return_value = None + get_redis.return_value = redis_client + stop_event = MagicMock() + heartbeat_thread = MagicMock() + start_heartbeat.return_value = (stop_event, heartbeat_thread) + + tasks.sync_listing_task.push_request(id="task-456") + try: + result = tasks.sync_listing_task.run() + finally: + tasks.sync_listing_task.pop_request() + + self.assertEqual(result["status"], "success") + clear_orphan.assert_called_once() + self.assertEqual(acquire_lock.call_count, 2) + release_lock.assert_called_once() + + def test_bootstrap_forces_full_scan(self) -> None: + with patch.object(tasks, "_get_persistence") as get_persistence, \ + patch.object(tasks, "_get_redis") as get_redis, \ + patch.object(tasks, "_acquire_lock", return_value=True), \ + patch.object(tasks, "_is_full_scan_done", return_value=False), \ + patch.object(tasks, "_set_full_scan_done") as set_done, \ + patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \ + patch.object(tasks, "_release_lock_if_owner"), \ + patch.object(tasks.sync_listing_task, "update_state"), \ + patch("openlane_scraper.worker.tasks.time.sleep"), \ + patch.object(tasks, "_run_browser_job", return_value={ + "run_id": 11, + "status": "success", + "cars_upserted": 5, + "cars_failed": 0, + "images_upserted": 10, + "skipped_existing": 0, + "elapsed_seconds": 3.0, + }): + persistence = MagicMock() + get_persistence.return_value = persistence + redis_client = MagicMock() + get_redis.return_value = redis_client + start_heartbeat.return_value = (MagicMock(), MagicMock()) + + tasks.sync_listing_task.push_request(id="task-bootstrap") + try: + result = tasks.sync_listing_task.run(only_new=True) + finally: + tasks.sync_listing_task.pop_request() + + self.assertEqual(result["status"], "success") + # Bootstrap should set full_scan_done to True on success. + set_done.assert_called_once_with(redis_client, True) + + +if __name__ == "__main__": + unittest.main()