Add implicit datatype conversions.

If a dataset is not the proper format required by a tool, but there is a converter available, the dataset will appear as a valid option in DataToolParameters.

The converted dataset will be created (and subsequently reused) when the job is executed. Conversion utilizes the job runner, and will run on the cluster.

Setting metadata will invalidate the converted dataset, and a new one will be automatically generated as needed.
This commit is contained in:
Daniel Blankenberg
2008-05-15 19:09:17 +00:00
parent c424782fb4
commit a7ea9357f2
9 changed files with 135 additions and 22 deletions
+11 -5
View File
@@ -183,11 +183,11 @@ class Data( object ):
"""Returns available converters by type for this dataset"""
return datatypes_registry.get_converters_by_datatype(original_dataset.ext)
def convert_dataset(self, trans, original_dataset, target_type):
def convert_dataset(self, trans, original_dataset, target_type, return_output = False, visible = True ):
"""This function adds a job to the queue to convert a dataset to another type. Returns a message about success/failure."""
converter = trans.app.datatypes_registry.get_converter_by_target_type(original_dataset.ext, target_type)
converter = trans.app.datatypes_registry.get_converter_by_target_type( original_dataset.ext, target_type )
if converter is None:
return "A converter does not exist for %s to %s." % (original_dataset.ext, target_type)
raise "A converter does not exist for %s to %s." % ( original_dataset.ext, target_type )
#Generate parameter dictionary
params = {}
@@ -200,11 +200,17 @@ class Data( object ):
params[input_name] = original_dataset
#Run converter, job is dispatched through Queue
converted_dataset = converter.execute(trans, incoming=params)
converted_dataset = converter.execute( trans, incoming = params, set_output_hid = visible )
if len(params) > 0:
trans.log_event( "Converter params: %s" % (str(params)), tool_id=converter.id )
if not visible:
for name, value in converted_dataset.iteritems():
value.visible = False
if return_output:
return converted_dataset
return "The file conversion of %s on data %s has been added to the Queue." % (converter.name, original_dataset.hid)
def before_edit( self, dataset ):
@@ -213,7 +219,7 @@ class Data( object ):
def after_edit( self, dataset ):
"""This function is called on the dataset after metadata is edited."""
pass
dataset.clear_associated_files( metadata_safe = True )
class Text( Data ):
+56 -4
View File
@@ -124,16 +124,16 @@ class History( object ):
def add_galaxy_session( self, galaxy_session ):
self.galaxy_sessions.append( GalaxySessionToHistoryAssociation( galaxy_session, self ) )
def add_dataset( self, dataset, parent_id=None, genome_build=None ):
def add_dataset( self, dataset, parent_id=None, genome_build=None, set_hid = True ):
if parent_id:
for data in self.datasets:
if data.id == parent_id:
dataset.hid = data.hid
break
else:
dataset.hid = self._next_hid()
if set_hid: dataset.hid = self._next_hid()
else:
dataset.hid = self._next_hid()
if set_hid: dataset.hid = self._next_hid()
self.genome_build = genome_build
self.datasets.append( dataset )
@@ -289,6 +289,7 @@ class Dataset( object ):
dbkey = property( get_dbkey, set_dbkey )
def change_datatype( self, new_ext ):
self.clear_associated_files()
datatypes_registry.change_datatype( self, new_ext )
def get_size( self ):
"""Returns the size of the data on disk"""
@@ -325,6 +326,7 @@ class Dataset( object ):
def init_meta( self, copy_from=None ):
return self.datatype.init_meta( self, copy_from=copy_from )
def set_meta( self, **kwd ):
self.clear_associated_files( metadata_safe = True )
return self.datatype.set_meta( self, **kwd )
def set_readonly_meta( self, **kwd ):
return self.datatype.set_readonly_meta( self, **kwd )
@@ -338,6 +340,17 @@ class Dataset( object ):
return self.datatype.display_name( self )
def display_info( self ):
return self.datatype.display_info( self )
def get_associated_files_by_type( self, file_type ):
valid = []
for assoc in self.associated_files:
if not assoc.deleted and assoc.type == file_type:
valid.append( assoc )
return valid
def clear_associated_files( self, metadata_safe = False, purge = False ):
#metadata_safe = True means to only clear when assoc.metadata_safe == False
for assoc in self.associated_files:
if not metadata_safe or not assoc.metadata_safe:
assoc.clear( purge = purge )
def get_child_by_designation(self, designation):
# if self.history:
# for data in self.history.datasets:
@@ -418,7 +431,46 @@ class DatasetChildAssociation( object ):
self.designation = designation
self.parent = None
self.child = None
class DatasetAssociatedFile( object ):
def __init__( self, id = None, dataset_id = None, file_type = None, parent_id = None, filename = None, deleted = False, purged = False, metadata_safe = True ):
self.id = id
self.dataset_id = dataset_id
self.type = file_type
self.parent_id = parent_id
self.filename = filename
self.deleted = deleted
self.purged = purged
self.metadata_safe = metadata_safe
def get_file_name( self ):
#return absolute path of the filename
if self.filename:
return os.path.abspath( self.filename )
if self.dataset_id is not None:
return self.dataset.file_name
else:
assert self.id is not None, "ID must be set before filename used (commit the object)"
assert self.parent_id is not None, "Parent ID must be set before filename used"
return os.path.abspath( "%s_accociated_%s" % ( self.parent.file_name, self.id ) )
def set_file_name ( self, filename ):
self.filename = filename
if self.dataset:
self.dataset.deleted = True
self.dataset = None
self.dataset_id = None
file_name = property( get_file_name, set_file_name )
def clear( self, purge = False ):
self.deleted = True
if self.dataset:
self.dataset.deleted = True
self.dataset.purged = purge
if purge:
self.purged = True
try: os.unlink( self.file_name )
except Exception, e: print "Failed to purge associated file (%s) from disk: %s" % ( self.file_name, e )
class Event( object ):
def __init__( self, message=None, history=None, user=None, galaxy_session=None ):
self.history = history
+25 -1
View File
@@ -86,6 +86,18 @@ Dataset.table = Table( "dataset", metadata,
Column( 'file_size', Numeric( 15, 0 ) ),
ForeignKeyConstraint(['parent_id'],['dataset.id'], ondelete="CASCADE") )
DatasetAssociatedFile.table = Table( "dataset_associated_file", metadata,
Column( "id", Integer, primary_key=True ),
Column( "create_time", DateTime, default=now ),
Column( "update_time", DateTime, default=now, onupdate=now ),
Column( "dataset_id", Integer, ForeignKey( "dataset.id" ), index=True, nullable=True ),
Column( "parent_id", Integer, ForeignKey( "dataset.id" ), index=True ),
Column( "filename", TEXT ),
Column( "deleted", Boolean, index=True, default=False ),
Column( "purged", Boolean, index=True, default=False ),
Column( "metadata_safe", Boolean, index=True, default=True ),
Column( "type", TrimmedString( 255 ) ) )
DatasetFileName.table = Table( "dataset_filename", metadata,
Column( "id", Integer, primary_key=True ),
Column( "create_time", DateTime, default=now ),
@@ -233,7 +245,10 @@ assign_mapper( context, Dataset, Dataset.table,
backref="parent" ),
dataset_file=relation(
DatasetFileName,
primaryjoin=( DatasetFileName.table.c.id == Dataset.table.c.filename_id ) )
primaryjoin=( DatasetFileName.table.c.id == Dataset.table.c.filename_id ) ),
associated_files=relation(
DatasetAssociatedFile,
primaryjoin=( DatasetAssociatedFile.table.c.parent_id == Dataset.table.c.id ) )
) )
assign_mapper( context, DatasetFileName, DatasetFileName.table )
@@ -241,6 +256,15 @@ assign_mapper( context, DatasetFileName, DatasetFileName.table )
assign_mapper( context, DatasetChildAssociation, DatasetChildAssociation.table,
properties=dict( child=relation( Dataset, primaryjoin=( DatasetChildAssociation.table.c.child_dataset_id == Dataset.table.c.id ) ) ) )
assign_mapper( context, DatasetAssociatedFile, DatasetAssociatedFile.table,
properties=dict( parent=relation(
Dataset,
primaryjoin=( DatasetAssociatedFile.table.c.parent_id == Dataset.table.c.id ) ),
dataset=relation(
Dataset,
primaryjoin=( DatasetAssociatedFile.table.c.dataset_id == Dataset.table.c.id ) ) ) )
# assign_mapper( model.Query, model.Query.table,
# properties=dict( datasets=relation( model.Dataset.mapper, backref="query") ) )
+2 -2
View File
@@ -845,14 +845,14 @@ class Tool:
raise Exception( "Unexpected parameter type" )
return args
def execute( self, trans, incoming={} ):
def execute( self, trans, incoming={}, set_output_hid = True ):
"""
Execute the tool using parameter values in `incoming`. This just
dispatches to the `ToolAction` instance specified by
`self.tool_action`. In general this will create a `Job` that
when run will build the tool's outputs, e.g. `DefaultToolAction`.
"""
return self.tool_action.execute( self, trans, incoming )
return self.tool_action.execute( self, trans, incoming, set_output_hid = set_output_hid )
def params_to_strings( self, params, app ):
return params_to_strings( self.inputs, params, app )
+22 -6
View File
@@ -16,7 +16,7 @@ class ToolAction( object ):
class DefaultToolAction( object ):
"""Default tool action is to run an external command"""
def collect_input_datasets( self, tool, param_values ):
def collect_input_datasets( self, tool, param_values, trans ):
"""
Collect any dataset inputs from incoming. Returns a mapping from
parameter name to Dataset instance for each tool parameter that is
@@ -24,22 +24,38 @@ class DefaultToolAction( object ):
"""
input_datasets = dict()
def visitor( prefix, input, value ):
def converted_dataset( data ):
if data and not isinstance( data.datatype, input.formats ):
for target_ext in input.extensions:
if target_ext in data.get_converter_types():
assoc = data.get_associated_files_by_type( "CONVERTED_%s" % target_ext )
if assoc: data = assoc[0].dataset
else:
#run converter here
assoc = trans.app.model.DatasetAssociatedFile( parent_id = data.id, file_type = "CONVERTED_%s" % target_ext, metadata_safe = False )
new_data = data.datatype.convert_dataset( trans, data, target_ext, return_output = True, visible = False ).values()[0]
new_data.hid = data.hid
new_data.name = data.name
assoc.dataset_id = new_data.id
data = new_data
break
return data
if isinstance( input, DataToolParameter ):
if isinstance( value, list ):
# If there are multiple inputs with the same name, they
# are stored as name1, name2, ...
for i, v in enumerate( value ):
input_datasets[ prefix + input.name + str( i + 1 ) ] = v
input_datasets[ prefix + input.name + str( i + 1 ) ] = converted_dataset( v )
else:
input_datasets[ prefix + input.name ] = value
input_datasets[ prefix + input.name ] = converted_dataset( value )
tool.visit_inputs( param_values, visitor )
return input_datasets
def execute(self, tool, trans, incoming={} ):
def execute(self, tool, trans, incoming={}, set_output_hid = True ):
out_data = {}
# Collect any input datasets from the incoming parameters
inp_data = self.collect_input_datasets( tool, incoming )
inp_data = self.collect_input_datasets( tool, incoming, trans )
# Deal with input metadata, 'dbkey', names, and types
@@ -133,7 +149,7 @@ class DefaultToolAction( object ):
for name in out_data.keys():
if name not in child_dataset_names:
data = out_data[ name ]
trans.history.add_dataset( data )
trans.history.add_dataset( data, set_hid = set_output_hid )
data.flush()
# Add all the children to their parents
+1 -1
View File
@@ -14,7 +14,7 @@ class UploadToolAction( object ):
self.empty = False
self.line_count = None
def execute( self, tool, trans, incoming={} ):
def execute( self, tool, trans, incoming={}, set_output_hid = True ):
data_file = incoming['file_data']
file_type = incoming['file_type']
dbkey = incoming['dbkey']
+13 -3
View File
@@ -1037,11 +1037,21 @@ class DataToolParameter( ToolParameter ):
hid = "%s.%d" % ( parent_hid, i + 1 )
else:
hid = str( data.hid )
if isinstance( data.datatype, self.formats) and not data.deleted and data.state not in [data.states.ERROR]:
if not data.deleted and data.state not in [data.states.ERROR] and data.visible:
if self.options and filter_key == 'dbkey' and data.get_dbkey() != filter_value:
continue
selected = ( value and ( data in value ) )
field.add_option( "%s: %s" % ( hid, data.name[:30] ), data.id, selected )
if isinstance( data.datatype, self.formats):
selected = ( value and ( data in value ) )
field.add_option( "%s: %s" % ( hid, data.name[:30] ), data.id, selected )
else:
for target_ext in self.extensions:
if target_ext in data.get_converter_types():
assoc = data.get_associated_files_by_type( "CONVERTED_%s" % target_ext )
if assoc:
data = assoc[0].dataset
selected = ( value and ( data in value ) )
field.add_option( "%s: (as %s) %s" % ( hid, target_ext, data.name[:30] ), data.id, selected )
break #we only report the first valid converter, assume self.extensions is a priority list
# Also collect children via association object
dataset_collector( [ assoc.child for assoc in data.children ], hid )
dataset_collector( history.datasets, None )
+3
View File
@@ -287,6 +287,7 @@ class RootController( BaseController ):
# assert data.parent == None, "You must delete the primary dataset first."
# history.datasets.remove( data )
data.deleted = True
data.clear_associated_files()
data.flush()
trans.log_event( "Dataset id %s marked as deleted" % str(id) )
if data.parent_id is None:
@@ -310,6 +311,7 @@ class RootController( BaseController ):
# assert data.parent == None, "You must delete the primary dataset first."
# history.datasets.remove( data )
data.deleted = True
data.clear_associated_files()
data.flush()
trans.log_event( "Dataset id %s marked as deleted async" % str(id) )
if data.parent_id is None:
@@ -366,6 +368,7 @@ class RootController( BaseController ):
history = trans.get_history()
for dataset in history.datasets:
dataset.deleted = True
dataset.clear_associated_files()
self.app.model.flush()
trans.log_event( "History id %s cleared" % (str(history.id)) )
trans.response.send_redirect( url_for("/index" ) )
@@ -106,6 +106,7 @@ def delete_userless_histories( h, cutoff_time ):
for dataset in history.datasets:
if not dataset.deleted:
dataset.deleted = True
dataset.clear_associated_files()
dataset.flush()
print "dataset_%d" %dataset.id
dataset_count += 1
@@ -258,6 +259,7 @@ def purge_dataset( dataset ):
os.unlink( dataset.file_name )
dataset.purged = True
dataset.file_size = 0
dataset.clear_associated_files( purge = True )
dataset.flush()
except Exception, exc:
return "# Error, exception: %s caught attempting to purge %s\n" %( str( exc ), dataset.file_name )