summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--worker.py6
1 files 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)