mirror of
https://github.com/django-q2/django-q2.git
synced 2026-09-26 01:38:11 +08:00
Fix concurrency issue in timeout timer value processing
According to multiprocessing documentation for Value (https://docs.python.org/3/library/multiprocessing.html#multiprocessing.Value) reads and writes are protected with lock when the lock argument is True (the default) or the lock argument is an instance of Lock or RLock. The documentation states that operations like += are not atomic as that involves reading and writing. On the worker side the critical section includes also storing finished task result because the timeout could happen after the task function has finished but before the result has been stored and timer.value has been updated to tell the guard process that the task has been finished. On the guard side the critical section includes all checks done to see if the worker has timed out or died and the actual reincarnation function because the worker could update timer value to -1 (idle) or -2 (recycle) after the guard has seen timer value 0 (timeout) and is going to terminate the worker.
This commit is contained in:
+19
-17
@@ -209,13 +209,14 @@ class Sentinel(object):
|
||||
while not self.stop_event.is_set() or not counter:
|
||||
# Check Workers
|
||||
for p in self.pool:
|
||||
# Are you alive?
|
||||
if not p.is_alive() or p.timer.value == 0:
|
||||
self.reincarnate(p)
|
||||
continue
|
||||
# Decrement timer if work is being done
|
||||
if p.timer.value > 0:
|
||||
p.timer.value -= cycle
|
||||
with p.timer.get_lock():
|
||||
# Are you alive?
|
||||
if not p.is_alive() or p.timer.value == 0:
|
||||
self.reincarnate(p)
|
||||
continue
|
||||
# Decrement timer if work is being done
|
||||
if p.timer.value > 0:
|
||||
p.timer.value -= cycle
|
||||
# Check Monitor
|
||||
if not self.monitor.is_alive():
|
||||
self.reincarnate(self.monitor)
|
||||
@@ -382,16 +383,17 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
|
||||
result = ('{} : {}'.format(e, traceback.format_exc()), False)
|
||||
if error_reporter:
|
||||
error_reporter.report()
|
||||
# Process result
|
||||
task['result'] = result[0]
|
||||
task['success'] = result[1]
|
||||
task['stopped'] = timezone.now()
|
||||
result_queue.put(task)
|
||||
timer.value = -1 # Idle
|
||||
# Recycle
|
||||
if task_count == Conf.RECYCLE:
|
||||
timer.value = -2 # Recycled
|
||||
break
|
||||
with timer.get_lock():
|
||||
# Process result
|
||||
task['result'] = result[0]
|
||||
task['success'] = result[1]
|
||||
task['stopped'] = timezone.now()
|
||||
result_queue.put(task)
|
||||
timer.value = -1 # Idle
|
||||
# Recycle
|
||||
if task_count == Conf.RECYCLE:
|
||||
timer.value = -2 # Recycled
|
||||
break
|
||||
logger.info(_('{} stopped doing work').format(name))
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user