From a7ea9357f20fe19dbb6a8f8f216570b41cf281cc Mon Sep 17 00:00:00 2001 From: Daniel Blankenberg Date: Thu, 15 May 2008 19:09:17 +0000 Subject: [PATCH] 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. --- lib/galaxy/datatypes/data.py | 16 ++++-- lib/galaxy/model/__init__.py | 60 ++++++++++++++++++-- lib/galaxy/model/mapping.py | 26 ++++++++- lib/galaxy/tools/__init__.py | 4 +- lib/galaxy/tools/actions/__init__.py | 28 +++++++-- lib/galaxy/tools/actions/upload.py | 2 +- lib/galaxy/tools/parameters/basic.py | 16 +++++- lib/galaxy/web/controllers/root.py | 3 + scripts/cleanup_datasets/cleanup_datasets.py | 2 + 9 files changed, 135 insertions(+), 22 deletions(-) diff --git a/lib/galaxy/datatypes/data.py b/lib/galaxy/datatypes/data.py index d1df4f183b3..3fc880f9143 100644 --- a/lib/galaxy/datatypes/data.py +++ b/lib/galaxy/datatypes/data.py @@ -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 ): diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index b45643e0b38..eba267448c4 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -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 diff --git a/lib/galaxy/model/mapping.py b/lib/galaxy/model/mapping.py index 0e21c2b1c50..34af54e94b7 100644 --- a/lib/galaxy/model/mapping.py +++ b/lib/galaxy/model/mapping.py @@ -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") ) ) diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 00696f28f12..22cead8f92d 100644 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -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 ) diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index 5518bd24d16..26d7ae1592d 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -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 diff --git a/lib/galaxy/tools/actions/upload.py b/lib/galaxy/tools/actions/upload.py index 62a4c59b48e..6e2e16ba285 100644 --- a/lib/galaxy/tools/actions/upload.py +++ b/lib/galaxy/tools/actions/upload.py @@ -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'] diff --git a/lib/galaxy/tools/parameters/basic.py b/lib/galaxy/tools/parameters/basic.py index 399131741e3..beb459958db 100644 --- a/lib/galaxy/tools/parameters/basic.py +++ b/lib/galaxy/tools/parameters/basic.py @@ -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 ) diff --git a/lib/galaxy/web/controllers/root.py b/lib/galaxy/web/controllers/root.py index 314f0c1e6cd..1ed2742308d 100644 --- a/lib/galaxy/web/controllers/root.py +++ b/lib/galaxy/web/controllers/root.py @@ -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" ) ) diff --git a/scripts/cleanup_datasets/cleanup_datasets.py b/scripts/cleanup_datasets/cleanup_datasets.py index 1764b6a7357..4738b1173d3 100644 --- a/scripts/cleanup_datasets/cleanup_datasets.py +++ b/scripts/cleanup_datasets/cleanup_datasets.py @@ -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 )