Allow tools to output collections with a dynamic number of datasets.

Models:

Track whether dataset collections have been populated yet.

Dataset collections are still effectively immutable once populated - but dynamic output collections require them to be sort of like `final` fields in Java (analogy courtesy of JJ) - allowing them to be declared before they are initialized or populated. This is tracked by the `populated_state` field.

Tools:

Output collections can now describe `discover_datasets` elements just like datasets - except in this case instead of dynamically populating new datasets in the history - they will comprise the collection. `designation` has been reused to serve as the element_identifier for the collection element corresponding to the dataset.

See Pull Request 356 for more information on the discover_datasets tag https://bitbucket.org/galaxy/galaxy-central/pull-request/356/enhancements-for-runtime-discovered.

Workflows:

Update workflow execution and recovery for dynamic output collections.

Galaxy workflow data flow before collections

* - * - * - * - * - *

Galaxy worfklow data flow after collections (iteration 1)

* - * - * \
           * - * - *
* - * - * /         \
                     * - * - *
* - * - * \         /
           * - * - *
* - * - * /

Galaxy worfklow data flow after this commit

              / * - * \
         * - *         * - *
        /     \ * - * /     \
       /                     \
      /                       \
     /        / * - * \        \
* - * -- * - *         * - * -- * - *
     \        \ * - * /        /
      \                       /
       \                     /
        \     / * - * \     /
         * - *         * - *
              \ * - * /
This commit is contained in:
John Chilton
2015-01-15 09:30:00 -05:00
parent 44f7317fa5
commit 4c5c8a47db
17 changed files with 417 additions and 7 deletions
+11
View File
@@ -1202,6 +1202,11 @@ class JobWrapper( object ):
out_data = dict( [ ( da.name, da.dataset ) for da in job.output_datasets ] )
inp_data.update( [ ( da.name, da.dataset ) for da in job.input_library_datasets ] )
out_data.update( [ ( da.name, da.dataset ) for da in job.output_library_datasets ] )
# TODO: eliminate overlap with tools/evaluation.py
out_collections = dict( [ ( obj.name, obj.dataset_collection_instance ) for obj in job.output_dataset_collection_instances ] )
out_collections.update( [ ( obj.name, obj.dataset_collection ) for obj in job.output_dataset_collections ] )
input_ext = 'data'
for _, data in inp_data.items():
# For loop odd, but sort simulating behavior in galaxy.tools.actions
@@ -1218,6 +1223,12 @@ class JobWrapper( object ):
'children': self.tool.collect_child_datasets(out_data, self.working_directory),
'primary': self.tool.collect_primary_datasets(out_data, self.working_directory, input_ext)
}
self.tool.collect_dynamic_collections(
out_collections,
job_working_directory=self.working_directory,
inp_data=inp_data,
job=job,
)
param_dict.update({'__collected_datasets__': collected_datasets})
# Certain tools require tasks to be completed after job execution
# ( this used to be performed in the "exec_after_process" hook, but hooks are deprecated ).
+18 -2
View File
@@ -28,6 +28,7 @@ class DatasetCollectionManager( object ):
Abstraction for interfacing with dataset collections instance - ideally abstarcts
out model and plugin details.
"""
ELEMENTS_UNINITIALIZED = object()
def __init__( self, app ):
self.type_registry = DatasetCollectionTypesRegistry( app )
@@ -129,11 +130,26 @@ class DatasetCollectionManager( object ):
elements = self.__load_elements( trans, element_identifiers )
# else if elements is set, it better be an ordered dict!
type_plugin = collection_type_description.rank_type_plugin()
dataset_collection = builder.build_collection( type_plugin, elements )
if elements is not self.ELEMENTS_UNINITIALIZED:
type_plugin = collection_type_description.rank_type_plugin()
dataset_collection = builder.build_collection( type_plugin, elements )
else:
dataset_collection = model.DatasetCollection( populated=False )
dataset_collection.collection_type = collection_type
return dataset_collection
def set_collection_elements( self, dataset_collection, dataset_instances ):
if dataset_collection.populated:
raise Exception("Cannot reset elements of an already populated dataset collection.")
collection_type = dataset_collection.collection_type
collection_type_description = self.collection_type_descriptions.for_collection_type( collection_type )
type_plugin = collection_type_description.rank_type_plugin()
builder.set_collection_elements( dataset_collection, type_plugin, dataset_instances )
dataset_collection.mark_as_populated()
return dataset_collection
def delete( self, trans, instance_type, id ):
dataset_collection_instance = self.get_dataset_collection_instance( trans, instance_type, id, check_ownership=True )
dataset_collection_instance.deleted = True
+6 -2
View File
@@ -77,7 +77,9 @@ def dictify_dataset_collection_instance( dataset_colleciton_instance, parent, se
# TODO: Work in progress - this end-point is not right yet...
dict_value[ 'url' ] = web.url_for( 'library_content', library_id=encoded_library_id, id=encoded_id, folder_id=encoded_folder_id )
if view == "element":
dict_value[ 'elements' ] = map( dictify_element, dataset_colleciton_instance.collection.elements )
collection = dataset_colleciton_instance.collection
dict_value[ 'elements' ] = map( dictify_element, collection.elements )
dict_value[ 'populated' ] = collection.populated
security.encode_all_ids( dict_value, recursive=True ) # TODO: Use Kyle's recusrive formulation of this.
return dict_value
@@ -87,7 +89,9 @@ def dictify_element( element ):
object_detials = element.element_object.to_dict()
if element.child_collection:
# Recursively yield elements for each nested collection...
object_detials[ "elements" ] = map( dictify_element, element.child_collection.elements )
child_collection = element.child_collection
object_detials[ "elements" ] = map( dictify_element, child_collection.elements )
object_detials[ "populated" ] = child_collection.populated
dictified[ "object" ] = object_detials
return dictified
+23
View File
@@ -2676,14 +2676,37 @@ class DatasetCollection( object, Dictifiable, UsesAnnotations ):
"""
dict_collection_visible_keys = ( 'id', 'collection_type' )
dict_element_visible_keys = ( 'id', 'collection_type' )
populated_states = Bunch(
NEW='new', # New dataset collection, unpopulated elements
OK='ok', # Collection elements populated (HDAs may or may not have errors)
FAILED='failed', # some problem populating state, won't be populated
)
def __init__(
self,
id=None,
collection_type=None,
populated=True,
):
self.id = id
self.collection_type = collection_type
if not populated:
self.populated_state = DatasetCollection.populated_states.NEW
@property
def populated( self ):
return self.populated_state == DatasetCollection.populated_states.OK
@property
def waiting_for_elements( self ):
return self.populated_state == DatasetCollection.populated_states.NEW
def mark_as_populated( self ):
self.populated_state = DatasetCollection.populated_states.OK
def handle_population_failed( self, message ):
self.populated_state = DatasetCollection.populated_states.FAILED
self.populated_state_message = message
@property
def dataset_instances( self ):
+2
View File
@@ -621,6 +621,8 @@ model.TransferJob.table = Table( "transfer_job", metadata,
model.DatasetCollection.table = Table( "dataset_collection", metadata,
Column( "id", Integer, primary_key=True ),
Column( "collection_type", Unicode(255), nullable=False ),
Column( "populated_state", TrimmedString( 64 ), default='ok', nullable=False ),
Column( "populated_state_message", TEXT ),
Column( "create_time", DateTime, default=now ),
Column( "update_time", DateTime, default=now, onupdate=now ),
)
@@ -38,6 +38,18 @@ def upgrade(migrate_engine):
for table in TABLES:
__create(table)
try:
dataset_collection_table = Table( "dataset_collection", metadata, autoload=True )
# need server_default because column in non-null
populated_state_column = Column( 'populated_state', TrimmedString( 64 ), default='ok', server_default="ok", nullable=False )
populated_state_column.create( dataset_collection_table )
populated_message_column = Column( 'populated_state_message', TEXT, nullable=True )
populated_message_column.create( dataset_collection_table )
except Exception as e:
print str(e)
log.exception( "Creating dataset collection populated column failed." )
def downgrade(migrate_engine):
metadata.bind = migrate_engine
@@ -46,6 +58,16 @@ def downgrade(migrate_engine):
for table in TABLES:
__drop(table)
try:
dataset_collection_table = Table( "dataset_collection", metadata, autoload=True )
populated_state_column = dataset_collection_table.c.populated_state
populated_state_column.drop()
populated_message_column = dataset_collection_table.c.populated_state_message
populated_message_column.drop()
except Exception as e:
print str(e)
log.exception( "Dropping dataset collection populated_state/ column failed." )
def __create(table):
try:
+32 -2
View File
@@ -297,7 +297,13 @@ class ToolOutputCollection( ToolOutputBase ):
self.structure = structure
self.outputs = odict()
# TODO:
self.metadata_source = None
def known_outputs( self, inputs ):
if self.dynamic_structure:
return []
def to_part( ( element_identifier, output ) ):
return ToolOutputCollectionPart( self, element_identifier, output )
@@ -316,14 +322,33 @@ class ToolOutputCollection( ToolOutputBase ):
return map( to_part, outputs.items() )
@property
def dynamic_structure(self):
return self.structure.dynamic
@property
def dataset_collectors(self):
if not self.dynamic_structure:
raise Exception("dataset_collectors called for output collection with static structure")
return self.structure.dataset_collectors
class ToolOutputCollectionStructure( object ):
def __init__( self, collection_type=None, structured_like=None ):
def __init__(
self,
collection_type,
structured_like,
dataset_collectors,
):
self.collection_type = collection_type
self.structured_like = structured_like
if collection_type is None and structured_like is None:
self.dataset_collectors = dataset_collectors
if collection_type is None and structured_like is None and dataset_collectors is None:
raise ValueError( "Output collection types must be specify type of structured_like" )
if dataset_collectors and structured_like:
raise ValueError( "Cannot specify dynamic structure (discovered_datasets) and structured_like attribute." )
self.dynamic = dataset_collectors is not None
class ToolOutputCollectionPart( object ):
@@ -2146,6 +2171,11 @@ class Tool( object, Dictifiable ):
"""
return output_collect.collect_primary_datasets( self, output, job_working_directory, input_ext )
def collect_dynamic_collections( self, output, **kwds ):
""" Find files corresponding to dynamically structured collections.
"""
return output_collect.collect_dynamic_collections( self, output, **kwds )
def to_dict( self, trans, link_details=False, io_details=False ):
""" Returns dict of tool. """
+4
View File
@@ -289,6 +289,10 @@ class DefaultToolAction( object ):
elements[ output_part_def.element_identifier ] = element
if output.dynamic_structure:
assert not elements # known_outputs must have been empty
elements = collections_manager.ELEMENTS_UNINITIALIZED
if mapping_over_collection:
dc = collections_manager.create_dataset_collection(
trans,
@@ -162,6 +162,11 @@ class DatasetCollectionMatcher( object ):
return self.dataset_collection_match( dataset_collection )
def dataset_collection_match( self, dataset_collection ):
# If dataset collection not yet populated, cannot determine if it
# would be a valid match for this parameter.
if not dataset_collection.populated:
return False
valid = True
for element in dataset_collection.elements:
if not self.__valid_element( element ):
@@ -12,6 +12,176 @@ from galaxy.util import odict
DATASET_ID_TOKEN = "DATASET_ID"
DEFAULT_EXTRA_FILENAME_PATTERN = r"primary_DATASET_ID_(?P<designation>[^_]+)_(?P<visible>[^_]+)_(?P<ext>[^_]+)(_(?P<dbkey>[^_]+))?"
import logging
log = logging.getLogger( __name__ )
def collect_dynamic_collections(
tool,
output_collections,
job_working_directory,
inp_data={},
job=None,
):
collections_service = tool.app.dataset_collections_service
job_context = JobContext(
tool,
job,
job_working_directory,
inp_data,
)
for name, has_collection in output_collections.items():
if name not in tool.output_collections:
continue
output_collection_def = tool.output_collections[ name ]
if not output_collection_def.dynamic_structure:
continue
# Could be HDCA for normal jobs or a DC for mapping
# jobs.
if hasattr(has_collection, "collection"):
collection = has_collection.collection
else:
collection = has_collection
try:
elements = job_context.build_collection_elements(
collection,
output_collection_def,
)
collections_service.set_collection_elements(
collection,
elements
)
except Exception:
log.info("Problem gathering output collection.")
collection.handle_population_failed("Problem building datasets for collection.")
class JobContext( object ):
def __init__( self, tool, job, job_working_directory, inp_data ):
self.inp_data = inp_data
self.app = tool.app
self.sa_session = tool.sa_session
self.job = job
self.job_working_directory = job_working_directory
@property
def permissions( self ):
inp_data = self.inp_data
existing_datasets = [ inp for inp in inp_data.values() if inp ]
if existing_datasets:
permissions = self.app.security_agent.guess_derived_permissions_for_datasets( existing_datasets )
else:
# No valid inputs, we will use history defaults
permissions = self.app.security_agent.history_get_default_permissions( self.job.history )
return permissions
def find_files( self, collection, dataset_collectors ):
filenames = odict.odict()
for path, extra_file_collector in walk_over_extra_files( dataset_collectors, self.job_working_directory, collection ):
filenames[ path ] = extra_file_collector
return filenames
def build_collection_elements( self, collection, output_collection_def ):
datasets = self.create_datasets(
collection,
output_collection_def,
)
elements = odict.odict()
# TODO: allow configurable sorting.
# <sort by="lexical" /> <!-- default -->
# <sort by="reverse_lexical" />
# <sort regex="example.(\d+).fastq" by="1:numerical" />
# <sort regex="part_(\d+)_sample_([^_]+).fastq" by="2:lexical,1:numerical" />
# TODO: allow nested structure
for designation in datasets.keys():
elements[ designation ] = datasets[ designation ]
return elements
def create_datasets( self, collection, output_collection_def ):
dataset_collectors = output_collection_def.dataset_collectors
filenames = self.find_files( collection, dataset_collectors )
datasets = {}
for filename, extra_file_collector in filenames.iteritems():
fields_match = extra_file_collector.match( collection, os.path.basename( filename ) )
if not fields_match:
raise Exception( "Problem parsing metadata fields for file %s" % filename )
designation = fields_match.designation
visible = fields_match.visible
ext = fields_match.ext
dbkey = fields_match.dbkey
# Create new primary dataset
name = fields_match.name or designation
dataset = self.create_dataset(
ext=ext,
designation=designation,
visible=visible,
dbkey=dbkey,
name=name,
filename=filename,
metadata_source_name=output_collection_def.metadata_source,
)
datasets[ designation ] = dataset
return datasets
def create_dataset(
self,
ext,
designation,
visible,
dbkey,
name,
filename,
metadata_source_name,
):
app = self.app
sa_session = self.sa_session
# Copy metadata from one of the inputs if requested.
metadata_source = None
if metadata_source_name:
metadata_source = self.inp_data[ metadata_source_name ]
# Create new primary dataset
primary_data = app.model.HistoryDatasetAssociation( extension=ext,
designation=designation,
visible=visible,
dbkey=dbkey,
create_dataset=True,
sa_session=sa_session )
app.security_agent.set_all_dataset_permissions( primary_data.dataset, self.permissions )
sa_session.add( primary_data )
sa_session.flush()
# Move data from temp location to dataset location
app.object_store.update_from_file(primary_data.dataset, file_name=filename, create=True)
primary_data.set_size()
# If match specified a name use otherwise generate one from
# designation.
primary_data.name = name
if metadata_source:
primary_data.init_meta( copy_from=metadata_source )
else:
primary_data.init_meta()
# Associate new dataset with job
if self.job:
assoc = app.model.JobToOutputDatasetAssociation( '__new_primary_file_%s|%s__' % ( name, designation ), primary_data )
assoc.job = self.job
sa_session.add( assoc )
sa_session.flush()
primary_data.state = 'ok'
return primary_data
def collect_primary_datasets( tool, output, job_working_directory, input_ext ):
app = tool.app
+4
View File
@@ -153,9 +153,13 @@ class XmlToolSource(ToolSource):
default_format = collection_elem.get( "format", "data" )
collection_type = collection_elem.get( "type", None )
structured_like = collection_elem.get( "structured_like", None )
dataset_collectors = None
if collection_elem.find( "discover_datasets" ) is not None:
dataset_collectors = output_collect.dataset_collectors_from_elem( collection_elem )
structure = galaxy.tools.ToolOutputCollectionStructure(
collection_type=collection_type,
structured_like=structured_like,
dataset_collectors=dataset_collectors,
)
output_collection = galaxy.tools.ToolOutputCollection(
name,
+5
View File
@@ -901,6 +901,11 @@ class ToolModule( WorkflowModule ):
if replacement_value.hidden_beneath_collection_instance:
replacement_value = replacement_value.hidden_beneath_collection_instance
outputs[ replacement_name ] = replacement_value
for job_output_collection in job_0.output_dataset_collection_instances:
replacement_name = job_output_collection.name
replacement_value = job_output_collection.dataset_collection_instance
outputs[ replacement_name ] = replacement_value
progress.set_step_outputs( step, outputs )
+12 -1
View File
@@ -251,7 +251,18 @@ class WorkflowProgress( object ):
step_outputs = self.outputs[ connection.output_step.id ]
if step_outputs is STEP_OUTPUT_DELAYED:
raise modules.DelayedWorkflowEvaluation()
return step_outputs[ connection.output_name ]
replacement = step_outputs[ connection.output_name ]
if isinstance( replacement, model.HistoryDatasetCollectionAssociation ):
if not replacement.collection.populated:
if not replacement.collection.waiting_for_elements:
# If we are not waiting for elements, there was some
# problem creating the collection. Collection will never
# be populated.
# TODO: consider distinguish between cancelled and failed?
raise modules.CancelWorkflowEvaluation()
raise modules.DelayedWorkflowEvaluation()
return replacement
def set_outputs_for_input( self, step, outputs={} ):
if self.inputs_by_step_id:
+34
View File
@@ -276,6 +276,40 @@ class ToolsTestCase( api.ApiTestCase ):
contents1 = self.dataset_populator.get_history_dataset_content( history_id, dataset_id=element1["object"]["id"])
assert contents1 == "1\n", contents1
@skip_without_tool( "collection_split_on_column" )
def test_dynamic_list_output( self ):
history_id = self.dataset_populator.new_history()
new_dataset1 = self.dataset_populator.new_dataset( history_id, content='samp1\t1\nsamp1\t3\nsamp2\t2\nsamp2\t4\n' )
inputs = {
'input1': dataset_to_param( new_dataset1 ),
}
self.dataset_populator.wait_for_history( history_id, assert_ok=True )
create = self._run( "collection_split_on_column", history_id, inputs, assert_ok=True )
jobs = create[ 'jobs' ]
implicit_collections = create[ 'implicit_collections' ]
collections = create[ 'output_collections' ]
self.assertEquals( len( jobs ), 1 )
job_id = jobs[ 0 ][ "id" ]
self.assertEquals( len( implicit_collections ), 0 )
self.assertEquals( len( collections ), 1 )
output_collection = collections[0]
self._assert_has_keys( output_collection, "id", "name", "elements", "populated" )
assert not output_collection[ "populated" ]
assert len( output_collection[ "elements" ] ) == 0
self.dataset_populator.wait_for_job( job_id, assert_ok=True )
get_collection_response = self._get( "dataset_collections/%s" % output_collection[ "id" ], data={"instance_type": "history"} )
self._assert_status_code_is( get_collection_response, 200 )
output_collection = get_collection_response.json()
self._assert_has_keys( output_collection, "id", "name", "elements", "populated" )
assert output_collection[ "populated" ]
assert len( output_collection[ "elements" ] ) == 2
@skip_without_tool( "cat1" )
def test_run_cat1_with_two_inputs( self ):
# Run tool with an multiple data parameter and grouping (repeat)
+38
View File
@@ -538,6 +538,44 @@ class WorkflowsApiTestCase( BaseWorkflowsApiTestCase ):
self.dataset_populator.wait_for_history( history_id, assert_ok=True )
self.assertEquals("a\nc\nb\nd\ne\ng\nf\nh\n", self.dataset_populator.get_history_dataset_content( history_id, hid=0 ) )
def test_workflow_run_dynamic_output_collections(self):
history_id = self.dataset_populator.new_history()
workflow_id = self._upload_yaml_workflow("""
- 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
tool_id: collection_split_on_column
state:
input1:
$link: cat_inputs#out_file1
- tool_id: cat_list
state:
input1:
$link: split_up#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),
}
self.__invoke_workflow( history_id, workflow_id, inputs )
# TODO: wait on workflow invocations
time.sleep(10)
self.dataset_populator.wait_for_history( history_id, assert_ok=True )
self.assertEquals("10.0\n30.0\n20.0\n40.0\n", self.dataset_populator.get_history_dataset_content( history_id, hid=0 ) )
def test_workflow_request( self ):
workflow = self.workflow_populator.load_workflow( name="test_for_queue" )
workflow_request, history_id = self._setup_workflow_run( workflow )
@@ -0,0 +1,30 @@
<tool id="collection_split_on_column" name="collection_split_on_column" version="0.1.0">
<command>
mkdir outputs; cd outputs; awk '{ print \$2 > \$1 ".tabular" }' $input1
</command>
<inputs>
<param name="input1" type="data" label="Input Table" help="Table to split on first column" format="tabular" />
</inputs>
<outputs>
<collection name="split_output" type="list" label="Table split on first column">
<discover_datasets pattern="__name_and_ext__" directory="outputs" />
</collection>
</outputs>
<tests>
<test>
<param name="input1" value="tinywga.fam" />
<output_collection name="split_output" type="list">
<element name="101">
<assert_contents>
<has_text_matching expression="^1\n2\n3\n$" />
</assert_contents>
</element>
<element name="1334">
<assert_contents>
<has_text_matching expression="^1\n10\n11\n12\n13\n2\n$" />
</assert_contents>
</element>
</output_collection>
</test>
</tests>
</tool>
@@ -37,6 +37,7 @@
<tool file="collection_creates_pair.xml" />
<tool file="collection_creates_list.xml" />
<tool file="collection_optional_param.xml" />
<tool file="collection_split_on_column.xml" />
<tool file="multiple_versions_v01.xml" />
<tool file="multiple_versions_v02.xml" />