diff --git a/lib/galaxy/managers/workflows.py b/lib/galaxy/managers/workflows.py index 0de745c801f..9ff58ab61e0 100644 --- a/lib/galaxy/managers/workflows.py +++ b/lib/galaxy/managers/workflows.py @@ -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 = {} diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index c4bfd5523ef..f9e1ddce76b 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -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, diff --git a/lib/galaxy/workflow/run.py b/lib/galaxy/workflow/run.py index f5415db15ba..d129fcf3f2b 100644 --- a/lib/galaxy/workflow/run.py +++ b/lib/galaxy/workflow/run.py @@ -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 diff --git a/test/api/test_workflows.py b/test/api/test_workflows.py index f2cb8c41bf9..35cd29addbf 100644 --- a/test/api/test_workflows.py +++ b/test/api/test_workflows.py @@ -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 ) diff --git a/test/api/test_workflows_from_yaml.py b/test/api/test_workflows_from_yaml.py index 3613f7e6c67..865459f3c85 100644 --- a/test/api/test_workflows_from_yaml.py +++ b/test/api/test_workflows_from_yaml.py @@ -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 diff --git a/test/api/yaml_to_workflow.py b/test/api/yaml_to_workflow.py index 2344a462c96..0d63deeef1a 100644 --- a/test/api/yaml_to_workflow.py +++ b/test/api/yaml_to_workflow.py @@ -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