Merge pull request #337 from jannero/timer-handling

Fix concurrency issue in timeout timer value processing
This commit is contained in:
Ilan Steemers
2019-02-11 13:28:50 +01:00
committed by GitHub
+19 -17
View File
@@ -209,13 +209,14 @@ class Sentinel(object):
while not self.stop_event.is_set() or not counter: while not self.stop_event.is_set() or not counter:
# Check Workers # Check Workers
for p in self.pool: for p in self.pool:
# Are you alive? with p.timer.get_lock():
if not p.is_alive() or p.timer.value == 0: # Are you alive?
self.reincarnate(p) if not p.is_alive() or p.timer.value == 0:
continue self.reincarnate(p)
# Decrement timer if work is being done continue
if p.timer.value > 0: # Decrement timer if work is being done
p.timer.value -= cycle if p.timer.value > 0:
p.timer.value -= cycle
# Check Monitor # Check Monitor
if not self.monitor.is_alive(): if not self.monitor.is_alive():
self.reincarnate(self.monitor) 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) result = ('{} : {}'.format(e, traceback.format_exc()), False)
if error_reporter: if error_reporter:
error_reporter.report() error_reporter.report()
# Process result with timer.get_lock():
task['result'] = result[0] # Process result
task['success'] = result[1] task['result'] = result[0]
task['stopped'] = timezone.now() task['success'] = result[1]
result_queue.put(task) task['stopped'] = timezone.now()
timer.value = -1 # Idle result_queue.put(task)
# Recycle timer.value = -1 # Idle
if task_count == Conf.RECYCLE: # Recycle
timer.value = -2 # Recycled if task_count == Conf.RECYCLE:
break timer.value = -2 # Recycled
break
logger.info(_('{} stopped doing work').format(name)) logger.info(_('{} stopped doing work').format(name))