diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index efce00417df..f6528cd66fd 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -3109,11 +3109,17 @@ class DatasetCollection( object, Dictifiable, UsesAnnotations ): @property def populated( self ): - return self.populated_state == DatasetCollection.populated_states.OK + top_level_populated = self.populated_state == DatasetCollection.populated_states.OK + if top_level_populated and self.has_subcollections: + return all(map(lambda e: e.child_collection.populated, self.elements)) + return top_level_populated @property def waiting_for_elements( self ): - return self.populated_state == DatasetCollection.populated_states.NEW + top_level_waiting = self.populated_state == DatasetCollection.populated_states.NEW + if not top_level_waiting and self.has_subcollections: + return any(map(lambda e: e.child_collection.waiting_for_elements, self.elements)) + return top_level_waiting def mark_as_populated( self ): self.populated_state = DatasetCollection.populated_states.OK @@ -3178,6 +3184,10 @@ class DatasetCollection( object, Dictifiable, UsesAnnotations ): # Nothing currently editable in this class. return {} + @property + def has_subcollections(self): + return ":" in self.collection_type + class DatasetCollectionInstance( object, HasName ): """ diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index 2e28e6aa014..20ba6dea2b9 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -83,6 +83,11 @@ def execute( trans, tool, param_combinations, history, rerun_remap_job_id=None, log.debug("Executed %d job(s) for tool %s request: %s" % (job_count, tool.id, all_jobs_timer)) if collection_info: history = history or tool.get_default_history_by_trans( trans ) + if len(param_combinations) == 0: + template = "Attempting to map over an empty collection, this is not yet implemented. colleciton_info is [%s]" + message = template % collection_info + log.warn(message) + raise Exception(message) params = param_combinations[0] execution_tracker.create_output_collections( trans, history, params ) diff --git a/test/api/test_workflows.py b/test/api/test_workflows.py index 03f939cc9d0..6ae340f96e1 100644 --- a/test/api/test_workflows.py +++ b/test/api/test_workflows.py @@ -763,6 +763,50 @@ steps: content = self.dataset_populator.get_history_dataset_content( history_id, hid=7 ) self.assertEqual(content.strip(), "samp1\t10.0\nsamp2\t20.0") + @skip_without_tool( "collection_split_on_column" ) + def test_workflow_run_dynamic_output_collections_3(self): + # Test a workflow that create a list:list:list followed by a mapping step. + history_id = self.dataset_populator.new_history() + workflow_id = self._upload_yaml_workflow(""" +class: GalaxyWorkflow +steps: + - label: text_input1 + type: input + - label: text_input2 + type: input + - label: cat_inputs + tool_id: cat1 + state: + input1: + $link: text_input1 + queries: + - input2: + $link: text_input2 + - label: split_up_1 + tool_id: collection_split_on_column + state: + input1: + $link: cat_inputs#out_file1 + - label: split_up_2 + tool_id: collection_split_on_column + state: + input1: + $link: split_up_1#split_output + - tool_id: cat + state: + input1: + $link: split_up_2#split_output +""") + hda1 = self.dataset_populator.new_dataset( history_id, content="samp1\t10.0\nsamp2\t20.0\n" ) + hda2 = self.dataset_populator.new_dataset( history_id, content="samp1\t30.0\nsamp2\t40.0\n" ) + self.dataset_populator.wait_for_history( history_id, assert_ok=True ) + inputs = { + '0': self._ds_entry(hda1), + '1': self._ds_entry(hda2), + } + invocation_id = self.__invoke_workflow( history_id, workflow_id, inputs ) + self.wait_for_invocation_and_jobs( history_id, workflow_id, invocation_id ) + @skip_without_tool( "mapper" ) @skip_without_tool( "pileup" ) def test_workflow_metadata_validation_0( self ): @@ -876,7 +920,11 @@ steps: index_map = { '0': self._ds_entry(hda1), } - invocation_id = self.__invoke_workflow( history_id, uploaded_workflow_id, index_map ) + invocation_id = self.__invoke_workflow( + history_id, + uploaded_workflow_id, + index_map, + ) # Give some time for workflow to get scheduled before scanning the history. time.sleep( 5 ) self.dataset_populator.wait_for_history( history_id, assert_ok=True ) @@ -888,10 +936,17 @@ steps: self.__review_paused_steps( uploaded_workflow_id, invocation_id, order_index=2, action=True ) - time.sleep( 5 ) + invocation_scheduled = False + for i in range( 25 ): + invocation = self._invocation_details( uploaded_workflow_id, invocation_id ) + if invocation[ 'state' ] == 'scheduled': + invocation_scheduled = True + break + + time.sleep( .5 ) + + assert invocation_scheduled, "Workflow state is not scheduled..." self.dataset_populator.wait_for_history( history_id, assert_ok=True ) - invocation = self._invocation_details( uploaded_workflow_id, invocation_id ) - assert invocation[ 'state' ] == 'scheduled', invocation @skip_without_tool( "cat" ) def test_workflow_pause_cancel( self ):