diff --git a/luigi/worker.py b/luigi/worker.py index 43e5dc0c0f..71cefe6a0c 100644 --- a/luigi/worker.py +++ b/luigi/worker.py @@ -209,12 +209,7 @@ def run(self): with self._forward_attributes(): new_deps = self._run_get_new_deps() if not new_deps: - if not self.check_complete_on_run: - # update the cache - if self.task_completion_cache is not None: - self.task_completion_cache[self.task.task_id] = True - status = DONE - elif self.check_complete(self.task): + if not self.check_complete_on_run or self.check_complete(self.task): status = DONE else: raise TaskException("Task finished running, but complete() is still returning false.") diff --git a/test/worker_test.py b/test/worker_test.py index 750f05b933..99ff6994e4 100644 --- a/test/worker_test.py +++ b/test/worker_test.py @@ -506,9 +506,9 @@ def run(self): w.run() # the complete methods of a's yielded first in b's run method were called equally often self.assertEqual(b0.complete_count, 1) - self.assertEqual(a0.complete_count, 2) - self.assertEqual(a1.complete_count, 2) - self.assertEqual(a2.complete_count, 2) + self.assertEqual(a0.complete_count, 3) + self.assertEqual(a1.complete_count, 3) + self.assertEqual(a2.complete_count, 3) # test with disabled cache_task_completion with Worker(scheduler=self.sch, worker_id='2', cache_task_completion=False) as w: