Merge pull request #1542 from jmchilton/did_fix

[16.01] Attempt to fix #1531.
This commit is contained in:
Martin Cech
2016-01-26 16:24:42 -05:00
3 changed files with 72 additions and 7 deletions
+18 -7
View File
@@ -402,23 +402,34 @@ class DeleteIntermediatesAction(DefaultJobAction):
# POTENTIAL ISSUES: When many outputs are being finish()ed
# concurrently, sometimes non-terminal steps won't be cleaned up
# because of the lag in job state updates.
sa_session.flush()
if not job.workflow_invocation_step:
log.debug("This job is not part of a workflow invocation, delete intermediates aborted.")
return
wfi = job.workflow_invocation_step.workflow_invocation
sa_session.refresh(wfi)
if wfi.active:
log.debug("Workflow still scheduling so new jobs may appear, skipping deletion of intermediate files.")
# Still evaluating workflow so we don't yet have all workflow invocation
# steps to start looking at.
return
if wfi.workflow.has_outputs_defined():
jobs_to_check = [wfistep.job for wfistep in wfi.steps if not wfistep.workflow_step.workflow_outputs]
outputs_defined = wfi.workflow.has_outputs_defined()
if outputs_defined:
wfi_steps = [wfistep for wfistep in wfi.steps if not wfistep.workflow_step.workflow_outputs and wfistep.workflow_step.type == "tool"]
jobs_to_check = []
for wfi_step in wfi_steps:
sa_session.refresh(wfi_step)
wfi_step_job = wfi_step.job
if wfi_step_job:
jobs_to_check.append(wfi_step_job)
else:
log.debug("No job found yet for wfi_step %s, (step %s)" % (wfi_step, wfi_step.workflow_step))
for j2c in jobs_to_check:
if j2c is None:
# Job not yet created, this will be re-evaluated after subsequent jobs in
# workflow.
return
for input_dataset in [x.dataset for x in j2c.input_datasets if x.dataset.creating_job.workflow_invocation_step and x.dataset.creating_job.workflow_invocation_step.workflow_invocation == wfi]:
creating_jobs = [(x, x.dataset.creating_job) for x in j2c.input_datasets if x.dataset.creating_job]
for (x, creating_job) in creating_jobs:
sa_session.refresh(creating_job)
sa_session.refresh(x)
for input_dataset in [x.dataset for (x, creating_job) in creating_jobs if creating_job.workflow_invocation_step and creating_job.workflow_invocation_step.workflow_invocation == wfi]:
safe_to_delete = True
for job_to_check in [d_j.job for d_j in input_dataset.dependent_jobs]:
if job_to_check != job and job_to_check.state not in [job.states.OK, job.states.DELETED]:
+44
View File
@@ -1236,6 +1236,50 @@ steps:
content = self.dataset_populator.get_history_dataset_details( history_id )
assert content[ "name" ] == "foo was replaced", content[ "name" ]
@skip_without_tool( "cat1" )
def test_delete_intermediate_datasets_pja_1( self ):
history_id = self.dataset_populator.new_history()
self._run_jobs("""
class: GalaxyWorkflow
inputs:
- id: input1
outputs:
- id: wf_output_1
source: third_cat#out_file1
steps:
- tool_id: cat1
label: first_cat
state:
input1:
$link: input1
- tool_id: cat1
label: second_cat
state:
input1:
$link: first_cat#out_file1
- tool_id: cat1
label: third_cat
state:
input1:
$link: second_cat#out_file1
outputs:
out_file1:
delete_intermediate_datasets: true
test_data:
input1: "hello world"
""", history_id=history_id)
hda1 = self.dataset_populator.get_history_dataset_details(history_id, hid=1)
hda2 = self.dataset_populator.get_history_dataset_details(history_id, hid=2)
hda3 = self.dataset_populator.get_history_dataset_details(history_id, hid=3)
hda4 = self.dataset_populator.get_history_dataset_details(history_id, hid=4)
assert not hda1["deleted"]
assert hda2["deleted"]
# I think hda3 should be deleted, but the inputs to
# steps with workflow outputs are not deleted.
# assert hda3["deleted"]
print hda3["deleted"]
assert not hda4["deleted"]
@skip_without_tool( "random_lines1" )
def test_run_replace_params_by_tool( self ):
workflow_request, history_id = self._setup_random_x2_workflow( "test_for_replace_tool_params" )
+10
View File
@@ -379,6 +379,16 @@ def transform_tool(context, step):
)
post_job_actions[action_name] = action
if output.get("delete_intermediate_datasets", None):
action_name = "DeleteIntermediatesAction%s" % name
arguments = dict()
action = __action(
"DeleteIntermediatesAction",
name,
arguments,
)
post_job_actions[action_name] = action
del step["outputs"]