From e4a87e167b5306c3bcd1a7847a795a1c7ac5bfa2 Mon Sep 17 00:00:00 2001 From: kdmukai Date: Wed, 9 Dec 2015 11:43:45 -0600 Subject: [PATCH] Fix for issue referenced in https://github.com/Koed00/django-q/issues/124 --- django_q/brokers/orm.py | 24 +++++++++++++++++------- 1 file changed, 17 insertions(+), 7 deletions(-) diff --git a/django_q/brokers/orm.py b/django_q/brokers/orm.py index f3a39ee..0e65863 100644 --- a/django_q/brokers/orm.py +++ b/django_q/brokers/orm.py @@ -2,10 +2,12 @@ from datetime import timedelta from time import sleep from django.utils import timezone +from django import db +from django.db import transaction from django_q.brokers import Broker from django_q.models import OrmQ -from django_q.conf import Conf +from django_q.conf import Conf, logger def _timeout(): @@ -15,16 +17,23 @@ def _timeout(): class ORM(Broker): @staticmethod def get_connection(list_key=Conf.PREFIX): + if transaction.get_autocommit(): # Only True when not in an atomic block + # Make sure stale connections in the broker thread are explicitly + # closed before attempting DB access. + # logger.debug("Broker thread calling close_old_connections") + db.close_old_connections() + else: + logger.debug("Broker in an atomic transaction") return OrmQ.objects.using(Conf.ORM) def queue_size(self): - return self.connection.filter(key=self.list_key, lock__lte=_timeout()).count() + return self.get_connection().filter(key=self.list_key, lock__lte=_timeout()).count() def lock_size(self): - return self.connection.filter(key=self.list_key, lock__gt=_timeout()).count() + return self.get_connection().filter(key=self.list_key, lock__gt=_timeout()).count() def purge_queue(self): - return self.connection.filter(key=self.list_key).delete() + return self.get_connection().filter(key=self.list_key).delete() def ping(self): return True @@ -38,11 +47,11 @@ class ORM(Broker): self.delete(task_id) def enqueue(self, task): - package = self.connection.create(key=self.list_key, payload=task, lock=_timeout()) + package = self.get_connection().create(key=self.list_key, payload=task, lock=_timeout()) return package.pk def dequeue(self): - tasks = self.connection.filter(key=self.list_key, lock__lt=_timeout())[0:Conf.BULK] + tasks = self.get_connection().filter(key=self.list_key, lock__lt=_timeout())[0:Conf.BULK] if tasks: task_list = [] lock = timezone.now() @@ -58,7 +67,8 @@ class ORM(Broker): return self.purge_queue() def delete(self, task_id): - self.connection.filter(pk=task_id).delete() + self.get_connection().filter(pk=task_id).delete() def acknowledge(self, task_id): return self.delete(task_id) +