From e0a5e82bdae407535b9d7c98e3dcf851b63d01a0 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 15 Dec 2014 22:20:11 -0500 Subject: [PATCH] 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. --- lib/galaxy/managers/workflows.py | 2 +- lib/galaxy/model/__init__.py | 22 ++++++++++ lib/galaxy/workflow/run.py | 34 ++++++++++++++++ test/api/test_workflows.py | 61 ++++++++++++++++++++++++++-- test/api/test_workflows_from_yaml.py | 31 ++++++++++++++ test/api/yaml_to_workflow.py | 4 ++ 6 files changed, 150 insertions(+), 4 deletions(-) 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