@@ -366,7 +366,6 @@ def wait(self, poll_interval=1, timeout=None, full_data=False):
366366 return job_data
367367
368368 time .sleep (poll_interval )
369-
370369 raise Exception ("Waited for job result for %s seconds, timeout." % timeout )
371370
372371 def save_retry (self , retry_exc ):
@@ -452,29 +451,6 @@ def _save_status(self, status, updates=None, exception=False, w=None, j=None):
452451 if self .stored is False and self .statuses_no_storage is not None and status in self .statuses_no_storage :
453452 return
454453
455- with context .connections .redis .pipeline (transaction = False ) as pipe :
456- queue = (updates or {}).get ("queue" ) or self .data ["queue" ]
457- if status != "started" :
458- # Queue change
459- if queue != self .data ["queue" ]:
460- pipe .decr ("queuesize:%s" % self .data ["queue" ])
461- if status == "queued" :
462- pipe .incr ("queuesize:%s" % queue )
463-
464- # Regular queues
465- elif status == "queued" and self .data .get ("status" ) != "started" :
466- pipe .incr ("queuesize:%s" % queue )
467-
468- elif status != "queued" and not self .data .get ("raw_queue" ):
469- pipe .decr ("queuesize:%s" % queue )
470-
471- # Raw queues retries
472- elif (updates or {}).get ("retry_count" , 0 ) > self .data .get ("retry_count" , 0 ):
473- pipe .incr ("queuesize:%s" % queue )
474-
475- pipe .expire ("queuesize:%s" % queue , context .get_current_config ().get ("queue_ttl" ))
476- pipe .execute ()
477-
478454 now = datetime .datetime .utcnow ()
479455 db_updates = {
480456 "status" : status ,
@@ -536,6 +512,28 @@ def _save_status(self, status, updates=None, exception=False, w=None, j=None):
536512 if exception :
537513 self ._save_traceback_history (status , trace , exc )
538514
515+ with context .connections .redis .pipeline (transaction = False ) as pipe :
516+ queue = (updates or {}).get ("queue" ) or self .data ["queue" ]
517+ if status != "started" :
518+ # Queue change
519+ if queue != self .data ["queue" ]:
520+ pipe .decr ("queuesize:%s" % self .data ["queue" ])
521+ if status == "queued" :
522+ pipe .incr ("queuesize:%s" % queue )
523+
524+ # Regular queues
525+ elif status == "queued" and self .data .get ("status" ) != "started" :
526+ pipe .incr ("queuesize:%s" % queue )
527+
528+ elif status != "queued" and not self .data .get ("raw_queue" ):
529+ pipe .decr ("queuesize:%s" % queue )
530+
531+ # Raw queues retries
532+ elif (updates or {}).get ("retry_count" , 0 ) > self .data .get ("retry_count" , 0 ):
533+ pipe .incr ("queuesize:%s" % queue )
534+
535+ pipe .expire ("queuesize:%s" % queue , context .get_current_config ().get ("queue_ttl" ))
536+ pipe .execute ()
539537
540538 def set_current_io (self , io_data ):
541539
0 commit comments