Skip to content

Commit f43a891

Browse files
author
maelorn
committed
fix data inconsistencies
1 parent 2474c24 commit f43a891

1 file changed

Lines changed: 17 additions & 11 deletions

File tree

mrq/job.py

Lines changed: 17 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -484,6 +484,13 @@ def _save_status(self, status, updates=None, exception=False, w=None, j=None):
484484
db_updates["traceback"] = trace
485485
db_updates["exceptiontype"] = exc.__name__
486486

487+
# get all data before updating them
488+
current_queue = (db_updates or {}).get("queue") or self.data["queue"]
489+
old_queue = self.data.get("queue")
490+
old_status = self.data.get("status")
491+
raw_queue = self.data.get("raw_queue")
492+
retry_count = self.data.get("retry_count", 0)
493+
487494
if self.data:
488495
self.data.update(db_updates)
489496

@@ -513,26 +520,25 @@ def _save_status(self, status, updates=None, exception=False, w=None, j=None):
513520
self._save_traceback_history(status, trace, exc)
514521

515522
with context.connections.redis.pipeline(transaction=False) as pipe:
516-
queue = (updates or {}).get("queue") or self.data["queue"]
517523
if status != "started":
518524
# Queue change
519-
if queue != self.data["queue"]:
520-
pipe.decr("queuesize:%s" % self.data["queue"])
525+
if current_queue != old_queue:
526+
pipe.decr("queuesize:%s" % old_queue)
521527
if status == "queued":
522-
pipe.incr("queuesize:%s" % queue)
528+
pipe.incr("queuesize:%s" % current_queue)
523529

524530
# Regular queues
525-
elif status == "queued" and self.data.get("status") != "started":
526-
pipe.incr("queuesize:%s" % queue)
531+
elif status == "queued" and old_status != "started":
532+
pipe.incr("queuesize:%s" % current_queue)
527533

528-
elif status != "queued" and not self.data.get("raw_queue"):
529-
pipe.decr("queuesize:%s" % queue)
534+
elif status != "queued" and not raw_queue:
535+
pipe.decr("queuesize:%s" % current_queue)
530536

531537
# Raw queues retries
532-
elif (updates or {}).get("retry_count", 0) > self.data.get("retry_count", 0):
533-
pipe.incr("queuesize:%s" % queue)
538+
elif (db_updates or {}).get("retry_count", 0) > retry_count:
539+
pipe.incr("queuesize:%s" % current_queue)
534540

535-
pipe.expire("queuesize:%s" % queue, context.get_current_config().get("queue_ttl"))
541+
pipe.expire("queuesize:%s" % current_queue, context.get_current_config().get("queue_ttl"))
536542
pipe.execute()
537543

538544
def set_current_io(self, io_data):

0 commit comments

Comments
 (0)