661 lines
31 KiB
Python
661 lines
31 KiB
Python
from __future__ import annotations
|
|
|
|
import unittest
|
|
from unittest.mock import MagicMock, patch
|
|
|
|
from iaai_scraper.worker import tasks
|
|
|
|
|
|
def make_settings_stub(*, parallel_segments: bool = False) -> MagicMock:
|
|
settings_stub = MagicMock()
|
|
settings_stub.celery.parallel_segments = parallel_segments
|
|
settings_stub.celery.task_soft_time_limit = 3300
|
|
settings_stub.celery.task_time_limit = 3600
|
|
settings_stub.celery.task_stall_timeout_seconds = 600
|
|
settings_stub.listing.listing_segments_json = "ignored"
|
|
settings_stub.discovery.mode = "listing"
|
|
settings_stub.discovery.hourly_mode = "rolling_refresh"
|
|
return settings_stub
|
|
|
|
|
|
class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|
def test_sync_listing_task_uses_hourly_sitemap_even_when_mode_is_not_sitemap(self) -> None:
|
|
settings_stub = make_settings_stub(parallel_segments=False)
|
|
settings_stub.discovery.mode = "listing"
|
|
settings_stub.discovery.hourly_mode = "rolling_refresh"
|
|
|
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
|
patch.object(tasks, "_get_redis") as get_redis, \
|
|
patch.object(tasks, "Settings", return_value=settings_stub), \
|
|
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, "_hourly_sitemap_rolling_refresh_sync", return_value={
|
|
"status": "success",
|
|
"run_id": None,
|
|
"cars_upserted": 25,
|
|
"cars_failed": 0,
|
|
"images_upserted": 50,
|
|
"skipped_existing": 100,
|
|
"elapsed_seconds": None,
|
|
"failures": [],
|
|
"hourly_mode": "sitemap_rolling_refresh",
|
|
"discovered_urls": 105,
|
|
"new_urls": 5,
|
|
"refresh_urls": 20,
|
|
"sold_marked": 2,
|
|
}) as hourly_sync, \
|
|
patch.object(tasks.sync_listing_task, "update_state"):
|
|
get_persistence.return_value = MagicMock()
|
|
redis_client = MagicMock()
|
|
redis_client.get.return_value = None
|
|
get_redis.return_value = redis_client
|
|
start_heartbeat.return_value = (MagicMock(), MagicMock())
|
|
|
|
tasks.sync_listing_task.push_request(id="task-hourly-force-sitemap")
|
|
try:
|
|
result = tasks.sync_listing_task.run()
|
|
finally:
|
|
tasks.sync_listing_task.pop_request()
|
|
|
|
self.assertEqual(result["status"], "success")
|
|
self.assertEqual(result["hourly_mode"], "sitemap_rolling_refresh")
|
|
hourly_sync.assert_called_once()
|
|
release_lock.assert_called_once()
|
|
|
|
def test_sync_listing_task_uses_hourly_sitemap_rolling_refresh_after_bootstrap(self) -> None:
|
|
settings_stub = make_settings_stub(parallel_segments=False)
|
|
settings_stub.discovery.mode = "sitemap"
|
|
settings_stub.discovery.hourly_mode = "rolling_refresh"
|
|
|
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
|
patch.object(tasks, "_get_redis") as get_redis, \
|
|
patch.object(tasks, "Settings", return_value=settings_stub), \
|
|
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, "_hourly_sitemap_rolling_refresh_sync", return_value={
|
|
"status": "success",
|
|
"run_id": None,
|
|
"cars_upserted": 25,
|
|
"cars_failed": 0,
|
|
"images_upserted": 50,
|
|
"skipped_existing": 100,
|
|
"elapsed_seconds": None,
|
|
"failures": [],
|
|
"hourly_mode": "sitemap_rolling_refresh",
|
|
"discovered_urls": 105,
|
|
"new_urls": 5,
|
|
"refresh_urls": 20,
|
|
"sold_marked": 2,
|
|
}) as hourly_sync, \
|
|
patch.object(tasks.sync_listing_task, "update_state"):
|
|
get_persistence.return_value = MagicMock()
|
|
redis_client = MagicMock()
|
|
redis_client.get.return_value = None
|
|
get_redis.return_value = redis_client
|
|
start_heartbeat.return_value = (MagicMock(), MagicMock())
|
|
|
|
tasks.sync_listing_task.push_request(id="task-hourly")
|
|
try:
|
|
result = tasks.sync_listing_task.run()
|
|
finally:
|
|
tasks.sync_listing_task.pop_request()
|
|
|
|
self.assertEqual(result["status"], "success")
|
|
self.assertEqual(result["hourly_mode"], "sitemap_rolling_refresh")
|
|
hourly_sync.assert_called_once()
|
|
release_lock.assert_called_once()
|
|
|
|
def test_sync_listing_task_can_use_hourly_sitemap_diff_when_configured(self) -> None:
|
|
settings_stub = make_settings_stub(parallel_segments=False)
|
|
settings_stub.discovery.mode = "sitemap"
|
|
settings_stub.discovery.hourly_mode = "diff"
|
|
|
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
|
patch.object(tasks, "_get_redis") as get_redis, \
|
|
patch.object(tasks, "Settings", return_value=settings_stub), \
|
|
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, "_hourly_sitemap_diff_sync", return_value={
|
|
"status": "success",
|
|
"run_id": None,
|
|
"cars_upserted": 5,
|
|
"cars_failed": 0,
|
|
"images_upserted": 10,
|
|
"skipped_existing": 100,
|
|
"elapsed_seconds": None,
|
|
"failures": [],
|
|
"hourly_mode": "sitemap_diff",
|
|
"discovered_urls": 105,
|
|
"new_urls": 5,
|
|
"sold_marked": 2,
|
|
}) as hourly_sync, \
|
|
patch.object(tasks.sync_listing_task, "update_state"):
|
|
get_persistence.return_value = MagicMock()
|
|
redis_client = MagicMock()
|
|
redis_client.get.return_value = None
|
|
get_redis.return_value = redis_client
|
|
start_heartbeat.return_value = (MagicMock(), MagicMock())
|
|
|
|
tasks.sync_listing_task.push_request(id="task-hourly-diff")
|
|
try:
|
|
result = tasks.sync_listing_task.run()
|
|
finally:
|
|
tasks.sync_listing_task.pop_request()
|
|
|
|
self.assertEqual(result["status"], "success")
|
|
self.assertEqual(result["hourly_mode"], "sitemap_diff")
|
|
hourly_sync.assert_called_once()
|
|
release_lock.assert_called_once()
|
|
|
|
def test_sync_listing_task_forces_sitemap_full_scan_in_bootstrap(self) -> None:
|
|
settings_stub = make_settings_stub(parallel_segments=False)
|
|
settings_stub.discovery.mode = "listing"
|
|
|
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
|
patch.object(tasks, "_get_redis") as get_redis, \
|
|
patch.object(tasks, "Settings", return_value=settings_stub), \
|
|
patch.object(tasks, "_acquire_lock", return_value=True), \
|
|
patch.object(tasks, "_is_full_scan_done", return_value=False), \
|
|
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
|
|
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
|
|
patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \
|
|
patch.object(tasks.sync_listing_task, "update_state"), \
|
|
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=[{"make": "HONDA"}]):
|
|
get_persistence.return_value = MagicMock()
|
|
redis_client = MagicMock()
|
|
redis_client.get.return_value = None
|
|
get_redis.return_value = redis_client
|
|
start_heartbeat.return_value = (MagicMock(), MagicMock())
|
|
|
|
sync_listing_mock = MagicMock(return_value={
|
|
"run_id": 99,
|
|
"status": "success",
|
|
"full_scan_completed": True,
|
|
"cars_upserted": 10,
|
|
"cars_failed": 0,
|
|
"images_upserted": 20,
|
|
"skipped_existing": 0,
|
|
"elapsed_seconds": 1.0,
|
|
"failures": [],
|
|
})
|
|
sync_listing_segmented_mock = MagicMock(side_effect=AssertionError("segmented path should not be used"))
|
|
scraper = MagicMock()
|
|
scraper.sync_listing = sync_listing_mock
|
|
scraper.sync_listing_segmented = sync_listing_segmented_mock
|
|
scraper_ctx = MagicMock()
|
|
scraper_ctx.__enter__.return_value = scraper
|
|
scraper_ctx.__exit__.return_value = None
|
|
|
|
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
|
|
tasks.sync_listing_task.push_request(id="task-bootstrap-sitemap")
|
|
try:
|
|
result = tasks.sync_listing_task.run()
|
|
finally:
|
|
tasks.sync_listing_task.pop_request()
|
|
|
|
self.assertEqual(result["status"], "success")
|
|
sync_listing_mock.assert_called_once_with(
|
|
make=None,
|
|
model=None,
|
|
lane="iaai_cars",
|
|
limit=None,
|
|
only_new=False,
|
|
)
|
|
sync_listing_segmented_mock.assert_not_called()
|
|
release_lock.assert_called_once()
|
|
|
|
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.object(tasks, "_run_browser_job", return_value={
|
|
"run_id": 7,
|
|
"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(make="Toyota")
|
|
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.object(tasks, "_run_browser_job", return_value={
|
|
"run_id": 9,
|
|
"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(make="Honda")
|
|
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_checkpoint_save_load_roundtrip(self) -> None:
|
|
storage: dict[str, str] = {}
|
|
|
|
def _fake_set(key, value, ex=None):
|
|
storage[key] = value
|
|
return True
|
|
|
|
redis_client = MagicMock()
|
|
redis_client.set.side_effect = _fake_set
|
|
redis_client.get.side_effect = lambda key: storage.get(key)
|
|
|
|
tasks._save_last_completed_segment(redis_client, 12)
|
|
|
|
self.assertEqual(tasks._load_last_completed_segment(redis_client), 12)
|
|
self.assertEqual(redis_client.set.call_args.kwargs["ex"], tasks.SYNC_LISTING_CHECKPOINT_TTL_SECONDS)
|
|
|
|
def test_load_checkpoint_clears_invalid_payload(self) -> None:
|
|
redis_client = MagicMock()
|
|
redis_client.get.return_value = "{not-a-number"
|
|
|
|
self.assertIsNone(tasks._load_last_completed_segment(redis_client))
|
|
redis_client.delete.assert_called_once_with(tasks.SYNC_LISTING_CHECKPOINT_KEY)
|
|
|
|
def test_load_checkpoint_empty_when_missing(self) -> None:
|
|
redis_client = MagicMock()
|
|
redis_client.get.return_value = None
|
|
|
|
self.assertIsNone(tasks._load_last_completed_segment(redis_client))
|
|
redis_client.delete.assert_not_called()
|
|
|
|
def test_sync_listing_bootstrap_prefers_sitemap_over_segment_resume(self) -> None:
|
|
segments = [
|
|
{"make": "ACURA"},
|
|
{"make": "AUDI"},
|
|
{"make": "BMW"},
|
|
{"make": "EAGLE"},
|
|
]
|
|
settings_stub = make_settings_stub()
|
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
|
patch.object(tasks, "_get_redis") as get_redis, \
|
|
patch.object(tasks, "Settings", return_value=settings_stub), \
|
|
patch.object(tasks, "_acquire_lock", return_value=True), \
|
|
patch.object(tasks, "_is_full_scan_done", return_value=False), \
|
|
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
|
|
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
|
|
patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \
|
|
patch.object(tasks.sync_listing_task, "update_state"), \
|
|
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=segments):
|
|
get_persistence.return_value = MagicMock()
|
|
redis_client = MagicMock()
|
|
redis_client.get.side_effect = lambda key: (
|
|
"1" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
|
|
)
|
|
get_redis.return_value = redis_client
|
|
start_heartbeat.return_value = (MagicMock(), MagicMock())
|
|
|
|
sync_segmented_mock = MagicMock(side_effect=AssertionError("segmented path should not be used"))
|
|
sync_listing_mock = MagicMock(return_value={
|
|
"run_id": 11, "status": "success", "full_scan_completed": True,
|
|
"cars_upserted": 1, "cars_failed": 0, "images_upserted": 0,
|
|
"skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [],
|
|
})
|
|
scraper_ctx = MagicMock()
|
|
scraper_ctx.__enter__.return_value.sync_listing_segmented = sync_segmented_mock
|
|
scraper_ctx.__enter__.return_value.sync_listing = sync_listing_mock
|
|
scraper_ctx.__exit__.return_value = None
|
|
|
|
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
|
|
tasks.sync_listing_task.push_request(id="task-resume-seg")
|
|
try:
|
|
result = tasks.sync_listing_task.run()
|
|
finally:
|
|
tasks.sync_listing_task.pop_request()
|
|
|
|
self.assertEqual(result["status"], "success")
|
|
sync_listing_mock.assert_called_once_with(
|
|
make=None,
|
|
model=None,
|
|
lane="iaai_cars",
|
|
limit=None,
|
|
only_new=False,
|
|
)
|
|
sync_segmented_mock.assert_not_called()
|
|
release_lock.assert_called_once()
|
|
|
|
def test_sync_listing_hourly_path_ignores_checkpoint_after_full_scan_completed(self) -> None:
|
|
segments = [{"make": "ACURA"}, {"make": "AUDI"}]
|
|
settings_stub = make_settings_stub()
|
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
|
patch.object(tasks, "_get_redis") as get_redis, \
|
|
patch.object(tasks, "Settings", return_value=settings_stub), \
|
|
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, "_clear_sync_checkpoint") as clear_checkpoint, \
|
|
patch.object(tasks, "_hourly_sitemap_rolling_refresh_sync", return_value={
|
|
"status": "success",
|
|
"run_id": None,
|
|
"cars_upserted": 1,
|
|
"cars_failed": 0,
|
|
"images_upserted": 0,
|
|
"skipped_existing": 0,
|
|
"elapsed_seconds": None,
|
|
"failures": [],
|
|
"hourly_mode": "sitemap_rolling_refresh",
|
|
"discovered_urls": 10,
|
|
"new_urls": 1,
|
|
"refresh_urls": 1,
|
|
"sold_marked": 0,
|
|
}) as hourly_sync, \
|
|
patch.object(tasks.sync_listing_task, "update_state"), \
|
|
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=segments):
|
|
get_persistence.return_value = MagicMock()
|
|
redis_client = MagicMock()
|
|
# Оставшийся чекпоинт не должен использоваться.
|
|
redis_client.get.side_effect = lambda key: (
|
|
"0" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
|
|
)
|
|
get_redis.return_value = redis_client
|
|
start_heartbeat.return_value = (MagicMock(), MagicMock())
|
|
|
|
tasks.sync_listing_task.push_request(id="task-792")
|
|
try:
|
|
result = tasks.sync_listing_task.run()
|
|
finally:
|
|
tasks.sync_listing_task.pop_request()
|
|
|
|
self.assertEqual(result["status"], "success")
|
|
self.assertEqual(result["hourly_mode"], "sitemap_rolling_refresh")
|
|
hourly_sync.assert_called_once()
|
|
clear_checkpoint.assert_called()
|
|
release_lock.assert_called_once()
|
|
|
|
def test_sync_listing_bootstrap_ignores_beyond_segment_checkpoint(self) -> None:
|
|
segments = [{"make": "ACURA"}, {"make": "AUDI"}]
|
|
settings_stub = make_settings_stub()
|
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
|
patch.object(tasks, "_get_redis") as get_redis, \
|
|
patch.object(tasks, "Settings", return_value=settings_stub), \
|
|
patch.object(tasks, "_acquire_lock", return_value=True), \
|
|
patch.object(tasks, "_is_full_scan_done", return_value=False), \
|
|
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
|
|
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
|
|
patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \
|
|
patch.object(tasks.sync_listing_task, "update_state"), \
|
|
patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=segments):
|
|
get_persistence.return_value = MagicMock()
|
|
redis_client = MagicMock()
|
|
redis_client.get.side_effect = lambda key: (
|
|
"99" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
|
|
)
|
|
get_redis.return_value = redis_client
|
|
start_heartbeat.return_value = (MagicMock(), MagicMock())
|
|
|
|
sync_segmented_mock = MagicMock(side_effect=AssertionError("segmented path should not be used"))
|
|
sync_listing_mock = MagicMock(return_value={
|
|
"run_id": 15, "status": "success", "full_scan_completed": True,
|
|
"cars_upserted": 0, "cars_failed": 0, "images_upserted": 0,
|
|
"skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [],
|
|
})
|
|
scraper_ctx = MagicMock()
|
|
scraper_ctx.__enter__.return_value.sync_listing_segmented = sync_segmented_mock
|
|
scraper_ctx.__enter__.return_value.sync_listing = sync_listing_mock
|
|
scraper_ctx.__exit__.return_value = None
|
|
|
|
with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx):
|
|
tasks.sync_listing_task.push_request(id="task-beyond")
|
|
try:
|
|
result = tasks.sync_listing_task.run()
|
|
finally:
|
|
tasks.sync_listing_task.pop_request()
|
|
|
|
self.assertEqual(result["status"], "success")
|
|
sync_listing_mock.assert_called_once_with(
|
|
make=None,
|
|
model=None,
|
|
lane="iaai_cars",
|
|
limit=None,
|
|
only_new=False,
|
|
)
|
|
sync_segmented_mock.assert_not_called()
|
|
release_lock.assert_called_once()
|
|
|
|
def test_sync_listing_task_clears_checkpoint_on_hourly_run(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, "_clear_sync_checkpoint") as clear_checkpoint, \
|
|
patch.object(tasks, "_run_browser_job", return_value={
|
|
"run_id": 13,
|
|
"status": "partial_success",
|
|
"full_scan_completed": True,
|
|
"cars_upserted": 2,
|
|
"cars_failed": 1,
|
|
"images_upserted": 3,
|
|
"skipped_existing": 0,
|
|
"elapsed_seconds": 4.0,
|
|
"failures": [{"vehicle_url": "v", "error": "e"}],
|
|
}), \
|
|
patch.object(tasks.sync_listing_task, "update_state"):
|
|
get_persistence.return_value = MagicMock()
|
|
redis_client = MagicMock()
|
|
redis_client.get.return_value = None
|
|
get_redis.return_value = redis_client
|
|
start_heartbeat.return_value = (MagicMock(), MagicMock())
|
|
|
|
tasks.sync_listing_task.push_request(id="task-791")
|
|
try:
|
|
result = tasks.sync_listing_task.run(make="Toyota")
|
|
finally:
|
|
tasks.sync_listing_task.pop_request()
|
|
|
|
self.assertEqual(result["status"], "partial_success")
|
|
# Hourly-ветка (full_scan_done=True) всегда удаляет оставшийся чекпоинт
|
|
# до браузерного job + после. Главное — вызов произошёл.
|
|
clear_checkpoint.assert_called()
|
|
release_lock.assert_called_once()
|
|
|
|
def test_try_set_followup_pending_deduplicates(self) -> None:
|
|
redis_client = MagicMock()
|
|
redis_client.set.side_effect = [True, False]
|
|
|
|
self.assertTrue(tasks._try_set_followup_pending(redis_client, ttl_seconds=120))
|
|
self.assertFalse(tasks._try_set_followup_pending(redis_client, ttl_seconds=120))
|
|
|
|
def test_bump_bootstrap_failure_streak_opens_circuit_breaker(self) -> None:
|
|
redis_client = MagicMock()
|
|
redis_client.incr.return_value = tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT
|
|
|
|
streak, should_enqueue = tasks._bump_bootstrap_failure_streak(
|
|
redis_client,
|
|
reason="max_retries_exceeded",
|
|
)
|
|
|
|
self.assertEqual(streak, tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT)
|
|
self.assertFalse(should_enqueue)
|
|
redis_client.expire.assert_called_once_with(
|
|
tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY,
|
|
tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS,
|
|
)
|
|
|
|
def test_bump_bootstrap_failure_streak_allows_retry_before_limit(self) -> None:
|
|
redis_client = MagicMock()
|
|
redis_client.incr.return_value = 1
|
|
|
|
streak, should_enqueue = tasks._bump_bootstrap_failure_streak(
|
|
redis_client,
|
|
reason="bootstrap_not_completed",
|
|
)
|
|
|
|
self.assertEqual(streak, 1)
|
|
self.assertTrue(should_enqueue)
|
|
|
|
def test_sync_listing_task_does_not_enqueue_followup_after_bootstrap_error_limit(self) -> None:
|
|
settings_stub = make_settings_stub(parallel_segments=False)
|
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
|
patch.object(tasks, "_get_redis") as get_redis, \
|
|
patch.object(tasks, "Settings", return_value=settings_stub), \
|
|
patch.object(tasks, "_acquire_lock", return_value=True), \
|
|
patch.object(tasks, "_is_full_scan_done", return_value=False), \
|
|
patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \
|
|
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
|
|
patch.object(tasks, "_try_set_followup_pending", return_value=True), \
|
|
patch.object(
|
|
tasks,
|
|
"_bump_bootstrap_failure_streak",
|
|
return_value=(tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT, False),
|
|
) as bump_streak, \
|
|
patch.object(tasks, "_clear_followup_pending") as clear_pending, \
|
|
patch.object(tasks.sync_listing_task, "update_state"), \
|
|
patch.object(
|
|
tasks,
|
|
"_run_browser_job",
|
|
return_value={
|
|
"run_id": 99,
|
|
"status": "failed",
|
|
"full_scan_completed": False,
|
|
"cars_upserted": 0,
|
|
"cars_failed": 0,
|
|
"images_upserted": 0,
|
|
"skipped_existing": 0,
|
|
"elapsed_seconds": 1.0,
|
|
"failures": [{"vehicle_url": "listing", "error": "bad resume"}],
|
|
"listing": {"vehicles_collected": 0},
|
|
},
|
|
):
|
|
get_persistence.return_value = MagicMock()
|
|
redis_client = MagicMock()
|
|
redis_client.get.return_value = None
|
|
get_redis.return_value = redis_client
|
|
start_heartbeat.return_value = (MagicMock(), MagicMock())
|
|
|
|
with patch.object(tasks.sync_listing_task, "app", new=MagicMock()) as task_app:
|
|
tasks.sync_listing_task.push_request(id="task-900")
|
|
try:
|
|
result = tasks.sync_listing_task.run()
|
|
finally:
|
|
tasks.sync_listing_task.pop_request()
|
|
|
|
self.assertEqual(result["status"], "failed")
|
|
bump_streak.assert_called_once()
|
|
clear_pending.assert_called()
|
|
task_app.send_task.assert_not_called()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main() |