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
« 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
7from opal_common.logger import logger
9DEFAULT_LOCK_ATTEMPT_INTERVAL = 5.0
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."""
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
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
29 async def __aexit__(self, exc_type, exc, tb):
30 """Releases the lock when exiting the lock context."""
31 await self.release()
33 async def acquire(self, timeout: Optional[int] = None):
34 """Tries to acquire the lock.
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")
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)
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
77 def _acquire(self) -> bool:
78 """Tries to acquire the lock, returns immediately regardless of
79 success.
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)
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