mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Allow implicit connections between workflow steps (no editor GUI yet).
This means steps that are not connecting an output of one step to the input of another. This could potentially address all sorts of untraditional (in a Galaxy sense) workflows where some sort of data is managed externally. The most important use I think I have heard discussed is that of data managers - this can be used in cases where data managers depend on one another (grab the fasta files in one step, index them in another) or workflows where a downstream analysis depends on index data populated via data managers in earlier steps. Not really sure how to represent these in the workflow editor - but this is a power user feature anyway so hopefully that is not super pressing. The YAML to workflow DSL supports the operation (see test cases) so these power users (a euphemism for Dan I guess) can just use that for now.
This commit is contained in:
@@ -529,7 +529,7 @@ class WorkflowContentsManager(UsesAnnotations):
|
||||
visit_input_values( module.tool.inputs, module.state.inputs, callback )
|
||||
# Filter
|
||||
# FIXME: this removes connection without displaying a message currently!
|
||||
input_connections = [ conn for conn in input_connections if conn.input_name in data_input_names ]
|
||||
input_connections = [ conn for conn in input_connections if (conn.input_name in data_input_names or conn.non_data_connection) ]
|
||||
|
||||
# Encode input connections as dictionary
|
||||
input_conn_dict = {}
|
||||
|
||||
@@ -3077,6 +3077,12 @@ class WorkflowStep( object ):
|
||||
|
||||
|
||||
class WorkflowStepConnection( object ):
|
||||
# Constant used in lieu of output_name and input_name to indicate an
|
||||
# implicit connection between two steps that is not dependent on a dataset
|
||||
# or a dataset collection. Allowing for instance data manager steps to setup
|
||||
# index data before a normal tool runs or for workflows that manage data
|
||||
# outside of Galaxy.
|
||||
NON_DATA_CONNECTION = "__NO_INPUT_OUTPUT_NAME__"
|
||||
|
||||
def __init__( self ):
|
||||
self.output_step_id = None
|
||||
@@ -3084,6 +3090,15 @@ class WorkflowStepConnection( object ):
|
||||
self.input_step_id = None
|
||||
self.input_name = None
|
||||
|
||||
def set_non_data_connection(self):
|
||||
self.output_name = WorkflowStepConnection.NON_DATA_CONNECTION
|
||||
self.input_name = WorkflowStepConnection.NON_DATA_CONNECTION
|
||||
|
||||
@property
|
||||
def non_data_connection(self):
|
||||
return (self.output_name == WorkflowStepConnection.NON_DATA_CONNECTION and
|
||||
self.input_name == WorkflowStepConnection.NON_DATA_CONNECTION)
|
||||
|
||||
|
||||
class WorkflowOutput(object):
|
||||
|
||||
@@ -3153,6 +3168,13 @@ class WorkflowInvocation( object, Dictifiable ):
|
||||
step_invocations[ step_id ].append( invocation_step )
|
||||
return step_invocations
|
||||
|
||||
def step_invocations_for_step_id( self, step_id ):
|
||||
step_invocations = []
|
||||
for invocation_step in self.steps:
|
||||
if step_id == invocation_step.workflow_step_id:
|
||||
step_invocations.append( invocation_step )
|
||||
return step_invocations
|
||||
|
||||
@staticmethod
|
||||
def poll_active_workflow_ids(
|
||||
sa_session,
|
||||
|
||||
@@ -136,6 +136,8 @@ class WorkflowInvoker( object ):
|
||||
for step in remaining_steps:
|
||||
jobs = None
|
||||
try:
|
||||
self.__check_implicitly_dependent_steps(step)
|
||||
|
||||
jobs = self._invoke_step( step )
|
||||
for job in (util.listify( jobs ) or [None]):
|
||||
# Record invocation
|
||||
@@ -160,6 +162,38 @@ class WorkflowInvoker( object ):
|
||||
# invocations.
|
||||
return self.progress.outputs
|
||||
|
||||
def __check_implicitly_dependent_steps( self, step ):
|
||||
""" Method will delay the workflow evaluation if implicitly dependent
|
||||
steps (steps dependent but not through an input->output way) are not
|
||||
yet complete.
|
||||
"""
|
||||
for input_connection in step.input_connections:
|
||||
if input_connection.non_data_connection:
|
||||
output_id = input_connection.output_step.id
|
||||
self.__check_implicitly_dependent_step( output_id )
|
||||
|
||||
def __check_implicitly_dependent_step( self, output_id ):
|
||||
step_invocations = self.workflow_invocation.step_invocations_for_step_id( output_id )
|
||||
|
||||
# No steps created yet - have to delay evaluation.
|
||||
if not step_invocations:
|
||||
raise modules.DelayedWorkflowEvaluation()
|
||||
|
||||
for step_invocation in step_invocations:
|
||||
job = step_invocation.job
|
||||
if job:
|
||||
# At least one job in incomplete.
|
||||
if not job.finished:
|
||||
raise modules.DelayedWorkflowEvaluation()
|
||||
|
||||
if job.state != job.states.OK:
|
||||
raise modules.CancelWorkflowEvaluation()
|
||||
|
||||
else:
|
||||
# TODO: Handle implicit dependency on stuff like
|
||||
# pause steps.
|
||||
pass
|
||||
|
||||
def _invoke_step( self, step ):
|
||||
jobs = step.module.execute( self.trans, self.progress, self.workflow_invocation, step )
|
||||
return jobs
|
||||
|
||||
@@ -107,7 +107,7 @@ class BaseWorkflowsApiTestCase( api.ApiTestCase ):
|
||||
invocation_details = invocation_details_response.json()
|
||||
return invocation_details
|
||||
|
||||
def _run_jobs( self, jobs_yaml, history_id=None ):
|
||||
def _run_jobs( self, jobs_yaml, history_id=None, wait=True ):
|
||||
if history_id is None:
|
||||
history_id = self.history_id
|
||||
workflow_id = self._upload_yaml_workflow(
|
||||
@@ -153,8 +153,9 @@ class BaseWorkflowsApiTestCase( api.ApiTestCase ):
|
||||
invocation_id = invocation[ "id" ]
|
||||
# Wait for workflow to become fully scheduled and then for all jobs
|
||||
# complete.
|
||||
self.wait_for_invocation( workflow_id, invocation_id )
|
||||
self.dataset_populator.wait_for_history( history_id, assert_ok=True )
|
||||
if wait:
|
||||
self.wait_for_invocation( workflow_id, invocation_id )
|
||||
self.dataset_populator.wait_for_history( history_id, assert_ok=True )
|
||||
jobs = self._history_jobs( history_id )
|
||||
return RunJobsSummary(
|
||||
history_id=history_id,
|
||||
@@ -516,6 +517,60 @@ class WorkflowsApiTestCase( BaseWorkflowsApiTestCase ):
|
||||
invocation = self._invocation_details( uploaded_workflow_id, invocation_id )
|
||||
assert invocation[ 'state' ] == 'cancelled'
|
||||
|
||||
def test_run_with_implicit_connection( self ):
|
||||
history_id = self.dataset_populator.new_history()
|
||||
run_summary = self._run_jobs("""
|
||||
steps:
|
||||
- label: test_input
|
||||
type: input
|
||||
- label: first_cat
|
||||
tool_id: cat1
|
||||
state:
|
||||
input1:
|
||||
$link: test_input
|
||||
- label: the_pause
|
||||
type: pause
|
||||
connect:
|
||||
input:
|
||||
- first_cat#out_file1
|
||||
- label: second_cat
|
||||
tool_id: cat1
|
||||
state:
|
||||
input1:
|
||||
$link: the_pause
|
||||
- label: third_cat
|
||||
tool_id: random_lines1
|
||||
connect:
|
||||
$step: second_cat
|
||||
state:
|
||||
num_lines: 1
|
||||
input:
|
||||
$link: test_input
|
||||
seed_source:
|
||||
seed_source_selector: set_seed
|
||||
seed: asdf
|
||||
__current_case__: 1
|
||||
test_data:
|
||||
test_input: "hello world"
|
||||
""", history_id=history_id, wait=False)
|
||||
time.sleep( 2 )
|
||||
history_id = run_summary.history_id
|
||||
workflow_id = run_summary.workflow_id
|
||||
invocation_id = run_summary.workflow_id
|
||||
self.dataset_populator.wait_for_history( history_id, assert_ok=True )
|
||||
invocation = self._invocation_details( workflow_id, invocation_id )
|
||||
assert invocation[ 'state' ] != 'scheduled'
|
||||
# Expect two jobs - the upload and first cat. randomlines shouldn't run
|
||||
# it is implicitly dependent on second cat.
|
||||
assert len( self._history_jobs( history_id ) ) == 2
|
||||
|
||||
self.__review_paused_steps( workflow_id, invocation_id, order_index=2, action=True )
|
||||
self.wait_for_invocation( workflow_id, invocation_id )
|
||||
time.sleep(1)
|
||||
self.dataset_populator.wait_for_history( history_id, assert_ok=True )
|
||||
time.sleep(1)
|
||||
assert len( self._history_jobs( history_id ) ) == 4
|
||||
|
||||
def test_cannot_run_inaccessible_workflow( self ):
|
||||
workflow = self.workflow_populator.load_workflow( name="test_for_run_cannot_access" )
|
||||
workflow_request, history_id = self._setup_workflow_run( workflow )
|
||||
|
||||
@@ -102,3 +102,34 @@ test_data:
|
||||
print self._get("workflows/%s/download" % workflow_id).json()
|
||||
assert False
|
||||
# TODO: fill out test...
|
||||
|
||||
def test_implicit_connections( self ):
|
||||
workflow_id = self._upload_yaml_workflow("""
|
||||
- label: test_input
|
||||
type: input
|
||||
- label: first_cat
|
||||
tool_id: cat1
|
||||
state:
|
||||
input1:
|
||||
$link: test_input
|
||||
- label: the_pause
|
||||
type: pause
|
||||
connect:
|
||||
input:
|
||||
- first_cat#out_file1
|
||||
- label: second_cat
|
||||
tool_id: cat1
|
||||
state:
|
||||
input1:
|
||||
$link: the_pause
|
||||
- label: third_cat
|
||||
tool_id: cat1
|
||||
connect:
|
||||
$step: second_cat
|
||||
state:
|
||||
input1:
|
||||
$link: test_input
|
||||
""")
|
||||
workflow = self._get("workflows/%s/download" % workflow_id).json()
|
||||
print workflow
|
||||
assert False
|
||||
|
||||
@@ -262,6 +262,8 @@ def __populate_input_connections(context, step, connect):
|
||||
values = [ values ]
|
||||
for value in values:
|
||||
if not isinstance(value, dict):
|
||||
if key == "$step":
|
||||
value += "#__NO_INPUT_OUTPUT_NAME__"
|
||||
value_parts = str(value).split("#")
|
||||
if len(value_parts) == 1:
|
||||
value_parts.append("output")
|
||||
@@ -270,6 +272,8 @@ def __populate_input_connections(context, step, connect):
|
||||
id = context.labels[id]
|
||||
value = {"id": int(id), "output_name": value_parts[1]}
|
||||
input_connection_value.append(value)
|
||||
if key == "$step":
|
||||
key = "__NO_INPUT_OUTPUT_NAME__"
|
||||
input_connections[key] = input_connection_value
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user