locks.py 14.8 KB
Newer Older
1
2
3
4
# -*- coding: utf-8 -*-
__author__ = 'Daniel Scheffler'

import time
5
from redis import StrictRedis
6
from redis_semaphore import Semaphore
7
from redis_lock import Lock
Daniel Scheffler's avatar
Daniel Scheffler committed
8
from redis.exceptions import ConnectionError as RedisConnectionError
9
import functools
10
from psutil import virtual_memory
Daniel Scheffler's avatar
Daniel Scheffler committed
11
import timeout_decorator
12
13
14
15

from ..misc.logging import GMS_logger
from ..options.config import GMS_config as CFG

16
try:
17
    redis_conn = StrictRedis(host='localhost', db=0)
18
    redis_conn.keys()  # may raise ConnectionError, e.g., if redis server is not installed or not running
Daniel Scheffler's avatar
Daniel Scheffler committed
19
except RedisConnectionError:
20
21
22
    redis_conn = None


23
24
class MultiSlotLock(Semaphore):
    def __init__(self, name='MultiSlotLock', allowed_slots=1, logger=None, **kwargs):
25
        self.disabled = redis_conn is None or allowed_slots in [None, False]
26
        self.namespace = name
27
28
29
30
        self.allowed_slots = allowed_slots
        self.logger = logger or GMS_logger("RedisLock: '%s'" % name)

        if not self.disabled:
31
            super(MultiSlotLock, self).__init__(client=redis_conn, count=allowed_slots, namespace=name, **kwargs)
32
33
34
35

    def acquire(self, timeout=0, target=None):
        if not self.disabled:
            if self.available_count == 0:
36
                self.logger.info("Waiting for free lock '%s'." % self.namespace)
37

38
            token = super(MultiSlotLock, self).acquire(timeout=timeout, target=target)
39

40
            self.logger.info("Acquired lock '%s'" % self.namespace +
41
42
43
44
45
46
                             ('.' if self.allowed_slots == 1 else ', slot #%s.' % (int(token) + 1)))

            return token

    def release(self):
        if not self.disabled:
47
            token = super(MultiSlotLock, self).release()
48
            if token:
49
                self.logger.info("Released lock '%s'" % self.namespace +
50
51
52
53
                                 ('.' if self.allowed_slots == 1 else ', slot #%s.' % (int(token) + 1)))

    def delete(self):
        if not self.disabled:
54
55
56
            self.client.delete(self.check_exists_key)
            self.client.delete(self.available_key)
            self.client.delete(self.grabbed_key)
57
58


59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
class SharedResourceLock(MultiSlotLock):
    def acquire(self, timeout=0, target=None):
        if not self.disabled:
            token = super(SharedResourceLock, self).acquire(timeout=timeout, target=target)
            self.client.hset(self.grabbed_key_jobID, token, self.current_time)

    def release_all_jobID_tokens(self):
        if not self.disabled:
            for token in self.client.hkeys(self.grabbed_key_jobID):
                self.signal(token)

            self.client.delete(self.grabbed_key_jobID)

    @property
    def grabbed_key_jobID(self):
        return self._get_and_set_key('_grabbed_key_jobID', 'GRABBED_BY_GMSJOB_%s' % CFG.ID)

    def signal(self, token):
        if token is None:
            return None
        with self.client.pipeline() as pipe:
            pipe.multi()
            pipe.hdel(self.grabbed_key, token)
            pipe.hdel(self.grabbed_key_jobID, token)  # only difference to Semaphore.signal()
            pipe.lpush(self.available_key, token)
            pipe.execute()
            return token

    def delete(self):
        if not self.disabled:
            super(SharedResourceLock, self).delete()
            self.client.delete(self.grabbed_key_jobID)


class IOLock(SharedResourceLock):
94
95
    def __init__(self, allowed_slots=1, logger=None, **kwargs):
        super(IOLock, self).__init__(name='IOLock', allowed_slots=allowed_slots, logger=logger, **kwargs)
96
97


98
class ProcessLock(SharedResourceLock):
99
100
    def __init__(self, allowed_slots=1, logger=None, **kwargs):
        super(ProcessLock, self).__init__(name='ProcessLock', allowed_slots=allowed_slots, logger=logger, **kwargs)
101
102


103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
# class MemoryLock(SharedResourceLock):
#     def __init__(self, allowed_slots=1, usage_threshold=80, blocking_threshold=90, logger=None, **kwargs):
#         """
#
#         :param allowed_slots:
#         :param usage_th:        Memory usage threshold above which the memory lock is active (percent).
#         :param logger:
#         :param kwargs:
#         """
#         self.usage_th = usage_threshold
#         self.blocking_th = blocking_threshold
#         self.logger = logger or GMS_logger("RedisLock: '%s'" % self.namespace)
#         self.mem_usage = virtual_memory().percent
#
#         self._acquired = None
#
#         self.check_system_overload()
#
#         if self.mem_usage >= usage_threshold:
#             self.logger.info('Memory usage is %.1f percent. Acquiring memory lock..' % self.mem_usage)
#             self.disabled = False
#         else:
#             self.logger.debug('Memory usage is %.1f percent. Memory lock not needed -> skipped.' % self.mem_usage)
#             self.disabled = True
#
#         super(MemoryLock, self).__init__(name='MemoryLock', allowed_slots=allowed_slots, logger=self.logger,
#                                          disabled=self.disabled, **kwargs)
#
#     @timeout_decorator.timeout(seconds=15 * 60, timeout_exception=TimeoutError,
#                                exception_message='Memory lock could not be acquired after waiting 15 minutes '
#                                                  'due to memory overload.')
#     def check_system_overload(self):
#         logged = False
#         while self.mem_usage >= self.blocking_th:
#             if not logged:
#                 self.logger.info('Memory usage is %.1f percent. Waiting until memory usage is below %s percent.'
#                                  % (self.mem_usage, self.blocking_th))
#                 logged = True
#             time.sleep(5)
#             self.mem_usage = virtual_memory().percent
#
#     def acquire(self, blocking=True, timeout=None):
#         if self.mem_usage > self.usage_th:
#             # self._acquired = super(MemoryLock, self).acquire(blocking=blocking, timeout=timeout)
#             self._acquired = super(MemoryLock, self).acquire(timeout=timeout)
#         else:
#             self._acquired = True
#
#         return self._acquired
#
#     def release(self):
#         if self._acquired and not self.disabled:
#             super(MemoryLock, self).release()


158
class MemoryLock(SharedResourceLock):
159
    def __init__(self, mem2lock_gb, allowed_slots=1, usage_threshold=80, blocking_threshold=90, logger=None, **kwargs):
160
        """
161

162
163
164
165
        :param usage_th:        Memory usage threshold above which the memory lock is active (percent).
        :param logger:
        :param kwargs:
        """
166
        super(MemoryLock, self).__init__(name='MemoryLock', allowed_slots=allowed_slots, logger=logger, **kwargs)
167
168

        self.mem2lock_gb = mem2lock_gb
169
        self.usage_th = usage_threshold
Daniel Scheffler's avatar
Daniel Scheffler committed
170
171
        self.blocking_th = blocking_threshold
        self.mem_usage = virtual_memory().percent
172

173
174
175
176
177
178
179
180
181
182
183
184
    @property
    def mem_reserved_gb(self):
        val = redis_conn.get('GMS_mem_reserved')
        if val is not None:
            return val
        return 0

    @property
    def usable_memory_gb(self):
        return int((virtual_memory().total * self.usage_th / 100 - virtual_memory().used) / 1024**3) \
               - int(self.mem_reserved_gb)

185
    @timeout_decorator.timeout(seconds=15 * 60, timeout_exception=TimeoutError,
Daniel Scheffler's avatar
Daniel Scheffler committed
186
187
188
189
190
191
192
                               exception_message='Memory lock could not be acquired after waiting 15 minutes '
                                                 'due to memory overload.')
    def check_system_overload(self):
        logged = False
        while self.mem_usage >= self.blocking_th:
            if not logged:
                self.logger.info('Memory usage is %.1f percent. Waiting until memory usage is below %s percent.'
Daniel Scheffler's avatar
Bugfix.    
Daniel Scheffler committed
193
                                 % (self.mem_usage, self.blocking_th))
Daniel Scheffler's avatar
Bugfix.    
Daniel Scheffler committed
194
                logged = True
Daniel Scheffler's avatar
Daniel Scheffler committed
195
196
            time.sleep(5)
            self.mem_usage = virtual_memory().percent
197

198
199
200
    @timeout_decorator.timeout(seconds=15 * 60, timeout_exception=TimeoutError,
                               exception_message='Memory lock could not be acquired after waiting 15 minutes '
                                                 'because not enough memory could reserved.')
201
    def ensure_enough_memory(self):
202
203
        if not self.disabled:
            logged = False
204

205
            with Lock(self.client, 'GMS_mem_checker'):
206
207
208
209
210
                while self.mem2lock_gb < self.usable_memory_gb:
                    if not logged:
                        self.logger.info('Currently usable memory: %s GB. Waiting until at least %s GB are usable.'
                                         % (self.usable_memory_gb, self.mem2lock_gb))
                        logged = True
211
                    time.sleep(1)
212

213
                self.client.incr('GMS_mem_reserved', self.mem2lock_gb)
214

215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
    def acquire(self, blocking=True, timeout=None):
        # TODO vielleicht auch als float redis value per increment decrement + lock mit 1 slot
        if not self.disabled:
            with Lock(self.client, 'GMS_mem_checker'):
                self.check_system_overload()
                self.ensure_enough_memory()

                if self.usage_th < self.mem_usage < self.blocking_th:
                    self.acquire(blocking=blocking, timeout=timeout)

            # # TODO add an acquire lock here
            # with MultiSlotLock(name='GMS_mem_acquire_key', allowed_slots=1):
            #     if self.usable_memory_gb >= self.mem2lock_gb:
            #         for i in range(self.mem2lock_gb):
            #             super(MemoryLock, self).acquire(timeout=timeout)
            #         self.client.incr('GMS_mem_reserved', self.mem2lock_gb)
            #
            #     else:
            #         time.sleep(1)
            #         self.acquire(blocking=blocking, timeout=timeout)
235

236
237
    def release(self):
        if not self.disabled:
238
239
240
            # for token in self._local_tokens:
            #     self.signal(token)
            super(MemoryLock, self).release()
241
            self.client.decr('GMS_mem_reserved', self.mem2lock_gb)
242
243


244
245
246
247
248
249
def acquire_process_lock(**processlock_kwargs):
    """Decorator function for ProcessLock.

    :param processlock_kwargs:  Keywourd arguments to be passed to ProcessLock class.
    """

250
    def decorator(func):
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
        @functools.wraps(func)  # needed to avoid pickling errors
        def wrapped_func(*args, **kwargs):
            with ProcessLock(**processlock_kwargs):
                result = func(*args, **kwargs)

            return result

        return wrapped_func

    return decorator


def acquire_mem_lock(**memlock_kwargs):
    """Decorator function for MEMLock.

Daniel Scheffler's avatar
Daniel Scheffler committed
266
    :param memlock_kwargs:  Keywourd arguments to be passed to MemoryLock class.
267
    """
268

269
    def decorator(func):
270
271
        @functools.wraps(func)  # needed to avoid pickling errors
        def wrapped_func(*args, **kwargs):
Daniel Scheffler's avatar
Daniel Scheffler committed
272
            with MemoryLock(**memlock_kwargs):
273
274
275
276
277
278
279
280
281
                result = func(*args, **kwargs)

            return result

        return wrapped_func

    return decorator


282
def release_unclosed_locks():
283
    if redis_conn:
284
285
286
287
288
289
290
        for L in [IOLock, ProcessLock]:
            lock = L(allowed_slots=1)
            lock.release_all_jobID_tokens()

            # delete the complete redis namespace if no lock slot is acquired anymore
            if lock.client.hlen(lock.grabbed_key) == 0:
                lock.delete()
291

292
293
        ML = MemoryReserver(1)
        ML.release_all_jobID_tokens()
294

295
296
297
298
299
300
301
        # delete the complete redis namespace if no lock slot is acquired anymore
        if ML.client.hlen(ML.grabbed_key) == 0:
            ML.delete()


class MemoryReserver(Semaphore):
    def __init__(self, mem2lock_gb, max_usage=90, logger=None, **kwargs):
302
303
304
305
        """

        :param reserved_mem:    Amount of memory to be reserved during the lock is acquired (gigabytes).
        """
306
307
308
309
310
311
312
313
314
315
316
317
318
        self.disabled = redis_conn is None
        self.mem2lock_gb = mem2lock_gb
        self.max_usage = max_usage
        self._waiting = False

        if not self.disabled:
            mem_limit = int(virtual_memory().total * max_usage / 100 / 1024**3)
            super(MemoryReserver, self).__init__(client=redis_conn, count=mem_limit, namespace='MemoryReserver',
                                                 **kwargs)
            self.logger = logger or GMS_logger("RedisLock: 'MemoryReserver'")

    @property
    def mem_reserved_gb(self):
319
        return int(redis_conn.get('GMS_mem_reserved') or 0)
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347

    @property
    def usable_memory_gb(self):
        return int((virtual_memory().total * self.max_usage / 100 - virtual_memory().used) / 1024**3) \
               - int(self.mem_reserved_gb)

    @property
    def grabbed_key_jobID(self):
        return self._get_and_set_key('_grabbed_key_jobID', 'GRABBED_BY_GMSJOB_%s' % CFG.ID)

    @property
    def reserved_key(self):
        return self._get_and_set_key('_reserved_key', 'MEM_RESERVED')

    @property
    def reserved_key_jobID(self):
        return self._get_and_set_key('_reserved_key_jobID', 'MEM_RESERVED_BY_GMSJOB_%s' % CFG.ID)

    def acquire(self, timeout=0, target=None):
        if not self.disabled:
            with Lock(self.client, 'GMS_mem_acquire_lock'):
                if self.usable_memory_gb >= self.mem2lock_gb:
                    for i in range(self.mem2lock_gb):
                        token = super(MemoryReserver, self).acquire(timeout=timeout)
                        self.client.hset(self.grabbed_key_jobID, token, self.current_time)

                    self.client.incr(self.reserved_key, self.mem2lock_gb)
                    self.client.incr(self.reserved_key_jobID, self.mem2lock_gb)
348

349
350
                    self.logger.info('Reserved %s GB of memory.' % self.mem2lock_gb)
                    self._waiting = False
351

352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
                else:
                    if not self._waiting:
                        self.logger.info('Currently usable memory: %s GB. Waiting until at least %s GB are usable.'
                                         % (self.usable_memory_gb, self.mem2lock_gb))
                        self._waiting = True

                    time.sleep(1)
                    self.acquire(timeout=timeout)

    def release(self):
        if not self.disabled:
            for token in self._local_tokens:
                self.signal(token)
            self.client.decr(self.reserved_key, self.mem2lock_gb)
            self.client.decr(self.reserved_key_jobID, self.mem2lock_gb)

            self.logger.info('Released %s GB of reserved memory.' % self.mem2lock_gb)

    def release_all_jobID_tokens(self):
        mem_reserved = int(redis_conn.get(self.reserved_key_jobID) or 0)
        if mem_reserved:
            redis_conn.decr(self.reserved_key, mem_reserved)

        redis_conn.delete(self.reserved_key_jobID)

        for token in self.client.hkeys(self.grabbed_key_jobID):
            self.signal(token)

        self.client.delete(self.grabbed_key_jobID)

    def delete(self):
        if not self.disabled:
            self.client.delete(self.check_exists_key)
            self.client.delete(self.available_key)
            self.client.delete(self.grabbed_key)
387

388
389
            if self.mem_reserved_gb <= 0:
                self.client.delete(self.reserved_key)