Allow using data collection steps via workflow API.

Implement API test for this and fixup test for previous commit related improved workflow run endpoint.
This commit is contained in:
John Chilton
2014-07-28 19:18:40 -04:00
parent 7ba215ad13
commit 5979ac41c8
6 changed files with 197 additions and 27 deletions
+29 -14
View File
@@ -82,11 +82,17 @@ class WorkflowsAPIController(BaseAPIController, UsesStoredWorkflowMixin, UsesHis
latest_workflow = stored_workflow.latest_workflow
inputs = {}
for step in latest_workflow.steps:
if step.type == 'data_input':
step_type = step.type
if step_type in ['data_input', 'data_collection_input']:
if step.tool_inputs and "name" in step.tool_inputs:
inputs[step.id] = {'label': step.tool_inputs['name'], 'value': ""}
label = step.tool_inputs['name']
elif step_type == "data_input":
label = "Input Dataset"
elif step_type == "data_collection_input":
label = "Input Dataset Collection"
else:
inputs[step.id] = {'label': "Input Dataset", 'value': ""}
raise ValueError("Invalid step_type %s" % step_type)
inputs[step.id] = {'label': label, 'value': ""}
else:
pass
# Eventually, allow regular tool parameters to be inserted and modified at runtime.
@@ -258,26 +264,35 @@ class WorkflowsAPIController(BaseAPIController, UsesStoredWorkflowMixin, UsesHis
try:
if inputs[k]['src'] == 'ldda':
ldda = trans.sa_session.query(self.app.model.LibraryDatasetDatasetAssociation).get(
trans.security.decode_id(inputs[k]['id']))
trans.security.decode_id(inputs[k]['id']))
assert trans.user_is_admin() or trans.app.security_agent.can_access_dataset( trans.get_current_user_roles(), ldda.dataset )
hda = ldda.to_history_dataset_association(history, add_to_history=add_to_history)
content = ldda.to_history_dataset_association(history, add_to_history=add_to_history)
elif inputs[k]['src'] == 'ld':
ldda = trans.sa_session.query(self.app.model.LibraryDataset).get(
trans.security.decode_id(inputs[k]['id'])).library_dataset_dataset_association
trans.security.decode_id(inputs[k]['id'])).library_dataset_dataset_association
assert trans.user_is_admin() or trans.app.security_agent.can_access_dataset( trans.get_current_user_roles(), ldda.dataset )
hda = ldda.to_history_dataset_association(history, add_to_history=add_to_history)
content = ldda.to_history_dataset_association(history, add_to_history=add_to_history)
elif inputs[k]['src'] == 'hda':
# Get dataset handle, add to dict and history if necessary
hda = trans.sa_session.query(self.app.model.HistoryDatasetAssociation).get(
trans.security.decode_id(inputs[k]['id']))
assert trans.user_is_admin() or trans.app.security_agent.can_access_dataset( trans.get_current_user_roles(), hda.dataset )
content = trans.sa_session.query(self.app.model.HistoryDatasetAssociation).get(
trans.security.decode_id(inputs[k]['id']))
assert trans.user_is_admin() or trans.app.security_agent.can_access_dataset( trans.get_current_user_roles(), content.dataset )
elif inputs[k]['src'] == 'hdca':
content = self.app.dataset_collections_service.get_dataset_collection_instance(
trans,
'history',
inputs[k]['id']
)
else:
trans.response.status = 400
return "Unknown dataset source '%s' specified." % inputs[k]['src']
if add_to_history and hda.history != history:
hda = hda.copy()
history.add_dataset(hda)
inputs[k]['hda'] = hda
if add_to_history and content.history != history:
content = content.copy()
if isinstance( content, self.app.model.HistoryDatasetAssociation ):
history.add_dataset( content )
else:
history.add_dataset_collection( content )
inputs[k]['hda'] = content # TODO: rename key to 'content', prescreen input ensure not populated explicitly
except AssertionError:
trans.response.status = 400
return "Invalid Dataset '%s' Specified" % inputs[k]['id']
+1 -1
View File
@@ -290,7 +290,7 @@ class WorkflowInvoker( object ):
step.state = step.module.get_runtime_state()
# This is an input step. Make sure we have an available input.
if step.type == 'data_input':
if step.type in [ 'data_input', 'data_collection_input' ]:
if self.inputs_by == "step_id":
key = str( step.id )
elif self.inputs_by == "name":
+5
View File
@@ -11,6 +11,8 @@ workflow_str = resource_string( __name__, "test_workflow_1.ga" )
# Simple workflow that takes an input and filters with random lines twice in a
# row - first grabbing 8 lines at random and then 6.
workflow_random_x2_str = resource_string( __name__, "test_workflow_2.ga" )
workflow_two_paired_str = resource_string( __name__, "test_workflow_two_paired.ga" )
DEFAULT_HISTORY_TIMEOUT = 10 # Secs to wait on history to turn ok
@@ -140,6 +142,9 @@ class WorkflowPopulator( object ):
def load_random_x2_workflow( self, name ):
return self.load_workflow( name, content=workflow_random_x2_str )
def load_two_paired_workflow( self, name ):
return self.load_workflow( name, content=workflow_two_paired_str )
def simple_workflow( self, name, **create_kwds ):
workflow = self.load_workflow( name )
return self.create_workflow( workflow, **create_kwds )
+116
View File
@@ -0,0 +1,116 @@
{
"a_galaxy_workflow": "true",
"annotation": "",
"format-version": "0.1",
"name": "MultipairTest223",
"steps": {
"0": {
"annotation": "",
"id": 0,
"input_connections": {},
"inputs": [
{
"description": "",
"name": "f1"
}
],
"name": "Input dataset collection",
"outputs": [],
"position": {
"left": 302.3333435058594,
"top": 330
},
"tool_errors": null,
"tool_id": null,
"tool_state": "{\"collection_type\": \"paired\", \"name\": \"f1\"}",
"tool_version": null,
"type": "data_collection_input",
"user_outputs": []
},
"1": {
"annotation": "",
"id": 1,
"input_connections": {},
"inputs": [
{
"description": "",
"name": "f2"
}
],
"name": "Input dataset collection",
"outputs": [],
"position": {
"left": 288.3333435058594,
"top": 446
},
"tool_errors": null,
"tool_id": null,
"tool_state": "{\"collection_type\": \"paired\", \"name\": \"f2\"}",
"tool_version": null,
"type": "data_collection_input",
"user_outputs": []
},
"2": {
"annotation": "",
"id": 2,
"input_connections": {
"kind|f1": {
"id": 0,
"output_name": "output"
},
"kind|f2": {
"id": 1,
"output_name": "output"
}
},
"inputs": [],
"name": "collection_two_paired",
"outputs": [
{
"name": "out1",
"type": "txt"
}
],
"position": {
"left": 782.3333740234375,
"top": 200
},
"post_job_actions": {},
"tool_errors": null,
"tool_id": "collection_two_paired",
"tool_state": "{\"__page__\": 0, \"kind\": \"{\\\"f1\\\": null, \\\"f2\\\": null, \\\"collection_type\\\": \\\"paired\\\", \\\"__current_case__\\\": 0}\", \"__rerun_remap_job_id__\": null}",
"tool_version": "0.1.0",
"type": "tool",
"user_outputs": []
},
"3": {
"annotation": "",
"id": 3,
"input_connections": {
"cond1|input1": {
"id": 2,
"output_name": "out1"
}
},
"inputs": [],
"name": "Concatenate datasets",
"outputs": [
{
"name": "out_file1",
"type": "input"
}
],
"position": {
"left": 1239.3333740234375,
"top": 108.97916793823242
},
"post_job_actions": {},
"tool_errors": null,
"tool_id": "cat2",
"tool_state": "{\"__page__\": 0, \"__rerun_remap_job_id__\": null, \"cond1\": \"{\\\"datatype\\\": \\\"txt\\\", \\\"input1\\\": null, \\\"__current_case__\\\": 0}\"}",
"tool_version": "1.0.0",
"type": "tool",
"user_outputs": []
}
}
}
+42 -8
View File
@@ -85,6 +85,25 @@ class WorkflowsApiTestCase( api.ApiTestCase ):
self._assert_status_code_is( run_workflow_response, 200 )
self.dataset_populator.wait_for_history( history_id, assert_ok=True )
@skip_without_tool( "cat1" )
@skip_without_tool( "collection_two_paired" )
def test_run_workflow_collection_params( self ):
workflow = self.workflow_populator.load_two_paired_workflow( name="test_for_run_two_paired" )
workflow_id = self.workflow_populator.create_workflow( workflow )
history_id = self.dataset_populator.new_history()
hdca1 = self.dataset_collection_populator.create_pair_in_history( history_id, contents=["1 2 3", "4 5 6"] ).json()
hdca2 = self.dataset_collection_populator.create_pair_in_history( history_id, contents=["7 8 9", "0 a b"] ).json()
self.dataset_populator.wait_for_history( history_id, assert_ok=True )
label_map = { "f1": self._ds_entry( hdca1 ), "f2": self._ds_entry( hdca2 ) }
workflow_request = dict(
history="hist_id=%s" % history_id,
workflow_id=workflow_id,
ds_map=self._build_ds_map( workflow_id, label_map ),
)
run_workflow_response = self._post( "workflows", data=workflow_request )
self._assert_status_code_is( run_workflow_response, 200 )
self.dataset_populator.wait_for_history( history_id, assert_ok=True )
@skip_without_tool( "cat1" )
def test_extract_from_history( self ):
history_id = self.dataset_populator.new_history()
@@ -137,7 +156,6 @@ class WorkflowsApiTestCase( api.ApiTestCase ):
input1 = tool_step[ "input_connections" ][ "input1" ]
input2 = tool_step[ "input_connections" ][ "queries_0|input2" ]
print downloaded_workflow
self.assertEquals( input_steps[ 0 ][ "id" ], input1[ "id" ] )
self.assertEquals( input_steps[ 1 ][ "id" ], input2[ "id" ] )
@@ -421,21 +439,34 @@ class WorkflowsApiTestCase( api.ApiTestCase ):
# renamed to 'the_new_name'.
assert "the_new_name" in map( lambda hda: hda[ "name" ], contents )
def _setup_workflow_run( self, workflow, history_id=None ):
def _setup_workflow_run( self, workflow, inputs_by='step_id', history_id=None ):
uploaded_workflow_id = self.workflow_populator.create_workflow( workflow )
if not history_id:
history_id = self.dataset_populator.new_history()
hda1 = self.dataset_populator.new_dataset( history_id, content="1 2 3" )
hda2 = self.dataset_populator.new_dataset( history_id, content="4 5 6" )
workflow_request = dict(
history="hist_id=%s" % history_id,
workflow_id=uploaded_workflow_id,
)
label_map = {
'WorkflowInput1': self._ds_entry(hda1),
'WorkflowInput2': self._ds_entry(hda2)
}
workflow_request = dict(
history="hist_id=%s" % history_id,
workflow_id=uploaded_workflow_id,
ds_map=self._build_ds_map( uploaded_workflow_id, label_map ),
)
if inputs_by == 'step_id':
ds_map = self._build_ds_map( uploaded_workflow_id, label_map )
workflow_request[ "ds_map" ] = ds_map
elif inputs_by == "step_index":
index_map = {
'0': self._ds_entry(hda1),
'1': self._ds_entry(hda2)
}
workflow_request[ "inputs" ] = dumps( index_map )
workflow_request[ "inputs_by" ] = 'step_index'
elif inputs_by == "name":
workflow_request[ "inputs" ] = dumps( label_map )
workflow_request[ "inputs_by" ] = 'name'
return workflow_request, history_id
def _build_ds_map( self, workflow_id, label_map ):
@@ -471,7 +502,10 @@ class WorkflowsApiTestCase( api.ApiTestCase ):
return workflow_inputs
def _ds_entry( self, hda ):
return dict( src="hda", id=hda[ "id" ] )
src = 'hda'
if 'history_content_type' in hda and hda[ 'history_content_type' ] == "dataset_collection":
src = 'hdca'
return dict( src=src, id=hda[ "id" ] )
def _assert_user_has_workflow_with_name( self, name ):
names = self.__workflow_names()
@@ -16,12 +16,12 @@
<option value="list">List of Datasets</option>
</param>
<when value="paired">
<param name="f1" type="data_collection" collection_type="paired" />
<param name="f2" type="data_collection" collection_type="paired" />
<param name="f1" type="data_collection" collection_type="paired" label="F1" />
<param name="f2" type="data_collection" collection_type="paired" label="F2" />
</when>
<when value="list">
<param name="f1" type="data_collection" collection_type="list" />
<param name="f2" type="data_collection" collection_type="list" />
<param name="f1" type="data_collection" collection_type="list" label="F1" />
<param name="f2" type="data_collection" collection_type="list" label="F2" />
</when>
</conditional>
</inputs>