diff --git a/iaai_scraper/worker/celery_app.py b/iaai_scraper/worker/celery_app.py index 0f4294a..c48969e 100644 --- a/iaai_scraper/worker/celery_app.py +++ b/iaai_scraper/worker/celery_app.py @@ -80,7 +80,7 @@ celery_app.conf.update( "task": "iaai_scraper.worker.tasks.sync_listing_task", "schedule": settings.celery.beat_sync_interval_minutes * 60.0, "args": (), - "kwargs": {"limit": settings.celery.beat_sync_limit}, + "kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": False}, "options": {"queue": "scraping"}, } }, @@ -99,6 +99,6 @@ def _on_worker_ready(**kwargs): logger.info("Worker ready — dispatching initial sync_listing task") celery_app.send_task( "iaai_scraper.worker.tasks.sync_listing_task", - kwargs={"limit": settings.celery.beat_sync_limit}, + kwargs={"limit": settings.celery.beat_sync_limit, "only_new": False}, queue="scraping", ) \ No newline at end of file diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index 22217e7..f0e6f66 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -17,6 +17,8 @@ from ..storage.db import PersistenceService logger = logging.getLogger("iaai_scraper.worker.tasks") SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing" +SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done" +SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task" def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0): @@ -162,6 +164,70 @@ def _release_lock_if_owner(redis_client: Redis, key: str, owner_token: str) -> N logger.warning("Failed to release lock %s", key, exc_info=True) +def _has_running_sync_listing_tasks(celery_app) -> bool: + try: + inspector = celery_app.control.inspect(timeout=1.0) + snapshots = [ + inspector.active() or {}, + inspector.reserved() or {}, + inspector.scheduled() or {}, + ] + except Exception: + logger.warning("Failed to inspect Celery workers for running sync tasks", exc_info=True) + return True + + for snapshot in snapshots: + for entries in snapshot.values(): + for entry in entries or []: + task_name = str(entry.get("name") or entry.get("request", {}).get("name") or "") + if task_name == SYNC_LISTING_TASK_NAME: + return True + return False + + +def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app) -> bool: + try: + owner_token = redis_client.get(SYNC_LISTING_LOCK_KEY) + if not owner_token: + return False + except Exception: + logger.warning("Failed to read sync listing lock before cleanup", exc_info=True) + return False + + if _has_running_sync_listing_tasks(celery_app): + logger.info("sync_listing lock preserved: active task still detected") + return False + + try: + ttl = redis_client.ttl(SYNC_LISTING_LOCK_KEY) + redis_client.delete(SYNC_LISTING_LOCK_KEY) + logger.warning( + "Removed orphan sync_listing lock owner=%s ttl=%s after worker restart", + owner_token, + ttl, + ) + return True + except Exception: + logger.warning("Failed to clear orphan sync listing lock", exc_info=True) + return False + + +def _is_full_scan_done(redis_client: Redis) -> bool: + try: + value = redis_client.get(SYNC_FULL_SCAN_DONE_KEY) + except Exception: + logger.warning("Failed to read full scan state", exc_info=True) + return False + return str(value or "").strip() == "1" + + +def _set_full_scan_done(redis_client: Redis, done: bool) -> None: + try: + redis_client.set(SYNC_FULL_SCAN_DONE_KEY, "1" if done else "0") + except Exception: + logger.warning("Failed to persist full scan state", exc_info=True) + + def _start_lock_heartbeat( redis_client: Redis, key: str, @@ -240,9 +306,37 @@ def sync_listing_task( lock_ttl = _sync_listing_lock_ttl_seconds() heartbeat_stop: Event | None = None heartbeat_thread: Thread | None = None + force_bootstrap_full_scan = False + + def _enqueue_bootstrap_followup(reason: str, delay_seconds: int = 5) -> None: + try: + self.app.send_task( + "iaai_scraper.worker.tasks.sync_listing_task", + kwargs={ + "make": make, + "model": model, + "lane": lane, + "limit": limit, + "only_new": only_new, + }, + queue="scraping", + countdown=max(0, int(delay_seconds)), + ) + logger.info( + "Bootstrap follow-up sync queued in %ss (reason=%s)", + delay_seconds, + reason, + ) + except Exception: + logger.warning("Failed to enqueue bootstrap follow-up sync", exc_info=True) lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl) + if not lock_acquired: + orphan_cleared = _clear_orphan_sync_listing_lock(redis_client, self.app) + if orphan_cleared: + lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl) + if not lock_acquired: logger.info("sync_listing_task skipped: another sync is already running") return { @@ -252,6 +346,22 @@ def sync_listing_task( } try: + full_scan_done_before_run = _is_full_scan_done(redis_client) + force_bootstrap_full_scan = not full_scan_done_before_run + effective_limit = None if force_bootstrap_full_scan else limit + effective_only_new = False if force_bootstrap_full_scan else only_new + + if force_bootstrap_full_scan: + logger.info( + "Bootstrap mode: forcing full scan (only_new=False, limit=None) until first complete run", + ) + + logger.info( + "sync_listing options: only_new=%s, limit=%s", + effective_only_new, + effective_limit, + ) + heartbeat_stop, heartbeat_thread = _start_lock_heartbeat( redis_client, SYNC_LISTING_LOCK_KEY, @@ -265,12 +375,22 @@ def sync_listing_task( make=make, model=model, lane=lane, - limit=limit, - only_new=only_new, + limit=effective_limit, + only_new=effective_only_new, ) result = _run_browser_job(_job) + if force_bootstrap_full_scan: + bootstrap_completed = bool(result.get("full_scan_completed")) + if bootstrap_completed: + _set_full_scan_done(redis_client, True) + logger.info("Bootstrap full scan completed; hourly schedule continues") + else: + _set_full_scan_done(redis_client, False) + logger.info("Bootstrap full scan not complete yet; queuing immediate continuation") + _enqueue_bootstrap_followup("bootstrap_not_completed") + summary = { "task_id": task_id, "run_id": result.get("run_id"), @@ -295,22 +415,31 @@ def sync_listing_task( logger.warning( "sync_listing_task soft timeout exceeded — partial progress already saved to DB" ) + if force_bootstrap_full_scan: + _set_full_scan_done(redis_client, False) + _enqueue_bootstrap_followup("soft_time_limit_exceeded") # Partial progress уже записан в БД через finish_sync_run. - # Не retry — следующий beat подхватит новые URL автоматически. + # Не retry — следующий запуск продолжит обработку по расписанию. return { "status": "timed_out", "task_id": task_id, "reason": "soft_time_limit_exceeded", - "note": "partial progress saved to DB; next beat cycle will continue", + "note": "partial progress saved to DB; bootstrap continuation queued", } except Exception as exc: logger.error("sync_listing_task failed: %s", exc, exc_info=True) # Retry только на не-таймаутные ошибки (сеть, БД, браузер). try: + if force_bootstrap_full_scan: + _set_full_scan_done(redis_client, False) + raise self.retry(exc=exc, countdown=5) raise self.retry(exc=exc) except self.MaxRetriesExceededError: logger.error("sync_listing_task max retries exceeded, giving up") + if force_bootstrap_full_scan: + _set_full_scan_done(redis_client, False) + _enqueue_bootstrap_followup("max_retries_exceeded") return { "status": "failed", "task_id": task_id, diff --git a/tests/test_worker_tasks.py b/tests/test_worker_tasks.py index 91d778a..285163f 100644 --- a/tests/test_worker_tasks.py +++ b/tests/test_worker_tasks.py @@ -58,6 +58,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): 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, "_run_browser_job", return_value={ @@ -87,6 +88,74 @@ 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: + 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 + + cleared = tasks._clear_orphan_sync_listing_lock(redis_client, celery_app) + + self.assertTrue(cleared) + 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 + + cleared = tasks._clear_orphan_sync_listing_lock(redis_client, celery_app) + + self.assertFalse(cleared) + redis_client.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, "_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() + 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() + if __name__ == "__main__": unittest.main() \ No newline at end of file