Coverage for /usr/local/lib/python3.10/site-packages/opal_common-0.0.0-py3.10.egg/opal_common/synchronization/named_lock.py: 71%

44 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 11:54 +0000

1import asyncio 

2import fcntl 

3import os 

4import time 

5from typing import Optional 

6 

7from opal_common.logger import logger 

8 

9DEFAULT_LOCK_ATTEMPT_INTERVAL = 5.0 

10 

11 

12class NamedLock: 

13 """Creates a a file-lock (can be a normal file or a named pipe / fifo), and 

14 exposes a context manager to try to acquire the lock asynchronously.""" 

15 

16 def __init__( 

17 self, path: str, attempt_interval: float = DEFAULT_LOCK_ATTEMPT_INTERVAL 

18 ): 

19 self._lock_file: str = path 

20 self._lock_file_fd = None 

21 self._attempt_interval = attempt_interval 

22 

23 async def __aenter__(self): 

24 """Using the lock as a context manager will try to acquire the lock 

25 until successful (without timeout)""" 

26 await self.acquire() 

27 return self 

28 

29 async def __aexit__(self, exc_type, exc, tb): 

30 """Releases the lock when exiting the lock context.""" 

31 await self.release() 

32 

33 async def acquire(self, timeout: Optional[int] = None): 

34 """Tries to acquire the lock. 

35 

36 if unsuccessful, will sleep and then try again after the attempt 

37 interval. an optional timeout can be provided to give up before 

38 acquiring the lock, in case we reach timeout, function throws 

39 TimeoutError. 

40 """ 

41 logger.debug( 

42 "[{pid}] trying to acquire lock (lock={lock})", 

43 pid=os.getpid(), 

44 lock=self._lock_file, 

45 ) 

46 start_time = time.time() 

47 while True: 

48 if self._acquire(): 48 ↛ 55line 48 didn't jump to line 55 because the condition on line 48 was always true

49 logger.debug( 

50 "[{pid}] lock acquired! (lock={lock})", 

51 pid=os.getpid(), 

52 lock=self._lock_file, 

53 ) 

54 break 

55 await asyncio.sleep(self._attempt_interval) 

56 # potentially give up due to timeout (if timeout is set) 

57 if timeout is not None and time.time() - start_time > timeout: 

58 raise TimeoutError("could not acquire lock") 

59 

60 async def release(self): 

61 """Releases the lock.""" 

62 logger.debug( 

63 "[{pid}] releasing lock (lock={lock})", 

64 pid=os.getpid(), 

65 lock=self._lock_file, 

66 ) 

67 fd = self._lock_file_fd 

68 self._lock_file_fd = None 

69 fcntl.flock(fd, fcntl.LOCK_UN) 

70 os.close(fd) 

71 

72 @property 

73 def is_locked(self): 

74 """True, if the object holds the file lock.""" 

75 return self._lock_file_fd is not None 

76 

77 def _acquire(self) -> bool: 

78 """Tries to acquire the lock, returns immediately regardless of 

79 success. 

80 

81 returns True if lock was acquired successfully, False otherwise. 

82 """ 

83 fd = os.open(self._lock_file, os.O_RDWR | os.O_CREAT | os.O_TRUNC) 

84 

85 # try to acquire the lock, returns immediately 

86 try: 

87 fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) 

88 except (IOError, OSError): 

89 os.close(fd) 

90 else: 

91 self._lock_file_fd = fd 

92 return self.is_locked