From 542df395ce186ccc28d2c4122b7b12024ade8f47 Mon Sep 17 00:00:00 2001 From: schneefux Date: Fri, 10 Mar 2017 22:14:00 +0100 Subject: bugfix: do not commit failed jobs as finished --- worker.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/worker.py b/worker.py index 1e16237..e5f0b25 100644 --- a/worker.py +++ b/worker.py @@ -50,9 +50,11 @@ class Worker(object): await self._windup() critical_error = False failed = [] + succeeded = [] for jobid, payload, priority in jobs: try: await self._execute_job(jobid, payload, priority) + succeeded.append(jobid) except JobFailed as err: failed.append((jobid, err.args[0])) if len(err.args) > 1: @@ -63,10 +65,10 @@ class Worker(object): else: critical_error = True # default if critical_error: - await self._queue.reset([j[0] for j in jobs]) + await self._queue.reset(succeeded) logging.warning("batch failed, reset") else: - await self._queue.finish([j[0] for j in jobs]) + await self._queue.finish(succeeded) for jobid, err in failed: await self._queue.fail(jobid, err) -- cgit v1.3.1