tune vps config and update tests
This commit is contained in:
@@ -8,32 +8,22 @@ from iaai_scraper.worker import tasks
|
||||
|
||||
|
||||
class TestWorkerTaskLockHelpers(unittest.TestCase):
|
||||
def test_acquire_lock_returns_true_on_success(self) -> None:
|
||||
def test_lock_acquire_refresh_release(self) -> None:
|
||||
redis_client = MagicMock()
|
||||
|
||||
# Acquire.
|
||||
redis_client.set.return_value = True
|
||||
|
||||
acquired = tasks._acquire_lock(redis_client, "lock:key", "owner-token", 120)
|
||||
|
||||
self.assertTrue(acquired)
|
||||
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)
|
||||
|
||||
def test_refresh_lock_if_owner_extends_ttl(self) -> None:
|
||||
redis_client = MagicMock()
|
||||
# Refresh.
|
||||
redis_client.eval.return_value = 1
|
||||
self.assertTrue(tasks._refresh_lock_if_owner(redis_client, "lock:key", "owner-token", 120))
|
||||
|
||||
refreshed = tasks._refresh_lock_if_owner(redis_client, "lock:key", "owner-token", 120)
|
||||
|
||||
self.assertTrue(refreshed)
|
||||
redis_client.eval.assert_called_once()
|
||||
|
||||
def test_release_lock_if_owner_uses_owner_token(self) -> None:
|
||||
redis_client = MagicMock()
|
||||
|
||||
# Release.
|
||||
redis_client.eval.reset_mock()
|
||||
tasks._release_lock_if_owner(redis_client, "lock:key", "owner-token")
|
||||
|
||||
redis_client.eval.assert_called_once()
|
||||
args = redis_client.eval.call_args[0]
|
||||
self.assertEqual(args[1], 1)
|
||||
self.assertEqual(args[2], "lock:key")
|
||||
self.assertEqual(args[3], "owner-token")
|
||||
|
||||
@@ -62,6 +52,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
||||
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,
|
||||
@@ -90,7 +81,8 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
||||
heartbeat_thread.join.assert_called_once()
|
||||
release_lock.assert_called_once()
|
||||
|
||||
def test_clear_orphan_sync_listing_lock_deletes_when_no_tasks_running(self) -> None:
|
||||
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
|
||||
@@ -101,27 +93,21 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
||||
inspector.scheduled.return_value = {"worker@node": []}
|
||||
celery_app.control.inspect.return_value = inspector
|
||||
|
||||
cleared = tasks._clear_orphan_sync_listing_lock(redis_client, celery_app)
|
||||
|
||||
self.assertTrue(cleared)
|
||||
self.assertTrue(tasks._clear_orphan_sync_listing_lock(redis_client, celery_app))
|
||||
redis_client.delete.assert_called_once_with(tasks.SYNC_LISTING_LOCK_KEY)
|
||||
|
||||
def test_clear_orphan_sync_listing_lock_keeps_when_task_detected(self) -> None:
|
||||
redis_client = MagicMock()
|
||||
redis_client.get.return_value = "owner-token"
|
||||
celery_app = MagicMock()
|
||||
inspector = MagicMock()
|
||||
inspector.active.return_value = {
|
||||
"worker@node": [{"name": tasks.SYNC_LISTING_TASK_NAME}]
|
||||
}
|
||||
inspector.reserved.return_value = {"worker@node": []}
|
||||
inspector.scheduled.return_value = {"worker@node": []}
|
||||
celery_app.control.inspect.return_value = inspector
|
||||
# Задача активна → не удаляет.
|
||||
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
|
||||
|
||||
cleared = tasks._clear_orphan_sync_listing_lock(redis_client, celery_app)
|
||||
|
||||
self.assertFalse(cleared)
|
||||
redis_client.delete.assert_not_called()
|
||||
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, \
|
||||
@@ -131,6 +117,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
||||
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,
|
||||
@@ -159,28 +146,22 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
||||
self.assertEqual(acquire_lock.call_count, 2)
|
||||
release_lock.assert_called_once()
|
||||
|
||||
def test_sync_listing_checkpoint_roundtrip(self) -> None:
|
||||
def test_sync_listing_checkpoint_save_load_and_resume(self) -> None:
|
||||
# Roundtrip: save → load.
|
||||
redis_client = MagicMock()
|
||||
storage: dict[str, str] = {}
|
||||
redis_client.set.side_effect = lambda key, value: storage.__setitem__(key, value)
|
||||
redis_client.get.side_effect = lambda key: storage.get(key)
|
||||
|
||||
tasks._save_sync_checkpoint(
|
||||
redis_client,
|
||||
task_id="task-1",
|
||||
page_number=12,
|
||||
make=None,
|
||||
model=None,
|
||||
lane="iaai_cars",
|
||||
redis_client, task_id="task-1", page_number=12,
|
||||
make=None, model=None, lane="iaai_cars",
|
||||
)
|
||||
|
||||
checkpoint = tasks._load_sync_checkpoint(redis_client)
|
||||
|
||||
self.assertIsNotNone(checkpoint)
|
||||
self.assertEqual(checkpoint["last_successful_page"], 12)
|
||||
self.assertEqual(checkpoint["status"], "in_progress")
|
||||
|
||||
def test_sync_listing_task_resumes_from_checkpoint_page(self) -> None:
|
||||
# Resume: задача стартует со страницы checkpoint + 1.
|
||||
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), \
|
||||
@@ -188,32 +169,24 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
||||
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", side_effect=lambda fn: fn()):
|
||||
persistence = MagicMock()
|
||||
get_persistence.return_value = persistence
|
||||
redis_client = MagicMock()
|
||||
redis_client.get.side_effect = lambda key: json.dumps({
|
||||
"status": "in_progress",
|
||||
"last_successful_page": 9,
|
||||
"make": None,
|
||||
"model": None,
|
||||
"lane": "iaai_cars",
|
||||
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=[]):
|
||||
get_persistence.return_value = MagicMock()
|
||||
redis_client2 = MagicMock()
|
||||
redis_client2.get.side_effect = lambda key: json.dumps({
|
||||
"status": "in_progress", "last_successful_page": 9,
|
||||
"make": None, "model": None, "lane": "iaai_cars",
|
||||
}) if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None
|
||||
get_redis.return_value = redis_client
|
||||
get_redis.return_value = redis_client2
|
||||
stop_event = MagicMock()
|
||||
heartbeat_thread = MagicMock()
|
||||
start_heartbeat.return_value = (stop_event, heartbeat_thread)
|
||||
|
||||
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": [],
|
||||
"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 = sync_listing_mock
|
||||
|
||||
Reference in New Issue
Block a user