2013-10-17 22:25:19 +00:00
|
|
|
from datetime import datetime, timedelta
|
|
|
|
|
2013-11-15 21:44:33 +00:00
|
|
|
from data.database import QueueItem, db
|
2014-02-16 23:59:24 +00:00
|
|
|
from app import app
|
|
|
|
|
|
|
|
|
|
|
|
transaction_factory = app.config['DB_TRANSACTION_FACTORY']
|
2013-10-17 22:25:19 +00:00
|
|
|
|
|
|
|
|
|
|
|
class WorkQueue(object):
|
|
|
|
def __init__(self, queue_name):
|
|
|
|
self.queue_name = queue_name
|
|
|
|
|
2013-11-15 20:49:26 +00:00
|
|
|
def put(self, message, available_after=0, retries_remaining=5):
|
2013-10-17 22:25:19 +00:00
|
|
|
"""
|
|
|
|
Put an item, if it shouldn't be processed for some number of seconds,
|
|
|
|
specify that amount as available_after.
|
|
|
|
"""
|
|
|
|
|
|
|
|
params = {
|
|
|
|
'queue_name': self.queue_name,
|
|
|
|
'body': message,
|
2013-11-15 20:49:26 +00:00
|
|
|
'retries_remaining': retries_remaining,
|
2013-10-17 22:25:19 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
if available_after:
|
|
|
|
available_date = datetime.now() + timedelta(seconds=available_after)
|
|
|
|
params['available_after'] = available_date
|
|
|
|
|
|
|
|
QueueItem.create(**params)
|
|
|
|
|
2013-10-18 19:28:16 +00:00
|
|
|
def get(self, processing_time=300):
|
2013-10-17 22:25:19 +00:00
|
|
|
"""
|
|
|
|
Get an available item and mark it as unavailable for the default of five
|
|
|
|
minutes.
|
|
|
|
"""
|
|
|
|
now = datetime.now()
|
2013-10-20 07:06:11 +00:00
|
|
|
available_or_expired = ((QueueItem.available == True) |
|
|
|
|
(QueueItem.processing_expires <= now))
|
2013-10-17 22:25:19 +00:00
|
|
|
|
2014-02-16 23:59:24 +00:00
|
|
|
with transaction_factory(db):
|
2013-11-15 21:44:33 +00:00
|
|
|
avail = QueueItem.select().where(QueueItem.queue_name == self.queue_name,
|
|
|
|
QueueItem.available_after <= now,
|
|
|
|
available_or_expired,
|
|
|
|
QueueItem.retries_remaining > 0)
|
2013-10-17 22:25:19 +00:00
|
|
|
|
2013-11-15 21:44:33 +00:00
|
|
|
found = list(avail.limit(1).order_by(QueueItem.available_after))
|
2013-10-17 22:25:19 +00:00
|
|
|
|
2013-11-15 21:44:33 +00:00
|
|
|
if found:
|
|
|
|
item = found[0]
|
|
|
|
item.available = False
|
|
|
|
item.processing_expires = now + timedelta(seconds=processing_time)
|
|
|
|
item.retries_remaining -= 1
|
|
|
|
item.save()
|
2013-10-17 22:25:19 +00:00
|
|
|
|
2013-11-15 21:44:33 +00:00
|
|
|
return item
|
2013-10-17 22:25:19 +00:00
|
|
|
|
2013-11-15 21:44:33 +00:00
|
|
|
return None
|
2013-10-17 22:25:19 +00:00
|
|
|
|
2013-10-18 19:28:16 +00:00
|
|
|
def complete(self, completed_item):
|
2013-10-18 21:27:09 +00:00
|
|
|
completed_item.delete_instance()
|
|
|
|
|
2013-11-16 20:05:26 +00:00
|
|
|
def incomplete(self, incomplete_item, retry_after=300):
|
|
|
|
retry_date = datetime.now() + timedelta(seconds=retry_after)
|
|
|
|
incomplete_item.available_after = retry_date
|
2013-10-29 19:42:19 +00:00
|
|
|
incomplete_item.available = True
|
|
|
|
incomplete_item.save()
|
|
|
|
|
2013-10-18 21:27:09 +00:00
|
|
|
|
|
|
|
image_diff_queue = WorkQueue('imagediff')
|
2014-02-12 17:51:30 +00:00
|
|
|
dockerfile_build_queue = WorkQueue('dockerfilebuild2')
|
2013-11-15 20:49:26 +00:00
|
|
|
webhook_queue = WorkQueue('webhook')
|