fix worker lock
This commit is contained in:
@@ -80,7 +80,7 @@ celery_app.conf.update(
|
|||||||
"task": "iaai_scraper.worker.tasks.sync_listing_task",
|
"task": "iaai_scraper.worker.tasks.sync_listing_task",
|
||||||
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||||
"args": (),
|
"args": (),
|
||||||
"kwargs": {"limit": settings.celery.beat_sync_limit},
|
"kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": False},
|
||||||
"options": {"queue": "scraping"},
|
"options": {"queue": "scraping"},
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
@@ -99,6 +99,6 @@ def _on_worker_ready(**kwargs):
|
|||||||
logger.info("Worker ready — dispatching initial sync_listing task")
|
logger.info("Worker ready — dispatching initial sync_listing task")
|
||||||
celery_app.send_task(
|
celery_app.send_task(
|
||||||
"iaai_scraper.worker.tasks.sync_listing_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",
|
queue="scraping",
|
||||||
)
|
)
|
||||||
@@ -17,6 +17,8 @@ from ..storage.db import PersistenceService
|
|||||||
logger = logging.getLogger("iaai_scraper.worker.tasks")
|
logger = logging.getLogger("iaai_scraper.worker.tasks")
|
||||||
|
|
||||||
SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
|
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):
|
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)
|
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(
|
def _start_lock_heartbeat(
|
||||||
redis_client: Redis,
|
redis_client: Redis,
|
||||||
key: str,
|
key: str,
|
||||||
@@ -240,9 +306,37 @@ def sync_listing_task(
|
|||||||
lock_ttl = _sync_listing_lock_ttl_seconds()
|
lock_ttl = _sync_listing_lock_ttl_seconds()
|
||||||
heartbeat_stop: Event | None = None
|
heartbeat_stop: Event | None = None
|
||||||
heartbeat_thread: Thread | 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)
|
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:
|
if not lock_acquired:
|
||||||
logger.info("sync_listing_task skipped: another sync is already running")
|
logger.info("sync_listing_task skipped: another sync is already running")
|
||||||
return {
|
return {
|
||||||
@@ -252,6 +346,22 @@ def sync_listing_task(
|
|||||||
}
|
}
|
||||||
|
|
||||||
try:
|
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(
|
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
|
||||||
redis_client,
|
redis_client,
|
||||||
SYNC_LISTING_LOCK_KEY,
|
SYNC_LISTING_LOCK_KEY,
|
||||||
@@ -265,12 +375,22 @@ def sync_listing_task(
|
|||||||
make=make,
|
make=make,
|
||||||
model=model,
|
model=model,
|
||||||
lane=lane,
|
lane=lane,
|
||||||
limit=limit,
|
limit=effective_limit,
|
||||||
only_new=only_new,
|
only_new=effective_only_new,
|
||||||
)
|
)
|
||||||
|
|
||||||
result = _run_browser_job(_job)
|
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 = {
|
summary = {
|
||||||
"task_id": task_id,
|
"task_id": task_id,
|
||||||
"run_id": result.get("run_id"),
|
"run_id": result.get("run_id"),
|
||||||
@@ -295,22 +415,31 @@ def sync_listing_task(
|
|||||||
logger.warning(
|
logger.warning(
|
||||||
"sync_listing_task soft timeout exceeded — partial progress already saved to DB"
|
"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.
|
# Partial progress уже записан в БД через finish_sync_run.
|
||||||
# Не retry — следующий beat подхватит новые URL автоматически.
|
# Не retry — следующий запуск продолжит обработку по расписанию.
|
||||||
return {
|
return {
|
||||||
"status": "timed_out",
|
"status": "timed_out",
|
||||||
"task_id": task_id,
|
"task_id": task_id,
|
||||||
"reason": "soft_time_limit_exceeded",
|
"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:
|
except Exception as exc:
|
||||||
logger.error("sync_listing_task failed: %s", exc, exc_info=True)
|
logger.error("sync_listing_task failed: %s", exc, exc_info=True)
|
||||||
# Retry только на не-таймаутные ошибки (сеть, БД, браузер).
|
# Retry только на не-таймаутные ошибки (сеть, БД, браузер).
|
||||||
try:
|
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)
|
raise self.retry(exc=exc)
|
||||||
except self.MaxRetriesExceededError:
|
except self.MaxRetriesExceededError:
|
||||||
logger.error("sync_listing_task max retries exceeded, giving up")
|
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 {
|
return {
|
||||||
"status": "failed",
|
"status": "failed",
|
||||||
"task_id": task_id,
|
"task_id": task_id,
|
||||||
|
|||||||
@@ -58,6 +58,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|||||||
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
with patch.object(tasks, "_get_persistence") as get_persistence, \
|
||||||
patch.object(tasks, "_get_redis") as get_redis, \
|
patch.object(tasks, "_get_redis") as get_redis, \
|
||||||
patch.object(tasks, "_acquire_lock", return_value=True), \
|
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, "_start_lock_heartbeat") as start_heartbeat, \
|
||||||
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
|
patch.object(tasks, "_release_lock_if_owner") as release_lock, \
|
||||||
patch.object(tasks, "_run_browser_job", return_value={
|
patch.object(tasks, "_run_browser_job", return_value={
|
||||||
@@ -87,6 +88,74 @@ class TestWorkerTaskLockHelpers(unittest.TestCase):
|
|||||||
heartbeat_thread.join.assert_called_once()
|
heartbeat_thread.join.assert_called_once()
|
||||||
release_lock.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__":
|
if __name__ == "__main__":
|
||||||
unittest.main()
|
unittest.main()
|
||||||
Reference in New Issue
Block a user