Coverage for /usr/local/lib/python3.10/site-packages/opal_common-0.0.0-py3.10.egg/opal_common/synchronization/expiring_redis_lock.py: 0%
18 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 11:54 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 11:54 +0000
1import asyncio
3import redis.asyncio as redis
4from opal_common.logger import logger
7async def run_locked(
8 _redis: redis.Redis, lock_name: str, coro: asyncio.coroutine, timeout: int = 10
9):
10 """This function runs a coroutine wrapped in a redis lock, in a way that
11 prevents hanging locks. Hanging locks can happen when a process crashes
12 while holding a lock.
14 This function sets a redis enforced timeout, and reacquires the lock
15 every timeout * 0.8 (as long as it runs)
16 """
17 lock = _redis.lock(lock_name, timeout=timeout)
18 try:
19 logger.debug(f"Trying to acquire redis lock: {lock_name}")
20 await lock.acquire()
21 logger.debug(f"Acquired lock: {lock_name}")
23 locked_task = asyncio.create_task(coro)
25 while True:
26 done, _ = await asyncio.wait(
27 (locked_task,),
28 timeout=timeout * 0.8,
29 return_when=asyncio.FIRST_COMPLETED,
30 )
31 if locked_task in done:
32 break
33 else:
34 # Extend lock timeout as long as the coroutine is still running
35 await lock.reacquire()
36 logger.debug(f"Reacquired lock: {lock_name}")
38 finally:
39 await lock.release()
40 logger.debug(f"Released lock: {lock_name}")