Skip to content

Commit cec5dce

Browse files
committed
pipeliner: refactor handling task removal in 'deferred' pipeline failure mode.
1 parent 95cccd8 commit cec5dce

1 file changed

Lines changed: 9 additions & 23 deletions

File tree

auto_process_ngs/pipeliner.py

Lines changed: 9 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -2293,7 +2293,7 @@ def run(self,working_dir=None,tasks_work_dir=None,log_dir=None,
22932293
# Check requirements
22942294
run_task = reduce(lambda x,y: x and
22952295
(y not in self._running) and
2296-
(y.completed or not y.active) and
2296+
y.completed and
22972297
(y.exit_code == 0 or task.always_run()),
22982298
requirements,True)
22992299
if run_task:
@@ -2421,16 +2421,20 @@ def run(self,working_dir=None,tasks_work_dir=None,log_dir=None,
24212421
# the removal operations
24222422
pending = []
24232423
for t in self._pending:
2424-
task = t[0]
2424+
task,requires,kws = t
24252425
if task.id() in remove_ids:
2426-
task.disable()
2427-
msg = "-- disabling dependent task '%s'" \
2426+
msg = "-- removing dependent task '%s'" \
24282427
% task.name()
24292428
if verbose:
24302429
msg += " (%s)" % task.id()
24312430
self.report(msg)
24322431
else:
2433-
pending.append(t)
2432+
# Update the requirements for the
2433+
# task and add back into 'pending'
2434+
for req in requires:
2435+
if req.id() in remove_ids:
2436+
task.drop_required_task(req.id())
2437+
pending.append(self.get_task(task.id()))
24342438
self._pending = pending
24352439
self._failed.extend(failed)
24362440
self.report("There are failed tasks but pipeline exit "
@@ -2536,7 +2540,6 @@ def __init__(self,_name,*args,**kws):
25362540
self._task_name = "%s.%s" % (sanitize_name(self._name),
25372541
uuid.uuid4())
25382542
self._completed = False
2539-
self._active = True
25402543
self._stdout_files = []
25412544
self._stderr_files = []
25422545
self._exit_code = 0
@@ -2620,17 +2623,6 @@ def completed(self):
26202623
"""
26212624
return self._completed
26222625

2623-
@property
2624-
def active(self):
2625-
"""
2626-
Check if the task is marked as active
2627-
2628-
Returns:
2629-
Boolean: True if task is marked as active,
2630-
False if not
2631-
"""
2632-
return self._active
2633-
26342626
@property
26352627
def updated(self):
26362628
"""
@@ -2762,12 +2754,6 @@ def fail(self,exit_code=1,message=None):
27622754
self._exit_code = exit_code
27632755
self._completed = True
27642756

2765-
def disable(self):
2766-
"""
2767-
Register the task as disabled
2768-
"""
2769-
self._active = False
2770-
27712757
def report(self,s):
27722758
"""
27732759
Internal: report messages from the task

0 commit comments

Comments
 (0)