Merge pull request #77 from Koed00/dev

Adds stale db connection check before every transaction
This commit is contained in:
Ilan Steemers
2015-09-28 13:49:10 +02:00
+9 -6
View File
@@ -117,8 +117,8 @@ class Sentinel(object):
self.task_queue = Queue(maxsize=Conf.QUEUE_LIMIT) if Conf.QUEUE_LIMIT else Queue() self.task_queue = Queue(maxsize=Conf.QUEUE_LIMIT) if Conf.QUEUE_LIMIT else Queue()
self.result_queue = Queue() self.result_queue = Queue()
self.event_out = Event() self.event_out = Event()
self.monitor = Process() self.monitor = None
self.pusher = Process() self.pusher = None
if start: if start:
self.start() self.start()
@@ -163,7 +163,7 @@ class Sentinel(object):
def reincarnate(self, process): def reincarnate(self, process):
""" """
:param process: the process to reincarnate :param process: the process to reincarnate
:type process: Process :type process: Process or None
""" """
if process == self.monitor: if process == self.monitor:
self.monitor = self.spawn_monitor() self.monitor = self.spawn_monitor()
@@ -276,7 +276,7 @@ class Sentinel(object):
def pusher(task_queue, event, broker=None): def pusher(task_queue, event, broker=None):
""" """
Pulls tasks of the Redis List and puts them in the task queue Pulls tasks of the broker and puts them in the task queue
:type task_queue: multiprocessing.Queue :type task_queue: multiprocessing.Queue
:type event: multiprocessing.Event :type event: multiprocessing.Event
""" """
@@ -318,11 +318,12 @@ def monitor(result_queue, broker=None):
broker = get_broker() broker = get_broker()
name = current_process().name name = current_process().name
logger.info(_("{} monitoring at {}").format(name, current_process().pid)) logger.info(_("{} monitoring at {}").format(name, current_process().pid))
db.close_old_connections()
for task in iter(result_queue.get, 'STOP'): for task in iter(result_queue.get, 'STOP'):
# acknowledge
ack_id = task.pop('ack_id', False) ack_id = task.pop('ack_id', False)
if ack_id: if ack_id:
broker.acknowledge(ack_id) broker.acknowledge(ack_id)
# save the result
save_task(task) save_task(task)
if task['success']: if task['success']:
logger.info(_("Processed [{}]").format(task['name'])) logger.info(_("Processed [{}]").format(task['name']))
@@ -340,7 +341,6 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
""" """
name = current_process().name name = current_process().name
logger.info(_('{} ready for work at {}').format(name, current_process().pid)) logger.info(_('{} ready for work at {}').format(name, current_process().pid))
db.close_old_connections()
task_count = 0 task_count = 0
# Start reading the task queue # Start reading the task queue
for task in iter(task_queue.get, 'STOP'): for task in iter(task_queue.get, 'STOP'):
@@ -360,6 +360,7 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
result = (e, False) result = (e, False)
# We're still going # We're still going
if not result: if not result:
db.close_old_connections()
# execute the payload # execute the payload
timer.value = task['kwargs'].pop('timeout', timeout or 0) # Busy timer.value = task['kwargs'].pop('timeout', timeout or 0) # Busy
try: try:
@@ -388,6 +389,7 @@ def save_task(task):
if not task.get('save', Conf.SAVE_LIMIT > 0) and task['success']: if not task.get('save', Conf.SAVE_LIMIT > 0) and task['success']:
return return
# SAVE LIMIT > 0: Prune database, SAVE_LIMIT 0: No pruning # SAVE LIMIT > 0: Prune database, SAVE_LIMIT 0: No pruning
db.close_old_connections()
try: try:
if task['success'] and 0 < Conf.SAVE_LIMIT <= Success.objects.count(): if task['success'] and 0 < Conf.SAVE_LIMIT <= Success.objects.count():
Success.objects.last().delete() Success.objects.last().delete()
@@ -412,6 +414,7 @@ def scheduler(broker=None):
""" """
if not broker: if not broker:
broker = get_broker() broker = get_broker()
db.close_old_connections()
try: try:
for s in Schedule.objects.exclude(repeats=0).filter(next_run__lt=timezone.now()): for s in Schedule.objects.exclude(repeats=0).filter(next_run__lt=timezone.now()):
args = () args = ()