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

1import asyncio 

2 

3import redis.asyncio as redis 

4from opal_common.logger import logger 

5 

6 

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. 

13 

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}") 

22 

23 locked_task = asyncio.create_task(coro) 

24 

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}") 

37 

38 finally: 

39 await lock.release() 

40 logger.debug(f"Released lock: {lock_name}")