From da89304d38c1514e78bc85214d2bc0b6d37de108 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Tue, 18 Aug 2026 13:47:00 +0300 Subject: [PATCH] Keep enrichment locks alive during long batches --- docker-compose.yml | 2 +- mobilede_scraper/worker/tasks.py | 19 ++++++++++++++++- tests/test_mobilede_enrichment.py | 34 ++++++++++++++++++++++++++++++- 3 files changed, 52 insertions(+), 3 deletions(-) diff --git a/docker-compose.yml b/docker-compose.yml index f59a3ad..25f3dc1 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -57,7 +57,7 @@ MOBILEDE_DETAIL_TIMEOUT_SECONDS: ${MOBILEDE_DETAIL_TIMEOUT_SECONDS:-10} MOBILEDE_BEAT_ENRICH_ENABLED: ${MOBILEDE_BEAT_ENRICH_ENABLED:-true} MOBILEDE_ENRICH_INTERVAL_SECONDS: ${MOBILEDE_ENRICH_INTERVAL_SECONDS:-900} - MOBILEDE_ENRICH_LOCK_TTL_SECONDS: ${MOBILEDE_ENRICH_LOCK_TTL_SECONDS:-3600} + MOBILEDE_ENRICH_LOCK_TTL_SECONDS: ${MOBILEDE_ENRICH_LOCK_TTL_SECONDS:-300} MOBILEDE_STARTUP_SYNC_ENABLED: ${MOBILEDE_STARTUP_SYNC_ENABLED:-true} MOBILEDE_RUNTIME_CONFIG_FILE: ${MOBILEDE_RUNTIME_CONFIG_FILE:-/app/runtime_config.json} # Автовосстановление worker. diff --git a/mobilede_scraper/worker/tasks.py b/mobilede_scraper/worker/tasks.py index 0ff4b08..a47aed2 100644 --- a/mobilede_scraper/worker/tasks.py +++ b/mobilede_scraper/worker/tasks.py @@ -2505,7 +2505,7 @@ def mobilede_enrich_images_batch_task( redis_client = _get_redis() owner_token = str(getattr(self.request, "id", None) or uuid.uuid4().hex) - lock_ttl = max(300, int(float(os.getenv("MOBILEDE_ENRICH_LOCK_TTL_SECONDS", "3600")))) + lock_ttl = max(300, int(float(os.getenv("MOBILEDE_ENRICH_LOCK_TTL_SECONDS", "300")))) if not _acquire_lock(redis_client, MOBILEDE_ENRICH_IMAGES_LOCK_KEY, owner_token, lock_ttl): return { "status": "locked", @@ -2529,6 +2529,23 @@ def mobilede_enrich_images_batch_task( for batch_start in range(0, len(candidates), effective_batch_size): batch = candidates[batch_start:batch_start + effective_batch_size] for _car_id, origin_id, origin_url, image_count in batch: + if not _refresh_lock_if_owner( + redis_client, + MOBILEDE_ENRICH_IMAGES_LOCK_KEY, + owner_token, + lock_ttl, + ): + logger.warning( + "mobile.de image enrich lock ownership lost; stopping before origin_id=%s", + origin_id, + ) + return { + "status": "lock_lost", + "candidates": len(candidates), + "enriched": enriched, + "failed": failed, + "skipped": skipped, + } listing_id = _mobilede_extract_listing_id(origin_url) or str(origin_id).rsplit(":", 1)[-1] if not listing_id: skipped += 1 diff --git a/tests/test_mobilede_enrichment.py b/tests/test_mobilede_enrichment.py index 9e4749d..dab6bee 100644 --- a/tests/test_mobilede_enrichment.py +++ b/tests/test_mobilede_enrichment.py @@ -81,7 +81,39 @@ class TestMobileDeEnrichment(unittest.TestCase): self.assertGreaterEqual(selector_kwargs["stale_before"], before) self.assertLessEqual(selector_kwargs["stale_before"], after) self.assertEqual(scraper.sync_detail.call_count, 2) - redis_client.eval.assert_called() + # Two heartbeat refreshes plus owner-safe release. + self.assertGreaterEqual(redis_client.eval.call_count, 3) + + def test_task_stops_when_enrichment_lock_ownership_is_lost(self) -> None: + enrichment = SimpleNamespace( + enabled=True, + batch_size=1, + max_existing_images=1, + staleness_hours=48, + delay_seconds=0.0, + ) + runtime_config = SimpleNamespace(mobilede=SimpleNamespace(enrichment=enrichment)) + redis_client = MagicMock() + redis_client.set.return_value = True + redis_client.eval.side_effect = [0, 0] + persistence = MagicMock() + persistence.get_active_cars_batch_for_image_enrich.return_value = [ + (1, "mobile.de:111", "https://suchen.mobile.de/fahrzeuge/details.html?id=111", 0), + ] + scraper = MagicMock() + + with ( + 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, "_get_redis", return_value=redis_client), + patch.object(tasks, "_get_persistence", return_value=persistence), + patch.object(tasks, "MobileDeScraper", return_value=scraper), + ): + result = tasks.mobilede_enrich_images_batch_task.run() + + self.assertEqual(result["status"], "lock_lost") + self.assertEqual(result["enriched"], 0) + scraper.sync_detail.assert_not_called() def test_task_skips_when_enrichment_lock_is_held(self) -> None: enrichment = SimpleNamespace(