Implement dead job and dead worker cleanups, add discovered node subprocess as anonymous workers
This commit is contained in:
@@ -2,6 +2,7 @@ from __future__ import absolute_import
|
||||
|
||||
import atexit
|
||||
import logging
|
||||
import psutil
|
||||
from multiprocessing import active_children
|
||||
|
||||
from django.db import transaction, DatabaseError
|
||||
@@ -17,8 +18,8 @@ from common.utils import encapsulate
|
||||
from clustering.models import Node
|
||||
from clustering.signals import node_died, node_heartbeat
|
||||
|
||||
from .models import JobQueue, JobProcessingConfig, JobQueueItem
|
||||
from .tasks import job_queue_poll
|
||||
from .models import JobQueue, JobProcessingConfig, JobQueueItem, Worker
|
||||
from .tasks import job_queue_poll, house_keeping
|
||||
from .links import (node_workers, job_queues, tool_link,
|
||||
job_queue_items_pending, job_queue_items_error, job_queue_items_active,
|
||||
job_queue_config_edit, setup_link, job_queue_start, job_queue_stop,
|
||||
@@ -32,6 +33,7 @@ def add_job_queue_jobs():
|
||||
job_processor_scheduler = LocalScheduler('job_processor', _(u'Job processor'))
|
||||
try:
|
||||
job_processor_scheduler.add_interval_job('job_queue_poll', _(u'Poll a job queue for pending jobs.'), job_queue_poll, seconds=JobProcessingConfig.get().job_queue_poll_interval)
|
||||
job_processor_scheduler.add_interval_job('house_keeping', _(u'Poll a job queue for pending jobs.'), house_keeping, seconds=JobProcessingConfig.get().dead_job_removal_interval)
|
||||
except DatabaseError:
|
||||
transaction.rollback()
|
||||
|
||||
@@ -61,15 +63,25 @@ register_model_list_columns(Node, [
|
||||
def process_dead_workers(sender, **kwargs):
|
||||
logger.debug('kwargs')
|
||||
logger.debug(kwargs)
|
||||
#TODO: delete all of the dying node workers and requeue their jobs
|
||||
|
||||
|
||||
@receiver(node_heartbeat, dispatch_uid='node_processes')
|
||||
def node_processes(sender, **kwargs):
|
||||
logger.debug('kwargs')
|
||||
logger.debug(kwargs)
|
||||
print "PROCS"
|
||||
def node_processes(sender, node, **kwargs):
|
||||
logger.debug('update current node\'s processes')
|
||||
pids = [process.pid for process in active_children()]
|
||||
logger.debug('pids: %s' % pids)
|
||||
# Create empty entry for a new unknown worker
|
||||
for process in active_children():
|
||||
print process.pid
|
||||
worker = Worker.objects.get_or_create(
|
||||
node=node, pid=process.pid)
|
||||
|
||||
all_active_pids = psutil.get_pid_list()
|
||||
# Remove stale workers based on current child pids
|
||||
#for dead_worker in Worker.objects.filter(node=node).exclude(pid__in=all_active_pids):
|
||||
for dead_worker in Worker.objects.exclude(pid__in=all_active_pids):
|
||||
dead_worker.delete()
|
||||
# TODO: requeue worker job or delete job?
|
||||
|
||||
|
||||
def kill_all_node_processes():
|
||||
|
||||
@@ -27,3 +27,4 @@ JOB_QUEUE_STATE_CHOICES = (
|
||||
)
|
||||
|
||||
DEFAULT_JOB_QUEUE_POLL_INTERVAL = 2
|
||||
DEFAULT_DEAD_JOB_REMOVAL_INTERVAL = 5
|
||||
|
||||
@@ -0,0 +1,66 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
import datetime
|
||||
from south.db import db
|
||||
from south.v2 import SchemaMigration
|
||||
from django.db import models
|
||||
|
||||
|
||||
class Migration(SchemaMigration):
|
||||
|
||||
def forwards(self, orm):
|
||||
|
||||
# Changing field 'Worker.job_queue_item'
|
||||
db.alter_column('job_processor_worker', 'job_queue_item_id', self.gf('django.db.models.fields.related.ForeignKey')(to=orm['job_processor.JobQueueItem'], null=True))
|
||||
|
||||
def backwards(self, orm):
|
||||
|
||||
# User chose to not deal with backwards NULL issues for 'Worker.job_queue_item'
|
||||
raise RuntimeError("Cannot reverse this migration. 'Worker.job_queue_item' and its values cannot be restored.")
|
||||
|
||||
models = {
|
||||
'clustering.node': {
|
||||
'Meta': {'object_name': 'Node'},
|
||||
'cpuload': ('django.db.models.fields.FloatField', [], {'default': '0', 'blank': 'True'}),
|
||||
'heartbeat': ('django.db.models.fields.DateTimeField', [], {'default': 'datetime.datetime(2012, 8, 3, 0, 0)', 'blank': 'True'}),
|
||||
'hostname': ('django.db.models.fields.CharField', [], {'unique': 'True', 'max_length': '255'}),
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'memory_usage': ('django.db.models.fields.FloatField', [], {'default': '0', 'blank': 'True'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'d'", 'max_length': '4'})
|
||||
},
|
||||
'job_processor.jobprocessingconfig': {
|
||||
'Meta': {'object_name': 'JobProcessingConfig'},
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'job_queue_poll_interval': ('django.db.models.fields.PositiveIntegerField', [], {'default': '2'}),
|
||||
'lock_id': ('django.db.models.fields.CharField', [], {'default': '1', 'unique': 'True', 'max_length': '1'})
|
||||
},
|
||||
'job_processor.jobqueue': {
|
||||
'Meta': {'object_name': 'JobQueue'},
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'name': ('django.db.models.fields.CharField', [], {'unique': 'True', 'max_length': '32'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'r'", 'max_length': '4'}),
|
||||
'unique_jobs': ('django.db.models.fields.BooleanField', [], {'default': 'True'})
|
||||
},
|
||||
'job_processor.jobqueueitem': {
|
||||
'Meta': {'ordering': "('creation_datetime',)", 'object_name': 'JobQueueItem'},
|
||||
'creation_datetime': ('django.db.models.fields.DateTimeField', [], {}),
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'job_queue': ('django.db.models.fields.related.ForeignKey', [], {'to': "orm['job_processor.JobQueue']"}),
|
||||
'job_type': ('django.db.models.fields.CharField', [], {'max_length': '32'}),
|
||||
'kwargs': ('django.db.models.fields.TextField', [], {}),
|
||||
'result': ('django.db.models.fields.TextField', [], {'blank': 'True'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'p'", 'max_length': '4'}),
|
||||
'unique_id': ('django.db.models.fields.CharField', [], {'unique': 'True', 'max_length': '64', 'blank': 'True'})
|
||||
},
|
||||
'job_processor.worker': {
|
||||
'Meta': {'ordering': "('creation_datetime',)", 'object_name': 'Worker'},
|
||||
'creation_datetime': ('django.db.models.fields.DateTimeField', [], {'default': 'datetime.datetime(2012, 8, 3, 0, 0)'}),
|
||||
'heartbeat': ('django.db.models.fields.DateTimeField', [], {'default': 'datetime.datetime(2012, 8, 3, 0, 0)', 'blank': 'True'}),
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'job_queue_item': ('django.db.models.fields.related.ForeignKey', [], {'to': "orm['job_processor.JobQueueItem']", 'null': 'True', 'blank': 'True'}),
|
||||
'node': ('django.db.models.fields.related.ForeignKey', [], {'to': "orm['clustering.Node']"}),
|
||||
'pid': ('django.db.models.fields.PositiveIntegerField', [], {'max_length': '255'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'r'", 'max_length': '4'})
|
||||
}
|
||||
}
|
||||
|
||||
complete_apps = ['job_processor']
|
||||
@@ -0,0 +1,69 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
import datetime
|
||||
from south.db import db
|
||||
from south.v2 import SchemaMigration
|
||||
from django.db import models
|
||||
|
||||
|
||||
class Migration(SchemaMigration):
|
||||
|
||||
def forwards(self, orm):
|
||||
# Adding field 'JobProcessingConfig.dead_job_removal_interval'
|
||||
db.add_column('job_processor_jobprocessingconfig', 'dead_job_removal_interval',
|
||||
self.gf('django.db.models.fields.PositiveIntegerField')(default=5),
|
||||
keep_default=False)
|
||||
|
||||
|
||||
def backwards(self, orm):
|
||||
# Deleting field 'JobProcessingConfig.dead_job_removal_interval'
|
||||
db.delete_column('job_processor_jobprocessingconfig', 'dead_job_removal_interval')
|
||||
|
||||
|
||||
models = {
|
||||
'clustering.node': {
|
||||
'Meta': {'object_name': 'Node'},
|
||||
'cpuload': ('django.db.models.fields.FloatField', [], {'default': '0', 'blank': 'True'}),
|
||||
'heartbeat': ('django.db.models.fields.DateTimeField', [], {'default': 'datetime.datetime(2012, 8, 3, 0, 0)', 'blank': 'True'}),
|
||||
'hostname': ('django.db.models.fields.CharField', [], {'unique': 'True', 'max_length': '255'}),
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'memory_usage': ('django.db.models.fields.FloatField', [], {'default': '0', 'blank': 'True'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'d'", 'max_length': '4'})
|
||||
},
|
||||
'job_processor.jobprocessingconfig': {
|
||||
'Meta': {'object_name': 'JobProcessingConfig'},
|
||||
'dead_job_removal_interval': ('django.db.models.fields.PositiveIntegerField', [], {'default': '5'}),
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'job_queue_poll_interval': ('django.db.models.fields.PositiveIntegerField', [], {'default': '2'}),
|
||||
'lock_id': ('django.db.models.fields.CharField', [], {'default': '1', 'unique': 'True', 'max_length': '1'})
|
||||
},
|
||||
'job_processor.jobqueue': {
|
||||
'Meta': {'object_name': 'JobQueue'},
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'name': ('django.db.models.fields.CharField', [], {'unique': 'True', 'max_length': '32'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'r'", 'max_length': '4'}),
|
||||
'unique_jobs': ('django.db.models.fields.BooleanField', [], {'default': 'True'})
|
||||
},
|
||||
'job_processor.jobqueueitem': {
|
||||
'Meta': {'ordering': "('creation_datetime',)", 'object_name': 'JobQueueItem'},
|
||||
'creation_datetime': ('django.db.models.fields.DateTimeField', [], {}),
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'job_queue': ('django.db.models.fields.related.ForeignKey', [], {'to': "orm['job_processor.JobQueue']"}),
|
||||
'job_type': ('django.db.models.fields.CharField', [], {'max_length': '32'}),
|
||||
'kwargs': ('django.db.models.fields.TextField', [], {}),
|
||||
'result': ('django.db.models.fields.TextField', [], {'blank': 'True'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'p'", 'max_length': '4'}),
|
||||
'unique_id': ('django.db.models.fields.CharField', [], {'unique': 'True', 'max_length': '64', 'blank': 'True'})
|
||||
},
|
||||
'job_processor.worker': {
|
||||
'Meta': {'ordering': "('creation_datetime',)", 'object_name': 'Worker'},
|
||||
'creation_datetime': ('django.db.models.fields.DateTimeField', [], {'default': 'datetime.datetime(2012, 8, 3, 0, 0)'}),
|
||||
'heartbeat': ('django.db.models.fields.DateTimeField', [], {'default': 'datetime.datetime(2012, 8, 3, 0, 0)', 'blank': 'True'}),
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'job_queue_item': ('django.db.models.fields.related.ForeignKey', [], {'to': "orm['job_processor.JobQueueItem']", 'null': 'True', 'blank': 'True'}),
|
||||
'node': ('django.db.models.fields.related.ForeignKey', [], {'to': "orm['clustering.Node']"}),
|
||||
'pid': ('django.db.models.fields.PositiveIntegerField', [], {'max_length': '255'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'r'", 'max_length': '4'})
|
||||
}
|
||||
}
|
||||
|
||||
complete_apps = ['job_processor']
|
||||
@@ -0,0 +1,67 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
import datetime
|
||||
from south.db import db
|
||||
from south.v2 import SchemaMigration
|
||||
from django.db import models
|
||||
|
||||
|
||||
class Migration(SchemaMigration):
|
||||
|
||||
def forwards(self, orm):
|
||||
# Adding unique constraint on 'Worker', fields ['node', 'pid']
|
||||
db.create_unique('job_processor_worker', ['node_id', 'pid'])
|
||||
|
||||
|
||||
def backwards(self, orm):
|
||||
# Removing unique constraint on 'Worker', fields ['node', 'pid']
|
||||
db.delete_unique('job_processor_worker', ['node_id', 'pid'])
|
||||
|
||||
|
||||
models = {
|
||||
'clustering.node': {
|
||||
'Meta': {'object_name': 'Node'},
|
||||
'cpuload': ('django.db.models.fields.FloatField', [], {'default': '0', 'blank': 'True'}),
|
||||
'heartbeat': ('django.db.models.fields.DateTimeField', [], {'default': 'datetime.datetime(2012, 8, 3, 0, 0)', 'blank': 'True'}),
|
||||
'hostname': ('django.db.models.fields.CharField', [], {'unique': 'True', 'max_length': '255'}),
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'memory_usage': ('django.db.models.fields.FloatField', [], {'default': '0', 'blank': 'True'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'d'", 'max_length': '4'})
|
||||
},
|
||||
'job_processor.jobprocessingconfig': {
|
||||
'Meta': {'object_name': 'JobProcessingConfig'},
|
||||
'dead_job_removal_interval': ('django.db.models.fields.PositiveIntegerField', [], {'default': '5'}),
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'job_queue_poll_interval': ('django.db.models.fields.PositiveIntegerField', [], {'default': '2'}),
|
||||
'lock_id': ('django.db.models.fields.CharField', [], {'default': '1', 'unique': 'True', 'max_length': '1'})
|
||||
},
|
||||
'job_processor.jobqueue': {
|
||||
'Meta': {'object_name': 'JobQueue'},
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'name': ('django.db.models.fields.CharField', [], {'unique': 'True', 'max_length': '32'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'r'", 'max_length': '4'}),
|
||||
'unique_jobs': ('django.db.models.fields.BooleanField', [], {'default': 'True'})
|
||||
},
|
||||
'job_processor.jobqueueitem': {
|
||||
'Meta': {'ordering': "('creation_datetime',)", 'object_name': 'JobQueueItem'},
|
||||
'creation_datetime': ('django.db.models.fields.DateTimeField', [], {}),
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'job_queue': ('django.db.models.fields.related.ForeignKey', [], {'to': "orm['job_processor.JobQueue']"}),
|
||||
'job_type': ('django.db.models.fields.CharField', [], {'max_length': '32'}),
|
||||
'kwargs': ('django.db.models.fields.TextField', [], {}),
|
||||
'result': ('django.db.models.fields.TextField', [], {'blank': 'True'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'p'", 'max_length': '4'}),
|
||||
'unique_id': ('django.db.models.fields.CharField', [], {'unique': 'True', 'max_length': '64', 'blank': 'True'})
|
||||
},
|
||||
'job_processor.worker': {
|
||||
'Meta': {'ordering': "('creation_datetime',)", 'unique_together': "(('node', 'pid'),)", 'object_name': 'Worker'},
|
||||
'creation_datetime': ('django.db.models.fields.DateTimeField', [], {'default': 'datetime.datetime(2012, 8, 3, 0, 0)'}),
|
||||
'heartbeat': ('django.db.models.fields.DateTimeField', [], {'default': 'datetime.datetime(2012, 8, 3, 0, 0)', 'blank': 'True'}),
|
||||
'id': ('django.db.models.fields.AutoField', [], {'primary_key': 'True'}),
|
||||
'job_queue_item': ('django.db.models.fields.related.ForeignKey', [], {'to': "orm['job_processor.JobQueueItem']", 'null': 'True', 'blank': 'True'}),
|
||||
'node': ('django.db.models.fields.related.ForeignKey', [], {'to': "orm['clustering.Node']"}),
|
||||
'pid': ('django.db.models.fields.PositiveIntegerField', [], {'max_length': '255'}),
|
||||
'state': ('django.db.models.fields.CharField', [], {'default': "'r'", 'max_length': '4'})
|
||||
}
|
||||
}
|
||||
|
||||
complete_apps = ['job_processor']
|
||||
@@ -23,7 +23,7 @@ from .literals import (JOB_STATE_CHOICES, JOB_STATE_PENDING,
|
||||
JOB_STATE_PROCESSING, JOB_STATE_ERROR, WORKER_STATE_CHOICES,
|
||||
WORKER_STATE_RUNNING, DEFAULT_JOB_QUEUE_POLL_INTERVAL,
|
||||
JOB_QUEUE_STATE_STOPPED, JOB_QUEUE_STATE_STARTED,
|
||||
JOB_QUEUE_STATE_CHOICES)
|
||||
JOB_QUEUE_STATE_CHOICES, DEFAULT_DEAD_JOB_REMOVAL_INTERVAL)
|
||||
from .exceptions import (JobQueuePushError, JobQueueNoPendingJobs,
|
||||
JobQueueAlreadyStarted, JobQueueAlreadyStopped)
|
||||
|
||||
@@ -37,7 +37,9 @@ class Job(object):
|
||||
# Run sync or launch async subprocess
|
||||
# OR launch 2 processes: monitor & actual process
|
||||
node = Node.objects.myself()
|
||||
worker = Worker.objects.create(node=node, pid=os.getpid(), job_queue_item=job_queue_item)
|
||||
worker, created = Worker.objects.get_or_create(node=node, pid=os.getpid())
|
||||
worker.job_queue_item=job_queue_item
|
||||
worker.save()
|
||||
try:
|
||||
transaction.commit_on_success(function)(**loads(job_queue_item.kwargs))
|
||||
#function(**loads(job_queue_item.kwargs))
|
||||
@@ -163,6 +165,15 @@ class JobQueue(models.Model):
|
||||
verbose_name_plural = _(u'job queues')
|
||||
|
||||
|
||||
class JobQueueItemManager(models.Manager):
|
||||
def dead_job_queue_items(self):
|
||||
return self.model.objects.filter(state=JOB_STATE_PROCESSING).filter(worker__isnull=True)
|
||||
|
||||
def check_dead_job_queue_items(self):
|
||||
for job_item in self.dead_job_queue_items():
|
||||
job_item.requeue(force=True, at_top=True)
|
||||
|
||||
|
||||
class JobQueueItem(models.Model):
|
||||
# TODO: add re-queue
|
||||
job_queue = models.ForeignKey(JobQueue, verbose_name=_(u'job queue'))
|
||||
@@ -175,7 +186,9 @@ class JobQueueItem(models.Model):
|
||||
default=JOB_STATE_PENDING,
|
||||
verbose_name=_(u'state'))
|
||||
result = models.TextField(blank=True, verbose_name=_(u'result'))
|
||||
|
||||
|
||||
objects = JobQueueItemManager()
|
||||
|
||||
def __unicode__(self):
|
||||
return self.unique_id
|
||||
|
||||
@@ -215,10 +228,18 @@ class JobQueueItem(models.Model):
|
||||
def is_in_pending_state(self):
|
||||
return self.state == JOB_STATE_PENDING
|
||||
|
||||
def requeue(self):
|
||||
if self.is_in_error_state:
|
||||
def requeue(self, force=False, at_top=False):
|
||||
"""
|
||||
Requeue a job so that it is executed again
|
||||
force: requeue even if job is not in error state
|
||||
at_top: requeue at the top of the file usually for jobs that
|
||||
die and shouldn't be placed at the bottom of the queue
|
||||
"""
|
||||
if self.is_in_error_state or force==True:
|
||||
# TODO: raise exception if not in error state
|
||||
self.state = JOB_STATE_PENDING
|
||||
self.creation_datetime = datetime.datetime.now()
|
||||
if not at_top:
|
||||
self.creation_datetime = datetime.datetime.now()
|
||||
self.save()
|
||||
|
||||
class Meta:
|
||||
@@ -236,7 +257,7 @@ class Worker(models.Model):
|
||||
choices=WORKER_STATE_CHOICES,
|
||||
default=WORKER_STATE_RUNNING,
|
||||
verbose_name=_(u'state'))
|
||||
job_queue_item = models.ForeignKey(JobQueueItem, verbose_name=_(u'job queue item'))
|
||||
job_queue_item = models.ForeignKey(JobQueueItem, blank=True, null=True, verbose_name=_(u'job queue item'))
|
||||
|
||||
def __unicode__(self):
|
||||
return u'%s-%s' % (self.node.hostname, self.pid)
|
||||
@@ -245,10 +266,12 @@ class Worker(models.Model):
|
||||
ordering = ('creation_datetime',)
|
||||
verbose_name = _(u'worker')
|
||||
verbose_name_plural = _(u'workers')
|
||||
unique_together = ('node', 'pid')
|
||||
|
||||
|
||||
class JobProcessingConfig(Singleton):
|
||||
job_queue_poll_interval = models.PositiveIntegerField(verbose_name=(u'job queue poll interval (in seconds)'), default=DEFAULT_JOB_QUEUE_POLL_INTERVAL)
|
||||
dead_job_removal_interval = models.PositiveIntegerField(verbose_name=(u'dead job check and removal interval (in seconds)'), help_text=_(u'Interval of time to check the cluster for and remove unresponsive jobs.'), default=DEFAULT_DEAD_JOB_REMOVAL_INTERVAL)
|
||||
|
||||
def __unicode__(self):
|
||||
return ugettext(u'Job queues configuration')
|
||||
|
||||
@@ -3,9 +3,10 @@ from __future__ import absolute_import
|
||||
import logging
|
||||
|
||||
from lock_manager import Lock, LockError
|
||||
from lock_manager.decorators import simple_locking
|
||||
from clustering.models import Node
|
||||
|
||||
from .models import JobQueue
|
||||
from .models import JobQueue, JobQueueItem
|
||||
from .exceptions import JobQueueNoPendingJobs
|
||||
from .literals import JOB_QUEUE_STATE_STARTED
|
||||
|
||||
@@ -41,4 +42,8 @@ def job_queue_poll():
|
||||
else:
|
||||
logger.debug('CPU load or memory usage over limit')
|
||||
|
||||
|
||||
|
||||
@simple_locking('house_keeping', 10)
|
||||
def house_keeping():
|
||||
logger.debug('starting')
|
||||
JobQueueItem.objects.check_dead_job_queue_items()
|
||||
|
||||
@@ -1,20 +0,0 @@
|
||||
from __future__ import absolute_import
|
||||
|
||||
from django.db import models
|
||||
|
||||
#from .exceptions import AlreadyQueued
|
||||
|
||||
|
||||
class OCRProcessingManager(models.Manager):
|
||||
"""
|
||||
Module manager class to handle adding documents to an OCR queue
|
||||
"""
|
||||
def queue_document(self, document):
|
||||
pass
|
||||
#document_queue = self.model.objects.get(name=queue_name)
|
||||
#if document_queue.queuedocument_set.filter(document_version=document.latest_version):
|
||||
# raise AlreadyQueued
|
||||
|
||||
#document_queue.queuedocument_set.create(document_version=document.latest_version, delay=True)
|
||||
|
||||
#return document_queue
|
||||
@@ -200,7 +200,6 @@ def submit_document_to_queue(request, document, post_submit_redirect=None):
|
||||
|
||||
try:
|
||||
document.submit_for_ocr()
|
||||
#ocr_job_queue.push(ocr_job_type, document_version_pk=document.latest_version.pk)
|
||||
messages.success(request, _(u'Document: %(document)s was added to the OCR queue sucessfully.') % {
|
||||
'document': document})
|
||||
except JobQueuePushError:
|
||||
|
||||
Reference in New Issue
Block a user