diff --git a/CHANGES/7902.bugfix b/CHANGES/7902.bugfix new file mode 100644 index 0000000000..19a604af84 --- /dev/null +++ b/CHANGES/7902.bugfix @@ -0,0 +1 @@ +Fixed RedisWorker leaking Redis locks when task acquire, abort, cancel, or immediate dispatch cleanup failed. diff --git a/pulpcore/tasking/redis_locks.py b/pulpcore/tasking/redis_locks.py index 7e6231fb08..e38fdbc933 100644 --- a/pulpcore/tasking/redis_locks.py +++ b/pulpcore/tasking/redis_locks.py @@ -388,6 +388,7 @@ def release_resource_locks( except redis.RedisError as e: _logger.error("Error releasing locks: %s", e) + raise async def async_release_resource_locks( @@ -439,3 +440,4 @@ async def async_release_resource_locks( except redis.RedisError as e: _logger.error("Error releasing locks: %s", e) + raise diff --git a/pulpcore/tasking/redis_tasks.py b/pulpcore/tasking/redis_tasks.py index f8fc95b49e..ebe1de63b5 100644 --- a/pulpcore/tasking/redis_tasks.py +++ b/pulpcore/tasking/redis_tasks.py @@ -7,6 +7,7 @@ import contextvars import logging +import time import redis from asgiref.sync import sync_to_async @@ -19,6 +20,7 @@ async_safe_release_task_locks, extract_task_resources, get_task_lock_key, + release_resource_locks, safe_release_task_locks, ) from pulpcore.tasking.tasks import ( @@ -103,6 +105,61 @@ def clear_cancel_signal(task_id): _logger.error("Error clearing cancellation signal for task %s: %s", task_id, e) +def _release_task_locks_any_owner(task): + redis_conn = get_redis_connection() + if redis_conn is None: + return + task_lock_key = get_task_lock_key(task.pk) + try: + lock_owner = redis_conn.get(task_lock_key) + except redis.RedisError as e: + _logger.error("Error reading task lock for %s: %s", task.pk, e) + return + if not lock_owner: + return + if isinstance(lock_owner, bytes): + lock_owner = lock_owner.decode() + exclusive_resources, shared_resources = extract_task_resources(task) + try: + release_resource_locks( + redis_conn, lock_owner, task_lock_key, exclusive_resources, shared_resources + ) + except redis.RedisError as e: + _logger.error("Error releasing locks for canceled task %s: %s", task.pk, e) + + +def _retry_safe_release_task_locks(task, lock_owner, attempts=3): + for attempt in range(attempts): + try: + safe_release_task_locks(task, lock_owner) + return + except redis.RedisError: + if attempt < attempts - 1: + time.sleep(0.1 * (attempt + 1)) + else: + _logger.error( + "Failed to release locks for task %s after %s attempts.", + task.pk, + attempts, + ) + + +async def _aretry_safe_release_task_locks(task, lock_owner, attempts=3): + for attempt in range(attempts): + try: + await async_safe_release_task_locks(task, lock_owner) + return + except redis.RedisError: + if attempt < attempts - 1: + await sync_to_async(time.sleep)(0.1 * (attempt + 1)) + else: + _logger.error( + "Failed to release locks for task %s after %s attempts.", + task.pk, + attempts, + ) + + def cancel_task(task_id): """ Cancel a task using Redis-based signaling. @@ -142,6 +199,7 @@ def cancel_task(task_id): Task.objects.filter(pk=task.pk).update(app_lock=AppStatus.objects.current()) task.app_lock = AppStatus.objects.current() task.set_canceled() + _release_task_locks_any_owner(task) else: # Task is RUNNING — signal the supervising worker. publish_cancel_signal(task.pk) @@ -390,7 +448,7 @@ def dispatch( except Exception: # Release locks if using_workdir() failed before # execute_task() had a chance to run and release them - safe_release_task_locks(task, lock_owner) + _retry_safe_release_task_locks(task, lock_owner) raise elif deferred: # Locks not available, defer to worker @@ -442,7 +500,7 @@ async def adispatch( except Exception: # Release locks if using_workdir() failed before # aexecute_task() had a chance to run and release them - await async_safe_release_task_locks(task, lock_owner) + await _aretry_safe_release_task_locks(task, lock_owner) raise elif deferred: # Locks not available, defer to worker diff --git a/pulpcore/tasking/redis_worker.py b/pulpcore/tasking/redis_worker.py index afde0a4f29..f908a91131 100644 --- a/pulpcore/tasking/redis_worker.py +++ b/pulpcore/tasking/redis_worker.py @@ -551,6 +551,10 @@ def fetch_task(self): except Exception as e: _logger.error("Error processing task %s: %s", task.pk, e) + try: + safe_release_task_locks(task, lock_owner=self.name) + except Exception: + pass continue if len(waiting_tasks) < fetch_limit: @@ -668,25 +672,24 @@ def supervise_task(self, task): if cancel_state: from pulpcore.tasking._util import delete_incomplete_resources - # Reload task from database to get current state - task.refresh_from_db() - # Only clean up if task is not already in a final state - # (subprocess may have already handled cancellation) - if task.state not in TASK_FINAL_STATES: - # Release locks BEFORE setting canceled state - # Atomically release task lock + resource locks in a single operation + try: + # Reload task from database to get current state + task.refresh_from_db() + # Only clean up if task is not already in a final state + # (subprocess may have already handled cancellation) + if task.state not in TASK_FINAL_STATES: + task.set_canceling() + _logger.info( + "Cleaning up task %s in domain: %s and marking as %s.", + task.pk, + domain.name, + cancel_state, + ) + delete_incomplete_resources(task) + task.set_canceled(final_state=cancel_state, reason=cancel_reason) + finally: self._maybe_release_locks(task) - task.set_canceling() - _logger.info( - "Cleaning up task %s in domain: %s and marking as %s.", - task.pk, - domain.name, - cancel_state, - ) - delete_incomplete_resources(task) - task.set_canceled(final_state=cancel_state, reason=cancel_reason) - self.task = None def handle_tasks(self): @@ -714,9 +717,7 @@ def handle_tasks(self): finally: # Safety net: if _execute_task() crashed before releasing locks, # atomically release all locks here (task lock + resource locks) - # NOTE: Only for immediate tasks that execute in this process. - # Deferred tasks execute in subprocess which handles its own lock release. - if task and task.immediate: + if task: self._maybe_release_locks(task) def sleep(self):