Keep enrichment locks alive during long batches
This commit is contained in:
+1
-1
@@ -57,7 +57,7 @@
|
|||||||
MOBILEDE_DETAIL_TIMEOUT_SECONDS: ${MOBILEDE_DETAIL_TIMEOUT_SECONDS:-10}
|
MOBILEDE_DETAIL_TIMEOUT_SECONDS: ${MOBILEDE_DETAIL_TIMEOUT_SECONDS:-10}
|
||||||
MOBILEDE_BEAT_ENRICH_ENABLED: ${MOBILEDE_BEAT_ENRICH_ENABLED:-true}
|
MOBILEDE_BEAT_ENRICH_ENABLED: ${MOBILEDE_BEAT_ENRICH_ENABLED:-true}
|
||||||
MOBILEDE_ENRICH_INTERVAL_SECONDS: ${MOBILEDE_ENRICH_INTERVAL_SECONDS:-900}
|
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_STARTUP_SYNC_ENABLED: ${MOBILEDE_STARTUP_SYNC_ENABLED:-true}
|
||||||
MOBILEDE_RUNTIME_CONFIG_FILE: ${MOBILEDE_RUNTIME_CONFIG_FILE:-/app/runtime_config.json}
|
MOBILEDE_RUNTIME_CONFIG_FILE: ${MOBILEDE_RUNTIME_CONFIG_FILE:-/app/runtime_config.json}
|
||||||
# Автовосстановление worker.
|
# Автовосстановление worker.
|
||||||
|
|||||||
@@ -2505,7 +2505,7 @@ def mobilede_enrich_images_batch_task(
|
|||||||
|
|
||||||
redis_client = _get_redis()
|
redis_client = _get_redis()
|
||||||
owner_token = str(getattr(self.request, "id", None) or uuid.uuid4().hex)
|
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):
|
if not _acquire_lock(redis_client, MOBILEDE_ENRICH_IMAGES_LOCK_KEY, owner_token, lock_ttl):
|
||||||
return {
|
return {
|
||||||
"status": "locked",
|
"status": "locked",
|
||||||
@@ -2529,6 +2529,23 @@ def mobilede_enrich_images_batch_task(
|
|||||||
for batch_start in range(0, len(candidates), effective_batch_size):
|
for batch_start in range(0, len(candidates), effective_batch_size):
|
||||||
batch = candidates[batch_start:batch_start + effective_batch_size]
|
batch = candidates[batch_start:batch_start + effective_batch_size]
|
||||||
for _car_id, origin_id, origin_url, image_count in batch:
|
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]
|
listing_id = _mobilede_extract_listing_id(origin_url) or str(origin_id).rsplit(":", 1)[-1]
|
||||||
if not listing_id:
|
if not listing_id:
|
||||||
skipped += 1
|
skipped += 1
|
||||||
|
|||||||
@@ -81,7 +81,39 @@ class TestMobileDeEnrichment(unittest.TestCase):
|
|||||||
self.assertGreaterEqual(selector_kwargs["stale_before"], before)
|
self.assertGreaterEqual(selector_kwargs["stale_before"], before)
|
||||||
self.assertLessEqual(selector_kwargs["stale_before"], after)
|
self.assertLessEqual(selector_kwargs["stale_before"], after)
|
||||||
self.assertEqual(scraper.sync_detail.call_count, 2)
|
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:
|
def test_task_skips_when_enrichment_lock_is_held(self) -> None:
|
||||||
enrichment = SimpleNamespace(
|
enrichment = SimpleNamespace(
|
||||||
|
|||||||
Reference in New Issue
Block a user