From 412f779f768a003bd3160048a80216bfaa6264cf Mon Sep 17 00:00:00 2001 From: Richard Park Date: Wed, 1 Feb 2012 16:21:38 -0500 Subject: [PATCH 01/87] Updated scripts for API workflow enhancements and changing workflow parameters programatically --- .hgignore | 4 + lib/galaxy/web/api/workflows.py | 378 ++++++++++++++++++ lib/galaxy/web/buildapp.py | 14 + scripts/api/workflow_delete_workflow_rpark.py | 22 + scripts/api/workflow_execute_rpark.py | 70 ++++ .../api/workflow_import_from_file_rpark.py | 39 ++ 6 files changed, 527 insertions(+) create mode 100644 scripts/api/workflow_delete_workflow_rpark.py create mode 100644 scripts/api/workflow_execute_rpark.py create mode 100644 scripts/api/workflow_import_from_file_rpark.py diff --git a/.hgignore b/.hgignore index 070d85b9761..31f66fa33f2 100644 --- a/.hgignore +++ b/.hgignore @@ -78,3 +78,7 @@ static/june_2007_style/blue/base_sprites.less *.rej *~ +syntax: regexp +^database$ +syntax: regexp +^scripts/api/spp_submodule\.ga$ diff --git a/lib/galaxy/web/api/workflows.py b/lib/galaxy/web/api/workflows.py index ce356da86ca..193d19c9322 100644 --- a/lib/galaxy/web/api/workflows.py +++ b/lib/galaxy/web/api/workflows.py @@ -84,6 +84,20 @@ class WorkflowsAPIController(BaseAPIController): However, we will import them if installed_repository_file is specified """ + + # ------------------------------------------------------------------------------- # + ### RPARK: dictionary containing which workflows to change and edit ### + param_map = {}; + if (payload.has_key('parameters') ): + #if (payload['parameters']): + param_map = payload['parameters']; + print("PARAMETER MAP:"); + print(param_map); + # ------------------------------------------------------------------------------- # + + + + if 'workflow_id' not in payload: # create new if 'installed_repository_file' in payload: @@ -168,6 +182,30 @@ class WorkflowsAPIController(BaseAPIController): # are not persisted so we need to do it every time) step.module.add_dummy_datasets( connections=step.input_connections ) step.state = step.module.state + + #################################################### + #################################################### + #print("CHECKING WORKFLOW STEPS:") + #print(step.tool_id); + #print(step.state.inputs); + #print("upgard messages"); + #print(step.state); + #print("\n"); + # RPARK: IF TOOL_NAME IN PARAMETER MAP # + if step.tool_id in param_map: + #print("-------------------------FOUND IN PARAMETER DICTIONARY") + #print(param_map[step.tool_id]); + change_param = param_map[step.tool_id]['param']; + change_value = param_map[step.tool_id]['value']; + #step.state.inputs['refGenomeSource']['index'] = "crapolo"; + #print(step.state.inputs[change_param]); + step.state.inputs[change_param] = change_value; + #print(step.state.inputs[change_param]); + #print(param_map[step.tool_id][change_value]); + #print("--------------------------------------------------") + #################################################### + #################################################### + if step.tool_errors: trans.response.status = 400 return "Workflow cannot be run because of validation errors in some steps: %s" % step_errors @@ -220,3 +258,343 @@ class WorkflowsAPIController(BaseAPIController): trans.sa_session.flush() return rval + # ---------------------------------------------------------------------------------------------- # + # ---------------------------------------------------------------------------------------------- # + # ---- RPARK EDITS ---- # + # ---------------------------------------------------------------------------------------------- # + # ---------------------------------------------------------------------------------------------- # + @web.expose_api + @web.json + def workflow_dict( self, trans, workflow_id, **kwd ): + """ + GET /api/workflows/{encoded_workflow_id}/download + Returns a selected workflow as a json dictionary. + """ + print "workflow controller: workflow dict called" + print workflow_id + + try: + stored_workflow = trans.sa_session.query(self.app.model.StoredWorkflow).get(trans.security.decode_id(workflow_id)) + except Exception,e: + return ("Workflow with ID='%s' can not be found\n Exception: %s") % (workflow_id, str( e )) + + # check to see if user has permissions to selected workflow + if stored_workflow.user != trans.user and not trans.user_is_admin(): + if trans.sa_session.query(trans.app.model.StoredWorkflowUserShareAssociation).filter_by(user=trans.user, stored_workflow=stored_workflow).count() == 0: + trans.response.status = 400 + return("Workflow is not owned by or shared with current user") + + return self._workflow_to_dict( trans, stored_workflow ) + + @web.expose_api + def delete( self, trans, id, **kwd ): + """ + DELETE /api/workflows/{encoded_workflow_id} + Deletes a specified workflow + Author: rpark + + copied from galaxy.web.controllers.workflows.py (delete) + """ + workflow_id = id; + + try: + stored_workflow = trans.sa_session.query(self.app.model.StoredWorkflow).get(trans.security.decode_id(workflow_id)) + except Exception,e: + return ("Workflow with ID='%s' can not be found\n Exception: %s") % (workflow_id, str( e )) + + # check to see if user has permissions to selected workflow + if stored_workflow.user != trans.user and not trans.user_is_admin(): + if trans.sa_session.query(trans.app.model.StoredWorkflowUserShareAssociation).filter_by(user=trans.user, stored_workflow=stored_workflow).count() == 0: + trans.response.status = 400 + return("Workflow is not owned by or shared with current user") + + #Mark a workflow as deleted + stored_workflow.deleted = True + trans.sa_session.flush() + + # Python Debugger + #import pdb; pdb.set_trace() + + # TODO: Unsure of response message to let api know that a workflow was successfully deleted + #return 'OK' + return ( "Workflow '%s' successfully deleted" % stored_workflow.name ) + + @web.expose_api + def import_new_workflow(self, trans, payload, **kwd): + """ + POST /api/workflows + Importing dynamic workflows from the api. Return newly generated workflow id. + Author: rpark + + # currently assumes payload['workflow'] is a json representation of a workflow to be inserted into the database + """ + + #import pdb; pdb.set_trace() + + data = payload['workflow']; + workflow, missing_tool_tups = self._workflow_from_dict( trans, data, source="API" ) + + # galaxy workflow newly created id + workflow_id = workflow.id; + # api encoded, id + encoded_id = trans.security.encode_id(workflow_id); + + + + # return list + rval= []; + + item = workflow.get_api_value(value_mapper={'id':trans.security.encode_id}) + item['url'] = url_for('workflow', id=encoded_id) + + rval.append(item); + + return rval; + + + def _workflow_from_dict( self, trans, data, source=None ): + """ + RPARK: copied from galaxy.web.controllers.workflows.py + Creates a workflow from a dict. Created workflow is stored in the database and returned. + """ + # Put parameters in workflow mode + trans.workflow_building_mode = True + # Create new workflow from incoming dict + workflow = model.Workflow() + # If there's a source, put it in the workflow name. + if source: + name = "%s (imported from %s)" % ( data['name'], source ) + else: + name = data['name'] + workflow.name = name + # Assume no errors until we find a step that has some + workflow.has_errors = False + # Create each step + steps = [] + # The editor will provide ids for each step that we don't need to save, + # but do need to use to make connections + steps_by_external_id = {} + # Keep track of tools required by the workflow that are not available in + # the local Galaxy instance. Each tuple in the list of missing_tool_tups + # will be ( tool_id, tool_name, tool_version ). + missing_tool_tups = [] + # First pass to build step objects and populate basic values + for key, step_dict in data[ 'steps' ].iteritems(): + # Create the model class for the step + step = model.WorkflowStep() + steps.append( step ) + steps_by_external_id[ step_dict['id' ] ] = step + # FIXME: Position should be handled inside module + step.position = step_dict['position'] + module = module_factory.from_dict( trans, step_dict, secure=False ) + if module.type == 'tool' and module.tool is None: + # A required tool is not available in the local Galaxy instance. + missing_tool_tup = ( step_dict[ 'tool_id' ], step_dict[ 'name' ], step_dict[ 'tool_version' ] ) + if missing_tool_tup not in missing_tool_tups: + missing_tool_tups.append( missing_tool_tup ) + module.save_to_step( step ) + if step.tool_errors: + workflow.has_errors = True + # Stick this in the step temporarily + step.temp_input_connections = step_dict['input_connections'] + + # Save step annotation. + annotation = step_dict[ 'annotation' ] + if annotation: + annotation = sanitize_html( annotation, 'utf-8', 'text/html' ) + # ------------------------------------------ # + # RPARK REMOVING: user annotation b/c of API + #self.add_item_annotation( trans.sa_session, trans.get_user(), step, annotation ) + # ------------------------------------------ # + + # Unpack and add post-job actions. + post_job_actions = step_dict.get( 'post_job_actions', {} ) + for name, pja_dict in post_job_actions.items(): + pja = PostJobAction( pja_dict[ 'action_type' ], + step, pja_dict[ 'output_name' ], + pja_dict[ 'action_arguments' ] ) + # Second pass to deal with connections between steps + for step in steps: + # Input connections + for input_name, conn_dict in step.temp_input_connections.iteritems(): + if conn_dict: + conn = model.WorkflowStepConnection() + conn.input_step = step + conn.input_name = input_name + conn.output_name = conn_dict['output_name'] + conn.output_step = steps_by_external_id[ conn_dict['id'] ] + del step.temp_input_connections + # Order the steps if possible + attach_ordered_steps( workflow, steps ) + # Connect up + stored = model.StoredWorkflow() + stored.name = workflow.name + workflow.stored_workflow = stored + stored.latest_workflow = workflow + stored.user = trans.user + # Persist + trans.sa_session.add( stored ) + trans.sa_session.flush() + return stored, missing_tool_tups + + def _workflow_to_dict( self, trans, stored ): + """ + RPARK: copied from galaxy.web.controllers.workflows.py + Converts a workflow to a dict of attributes suitable for exporting. + """ + workflow = stored.latest_workflow + + ### ----------------------------------- ### + ## RPARK EDIT ## + workflow_annotation = self.get_item_annotation_obj( trans.sa_session, trans.user, stored ) + annotation_str = "" + if workflow_annotation: + annotation_str = workflow_annotation.annotation + ### ----------------------------------- ### + + + # Pack workflow data into a dictionary and return + data = {} + data['a_galaxy_workflow'] = 'true' # Placeholder for identifying galaxy workflow + data['format-version'] = "0.1" + data['name'] = workflow.name + ### ----------------------------------- ### + ## RPARK EDIT ## + data['annotation'] = annotation_str + ### ----------------------------------- ### + + data['steps'] = {} + # For each step, rebuild the form and encode the state + for step in workflow.steps: + # Load from database representation + module = module_factory.from_workflow_step( trans, step ) + + ### ----------------------------------- ### + ## RPARK EDIT ## + # Get user annotation. + step_annotation = self.get_item_annotation_obj(trans.sa_session, trans.user, step ) + annotation_str = "" + if step_annotation: + annotation_str = step_annotation.annotation + ### ----------------------------------- ### + + # Step info + step_dict = { + 'id': step.order_index, + 'type': module.type, + 'tool_id': module.get_tool_id(), + 'tool_version' : step.tool_version, + 'name': module.get_name(), + 'tool_state': module.get_state( secure=False ), + 'tool_errors': module.get_errors(), + ## 'data_inputs': module.get_data_inputs(), + ## 'data_outputs': module.get_data_outputs(), + + ### ----------------------------------- ### + ## RPARK EDIT ## + 'annotation' : annotation_str + ### ----------------------------------- ### + + } + # Add post-job actions to step dict. + if module.type == 'tool': + pja_dict = {} + for pja in step.post_job_actions: + pja_dict[pja.action_type+pja.output_name] = dict( action_type = pja.action_type, + output_name = pja.output_name, + action_arguments = pja.action_arguments ) + step_dict[ 'post_job_actions' ] = pja_dict + # Data inputs + step_dict['inputs'] = [] + if module.type == "data_input": + # Get input dataset name; default to 'Input Dataset' + name = module.state.get( 'name', 'Input Dataset') + step_dict['inputs'].append( { "name" : name, "description" : annotation_str } ) + else: + # Step is a tool and may have runtime inputs. + for name, val in module.state.inputs.items(): + input_type = type( val ) + if input_type == RuntimeValue: + step_dict['inputs'].append( { "name" : name, "description" : "runtime parameter for tool %s" % module.get_name() } ) + elif input_type == dict: + # Input type is described by a dict, e.g. indexed parameters. + for partname, partval in val.items(): + if type( partval ) == RuntimeValue: + step_dict['inputs'].append( { "name" : name, "description" : "runtime parameter for tool %s" % module.get_name() } ) + # User outputs + step_dict['user_outputs'] = [] + """ + module_outputs = module.get_data_outputs() + step_outputs = trans.sa_session.query( WorkflowOutput ).filter( step=step ) + for output in step_outputs: + name = output.output_name + annotation = "" + for module_output in module_outputs: + if module_output.get( 'name', None ) == name: + output_type = module_output.get( 'extension', '' ) + break + data['outputs'][name] = { 'name' : name, 'annotation' : annotation, 'type' : output_type } + """ + + # All step outputs + step_dict['outputs'] = [] + if type( module ) is ToolModule: + for output in module.get_data_outputs(): + step_dict['outputs'].append( { 'name' : output['name'], 'type' : output['extensions'][0] } ) + # Connections + input_connections = step.input_connections + if step.type is None or step.type == 'tool': + # Determine full (prefixed) names of valid input datasets + data_input_names = {} + def callback( input, value, prefixed_name, prefixed_label ): + if isinstance( input, DataToolParameter ): + data_input_names[ prefixed_name ] = True + visit_input_values( module.tool.inputs, module.state.inputs, callback ) + # Filter + # FIXME: this removes connection without displaying a message currently! + input_connections = [ conn for conn in input_connections if conn.input_name in data_input_names ] + # Encode input connections as dictionary + input_conn_dict = {} + for conn in input_connections: + input_conn_dict[ conn.input_name ] = \ + dict( id=conn.output_step.order_index, output_name=conn.output_name ) + step_dict['input_connections'] = input_conn_dict + # Position + step_dict['position'] = step.position + # Add to return value + data['steps'][step.order_index] = step_dict + return data + + def get_item_annotation_obj( self, db_session, user, item ): + """ + RPARK: copied from galaxy.model.item_attr.py + Returns a user's annotation object for an item. """ + # Get annotation association class. + annotation_assoc_class = self._get_annotation_assoc_class( item ) + if not annotation_assoc_class: + return None + + # Get annotation association object. + annotation_assoc = db_session.query( annotation_assoc_class ).filter_by( user=user ) + + # TODO: use filtering like that in _get_item_id_filter_str() + if item.__class__ == galaxy.model.History: + annotation_assoc = annotation_assoc.filter_by( history=item ) + elif item.__class__ == galaxy.model.HistoryDatasetAssociation: + annotation_assoc = annotation_assoc.filter_by( hda=item ) + elif item.__class__ == galaxy.model.StoredWorkflow: + annotation_assoc = annotation_assoc.filter_by( stored_workflow=item ) + elif item.__class__ == galaxy.model.WorkflowStep: + annotation_assoc = annotation_assoc.filter_by( workflow_step=item ) + elif item.__class__ == galaxy.model.Page: + annotation_assoc = annotation_assoc.filter_by( page=item ) + elif item.__class__ == galaxy.model.Visualization: + annotation_assoc = annotation_assoc.filter_by( visualization=item ) + return annotation_assoc.first() + + def _get_annotation_assoc_class( self, item ): + """ + RPARK: copied from galaxy.model.item_attr.py + Returns an item's item-annotation association class. """ + class_name = '%sAnnotationAssociation' % item.__class__.__name__ + return getattr( galaxy.model, class_name, None ) diff --git a/lib/galaxy/web/buildapp.py b/lib/galaxy/web/buildapp.py index b418134d4b3..92038157877 100644 --- a/lib/galaxy/web/buildapp.py +++ b/lib/galaxy/web/buildapp.py @@ -151,6 +151,20 @@ def app_factory( global_conf, **kwargs ): webapp.api_mapper.resource_with_deleted( 'history', 'histories', path_prefix='/api' ) #webapp.api_mapper.connect( 'run_workflow', '/api/workflow/{workflow_id}/library/{library_id}', controller='workflows', action='run', workflow_id=None, library_id=None, conditions=dict(method=["GET"]) ) + # ---------------------------------------------- # + # ---------------------------------------------- # + # RPARK EDIT + + # How to extend API: url_mapping + # "POST /api/workflows/import" => ``workflows.import_workflow()``. + # Defines a named route "import_workflow". + webapp.api_mapper.connect("import_workflow", "/api/workflows/upload", controller="workflows", action="import_new_workflow", conditions=dict(method=["POST"])) + webapp.api_mapper.connect("workflow_dict", '/api/workflows/download/{workflow_id}', controller='workflows', action='workflow_dict', conditions=dict(method=['GET'])) + + #import pdb; pdb.set_trace() + # ---------------------------------------------- # + # ---------------------------------------------- # + webapp.finalize_config() # Wrap the webapp in some useful middleware if kwargs.get( 'middleware', True ): diff --git a/scripts/api/workflow_delete_workflow_rpark.py b/scripts/api/workflow_delete_workflow_rpark.py new file mode 100644 index 00000000000..734e12ca8f5 --- /dev/null +++ b/scripts/api/workflow_delete_workflow_rpark.py @@ -0,0 +1,22 @@ +#!/usr/bin/env python +""" +# Author: RPARK +API script for deleting workflows +""" + +import os, sys +sys.path.insert( 0, os.path.dirname( __file__ ) ) +from common import delete + +try: + assert sys.argv[2] +except IndexError: + print 'usage: %s key url [purge (true/false)] ' % os.path.basename( sys.argv[0] ) + sys.exit( 1 ) +try: + data = {} + data[ 'purge' ] = sys.argv[3] +except IndexError: + pass + +delete( sys.argv[1], sys.argv[2], data ) diff --git a/scripts/api/workflow_execute_rpark.py b/scripts/api/workflow_execute_rpark.py new file mode 100644 index 00000000000..a8ca153d9dd --- /dev/null +++ b/scripts/api/workflow_execute_rpark.py @@ -0,0 +1,70 @@ +#!/usr/bin/env python +""" +Execute workflows from the command line. +Example calls: +python workflow_execute.py /api/workflows f2db41e1fa331b3e 'Test API History' '38=ldda=0qr350234d2d192f' +python workflow_execute.py /api/workflows f2db41e1fa331b3e 'hist_id=a912e9e5d84530d4' '38=hda=03501d7626bd192f' +""" + +""" +python workflow_execute.py /api/workflows f2db41e1fa331b3e 'hist_id=a912e9e5d84530d4' '38=hda=03501d7626bd192f' 'param=tool=name=value' + +'param=tool=name=value' + +Example +python workflow_execute_rpark.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=bwa' + +python workflow_execute_rpark.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=arachne' 'param=bowtie_wrapper=suppressHeader=True' + +python workflow_execute_rpark.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=bowtie' 'param=bowtie_wrapper=suppressHeader=True' 'param=peakcalling_spp=window_size=1000' + +""" + +import os, sys +sys.path.insert( 0, os.path.dirname( __file__ ) ) +from common import submit + + +def main(): + try: + print("workflow_execute:py:"); + data = {} + data['workflow_id'] = sys.argv[3] + data['history'] = sys.argv[4] + data['ds_map'] = {} + + ######################################################### + ### MY EDITS ############################################ + ### Trying to pass in parameter for my own dictionary ### + data['parameters'] = {}; + + # DBTODO If only one input is given, don't require a step + # mapping, just use it for everything? + for v in sys.argv[5:]: + print("Multiple arguments "); + print(v); + + try: + step, src, ds_id = v.split('='); + data['ds_map'][step] = {'src':src, 'id':ds_id}; + + except ValueError: + print("VALUE ERROR:"); + wtype, wtool, wparam, wvalue = v.split('='); + data['parameters'][wtool] = {'param':wparam, 'value':wvalue} + + + ######################################################### + ### MY EDITS ############################################ + ### Trying to pass in parameter for my own dictionary ### + #data['parameters']['bowtie'] = {'param':'stepSize', 'value':100} + #data['parameters']['sam_to_bam'] = {'param':'genome', 'value':'hg18'} + + except IndexError: + print 'usage: %s key url workflow_id history step=src=dataset_id' % os.path.basename(sys.argv[0]) + sys.exit(1) + submit( sys.argv[1], sys.argv[2], data ) + +if __name__ == '__main__': + main() + diff --git a/scripts/api/workflow_import_from_file_rpark.py b/scripts/api/workflow_import_from_file_rpark.py new file mode 100644 index 00000000000..3443d3989fa --- /dev/null +++ b/scripts/api/workflow_import_from_file_rpark.py @@ -0,0 +1,39 @@ +#!/usr/bin/env python + +""" + +python rpark_import_workflow_from_file.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows/import 'spp_submodule.ga' +python rpark_import_workflow_from_file.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows/import 'spp_submodule.ga' +""" + +import os, sys +sys.path.insert( 0, os.path.dirname( __file__ ) ) +from common import submit + +### Rpark edit ### +import simplejson + +def openWorkflow(in_file): + with open(in_file) as f: + temp_data = simplejson.load(f) + return temp_data; + + + +try: + assert sys.argv[2] +except IndexError: + print 'usage: %s key url [name] ' % os.path.basename( sys.argv[0] ) + sys.exit( 1 ) +try: + #data = {} + #data[ 'name' ] = sys.argv[3] + data = {}; + workflow_dict = openWorkflow(sys.argv[3]); + data ['workflow'] = workflow_dict; + + +except IndexError: + pass + +submit( sys.argv[1], sys.argv[2], data ) From 9fb4df9dd811d1c219c5865da31946d32137220a Mon Sep 17 00:00:00 2001 From: Richard Park Date: Wed, 1 Feb 2012 16:28:00 -0500 Subject: [PATCH 02/87] Updated import statements in workflows api controller --- lib/galaxy/web/api/workflows.py | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/lib/galaxy/web/api/workflows.py b/lib/galaxy/web/api/workflows.py index 193d19c9322..82f76aea20a 100644 --- a/lib/galaxy/web/api/workflows.py +++ b/lib/galaxy/web/api/workflows.py @@ -11,6 +11,21 @@ from galaxy.web.base.controller import BaseAPIController, url_for from galaxy.workflow.modules import module_factory from galaxy.jobs.actions.post import ActionBox +# ---------------------------------------------------------------------------------------------- # +# ---------------------------------------------------------------------------------------------- # +# ---- RPARK EDITS ---- # +import pkg_resources +pkg_resources.require( "simplejson" ) +from galaxy import model +from galaxy.web.controllers.workflow import attach_ordered_steps +from galaxy.util.sanitize_html import sanitize_html +from galaxy.workflow.modules import * +from galaxy.model.item_attrs import * + +# ---------------------------------------------------------------------------------------------- # +# ---------------------------------------------------------------------------------------------- # + + log = logging.getLogger(__name__) class WorkflowsAPIController(BaseAPIController): From cfef7805231c3e443ecc6ea981285eea61fada2c Mon Sep 17 00:00:00 2001 From: Richard Park Date: Wed, 1 Feb 2012 17:22:50 -0500 Subject: [PATCH 03/87] updated workflow_dict function for returning a selected workflow as a json object via API --- lib/galaxy/web/api/workflows.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/lib/galaxy/web/api/workflows.py b/lib/galaxy/web/api/workflows.py index 82f76aea20a..d670349e068 100644 --- a/lib/galaxy/web/api/workflows.py +++ b/lib/galaxy/web/api/workflows.py @@ -285,8 +285,6 @@ class WorkflowsAPIController(BaseAPIController): GET /api/workflows/{encoded_workflow_id}/download Returns a selected workflow as a json dictionary. """ - print "workflow controller: workflow dict called" - print workflow_id try: stored_workflow = trans.sa_session.query(self.app.model.StoredWorkflow).get(trans.security.decode_id(workflow_id)) From 6cc8c313ef8c9fddc389e101fb6025fba7bf3860 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Wed, 8 Feb 2012 17:23:58 -0500 Subject: [PATCH 04/87] Updated import new workflow function --- lib/galaxy/web/api/workflows.py | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/lib/galaxy/web/api/workflows.py b/lib/galaxy/web/api/workflows.py index d670349e068..3b3fef670e0 100644 --- a/lib/galaxy/web/api/workflows.py +++ b/lib/galaxy/web/api/workflows.py @@ -279,7 +279,7 @@ class WorkflowsAPIController(BaseAPIController): # ---------------------------------------------------------------------------------------------- # # ---------------------------------------------------------------------------------------------- # @web.expose_api - @web.json + #@web.json def workflow_dict( self, trans, workflow_id, **kwd ): """ GET /api/workflows/{encoded_workflow_id}/download @@ -297,8 +297,9 @@ class WorkflowsAPIController(BaseAPIController): trans.response.status = 400 return("Workflow is not owned by or shared with current user") - return self._workflow_to_dict( trans, stored_workflow ) - + ret_dict = self._workflow_to_dict( trans, stored_workflow ); + return ret_dict + @web.expose_api def delete( self, trans, id, **kwd ): """ @@ -352,8 +353,6 @@ class WorkflowsAPIController(BaseAPIController): # api encoded, id encoded_id = trans.security.encode_id(workflow_id); - - # return list rval= []; @@ -362,8 +361,7 @@ class WorkflowsAPIController(BaseAPIController): rval.append(item); - return rval; - + return item; def _workflow_from_dict( self, trans, data, source=None ): """ From 392625681d7746632589de167464ccb5c4773d08 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Fri, 16 Mar 2012 17:23:15 -0400 Subject: [PATCH 05/87] Updated workflow API --- lib/galaxy/web/api/workflows.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/web/api/workflows.py b/lib/galaxy/web/api/workflows.py index 3b3fef670e0..b56b17803ee 100644 --- a/lib/galaxy/web/api/workflows.py +++ b/lib/galaxy/web/api/workflows.py @@ -336,7 +336,7 @@ class WorkflowsAPIController(BaseAPIController): @web.expose_api def import_new_workflow(self, trans, payload, **kwd): """ - POST /api/workflows + POST /api/workflows/upload Importing dynamic workflows from the api. Return newly generated workflow id. Author: rpark @@ -344,7 +344,7 @@ class WorkflowsAPIController(BaseAPIController): """ #import pdb; pdb.set_trace() - + data = payload['workflow']; workflow, missing_tool_tups = self._workflow_from_dict( trans, data, source="API" ) From e6d5c55c8945a38640ebaf9bbdc416d205781526 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 10 Apr 2012 23:00:41 -0400 Subject: [PATCH 06/87] Updated notes on how to run api/workflow_execute_parameters.py --- ...flow_execute_rpark.py => workflow_execute_parameters.py} | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) rename scripts/api/{workflow_execute_rpark.py => workflow_execute_parameters.py} (73%) diff --git a/scripts/api/workflow_execute_rpark.py b/scripts/api/workflow_execute_parameters.py similarity index 73% rename from scripts/api/workflow_execute_rpark.py rename to scripts/api/workflow_execute_parameters.py index a8ca153d9dd..47a87980e16 100644 --- a/scripts/api/workflow_execute_rpark.py +++ b/scripts/api/workflow_execute_parameters.py @@ -12,11 +12,11 @@ python workflow_execute.py /api/workflows f2db41e1fa331b3e 'param=tool=name=value' Example -python workflow_execute_rpark.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=bwa' +python workflow_execute_parameters.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=bwa' -python workflow_execute_rpark.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=arachne' 'param=bowtie_wrapper=suppressHeader=True' +python workflow_execute_parameters.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=arachne' 'param=bowtie_wrapper=suppressHeader=True' -python workflow_execute_rpark.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=bowtie' 'param=bowtie_wrapper=suppressHeader=True' 'param=peakcalling_spp=window_size=1000' +python workflow_execute_parameters.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=bowtie' 'param=bowtie_wrapper=suppressHeader=True' 'param=peakcalling_spp=window_size=1000' """ From 1ce82abc842c65595712e9f9f658bfefec7f8ea4 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 10 Apr 2012 23:05:23 -0400 Subject: [PATCH 07/87] Cleaned up code for api/workflows.py --- lib/galaxy/web/api/workflows.py | 16 ---------------- 1 file changed, 16 deletions(-) diff --git a/lib/galaxy/web/api/workflows.py b/lib/galaxy/web/api/workflows.py index b56b17803ee..4d025d93fdf 100644 --- a/lib/galaxy/web/api/workflows.py +++ b/lib/galaxy/web/api/workflows.py @@ -104,10 +104,7 @@ class WorkflowsAPIController(BaseAPIController): ### RPARK: dictionary containing which workflows to change and edit ### param_map = {}; if (payload.has_key('parameters') ): - #if (payload['parameters']): param_map = payload['parameters']; - print("PARAMETER MAP:"); - print(param_map); # ------------------------------------------------------------------------------- # @@ -200,24 +197,11 @@ class WorkflowsAPIController(BaseAPIController): #################################################### #################################################### - #print("CHECKING WORKFLOW STEPS:") - #print(step.tool_id); - #print(step.state.inputs); - #print("upgard messages"); - #print(step.state); - #print("\n"); # RPARK: IF TOOL_NAME IN PARAMETER MAP # if step.tool_id in param_map: - #print("-------------------------FOUND IN PARAMETER DICTIONARY") - #print(param_map[step.tool_id]); change_param = param_map[step.tool_id]['param']; change_value = param_map[step.tool_id]['value']; - #step.state.inputs['refGenomeSource']['index'] = "crapolo"; - #print(step.state.inputs[change_param]); step.state.inputs[change_param] = change_value; - #print(step.state.inputs[change_param]); - #print(param_map[step.tool_id][change_value]); - #print("--------------------------------------------------") #################################################### #################################################### From 683982921142ec6d753239522ee9fdcbfd0f0170 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Fri, 3 Aug 2012 02:09:05 -0400 Subject: [PATCH 08/87] Cleaned up codebase for attempting adding api additions to the galaxy main branch --- ...e_workflow_rpark.py => workflow_delete.py} | 9 +++-- scripts/api/workflow_execute_parameters.py | 35 ++++++------------- 2 files changed, 17 insertions(+), 27 deletions(-) rename scripts/api/{workflow_delete_workflow_rpark.py => workflow_delete.py} (59%) diff --git a/scripts/api/workflow_delete_workflow_rpark.py b/scripts/api/workflow_delete.py similarity index 59% rename from scripts/api/workflow_delete_workflow_rpark.py rename to scripts/api/workflow_delete.py index 734e12ca8f5..3b6d5a872bc 100644 --- a/scripts/api/workflow_delete_workflow_rpark.py +++ b/scripts/api/workflow_delete.py @@ -1,7 +1,12 @@ #!/usr/bin/env python """ -# Author: RPARK -API script for deleting workflows +# ---------------------------------------------- # +# PARKLAB, Author: RPARK +API example script for deleting workflows +# ---------------------------------------------- # + +Example calls: +python workflow_delete.py /api/workflows/ True """ import os, sys diff --git a/scripts/api/workflow_execute_parameters.py b/scripts/api/workflow_execute_parameters.py index 47a87980e16..7ce37db885e 100644 --- a/scripts/api/workflow_execute_parameters.py +++ b/scripts/api/workflow_execute_parameters.py @@ -1,23 +1,13 @@ #!/usr/bin/env python """ +# ---------------------------------------------- # +# PARKLAB, Author: RPARK +# ---------------------------------------------- # + Execute workflows from the command line. Example calls: -python workflow_execute.py /api/workflows f2db41e1fa331b3e 'Test API History' '38=ldda=0qr350234d2d192f' -python workflow_execute.py /api/workflows f2db41e1fa331b3e 'hist_id=a912e9e5d84530d4' '38=hda=03501d7626bd192f' -""" - -""" -python workflow_execute.py /api/workflows f2db41e1fa331b3e 'hist_id=a912e9e5d84530d4' '38=hda=03501d7626bd192f' 'param=tool=name=value' - -'param=tool=name=value' - -Example -python workflow_execute_parameters.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=bwa' - -python workflow_execute_parameters.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=arachne' 'param=bowtie_wrapper=suppressHeader=True' - -python workflow_execute_parameters.py 35a24ae2643785ff3d046c98ea362c7f http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=bowtie' 'param=bowtie_wrapper=suppressHeader=True' 'param=peakcalling_spp=window_size=1000' - +python workflow_execute.py /api/workflows 'hist_id=' '38=hda=' 'param=tool=name=value' +python workflow_execute_parameters.py http://localhost:8080/api/workflows 1cd8e2f6b131e891 'Test API' '69=ld=a799d38679e985db' '70=ld=33b43b4e7093c91f' 'param=peakcalling_spp=aligner=bowtie' 'param=bowtie_wrapper=suppressHeader=True' 'param=peakcalling_spp=window_size=1000' """ import os, sys @@ -34,7 +24,6 @@ def main(): data['ds_map'] = {} ######################################################### - ### MY EDITS ############################################ ### Trying to pass in parameter for my own dictionary ### data['parameters'] = {}; @@ -51,14 +40,10 @@ def main(): except ValueError: print("VALUE ERROR:"); wtype, wtool, wparam, wvalue = v.split('='); - data['parameters'][wtool] = {'param':wparam, 'value':wvalue} - - - ######################################################### - ### MY EDITS ############################################ - ### Trying to pass in parameter for my own dictionary ### - #data['parameters']['bowtie'] = {'param':'stepSize', 'value':100} - #data['parameters']['sam_to_bam'] = {'param':'genome', 'value':'hg18'} + try: + data['parameters'][wtool] = {'param':wparam, 'value':wvalue} + except ValueError: + print("TOOL ID ERROR:"); except IndexError: print 'usage: %s key url workflow_id history step=src=dataset_id' % os.path.basename(sys.argv[0]) From 47b63cbf9121a4b0591aa4cd67cd268030624419 Mon Sep 17 00:00:00 2001 From: Enis Afgan Date: Tue, 14 Aug 2012 10:04:25 +1000 Subject: [PATCH 09/87] Add the ability for Galaxy's ObjectStore to use OpenStack's SWIFT object store as the backend data storage --- lib/galaxy/config.py | 12 ++++++--- lib/galaxy/objectstore/__init__.py | 39 ++++++++++++++++++++++++------ universe_wsgi.ini.sample | 22 +++++++++++------ 3 files changed, 53 insertions(+), 20 deletions(-) diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index e0893cb2c43..c55516ec843 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -169,10 +169,14 @@ class Configuration( object ): if self.nginx_upload_store: self.nginx_upload_store = os.path.abspath( self.nginx_upload_store ) self.object_store = kwargs.get( 'object_store', 'disk' ) - self.aws_access_key = kwargs.get( 'aws_access_key', None ) - self.aws_secret_key = kwargs.get( 'aws_secret_key', None ) - self.s3_bucket = kwargs.get( 's3_bucket', None) - self.use_reduced_redundancy = kwargs.get( 'use_reduced_redundancy', False ) + self.os_access_key = kwargs.get( 'os_access_key', None ) + self.os_secret_key = kwargs.get( 'os_secret_key', None ) + self.os_bucket_name = kwargs.get( 'os_bucket_name', None ) + self.os_host = kwargs.get( 'os_host', None ) + self.os_port = kwargs.get( 'os_port', None ) + self.os_is_secure = string_as_bool( kwargs.get( 'os_is_secure', True ) ) + self.os_conn_path = kwargs.get( 'os_conn_path', '/' ) + self.os_use_reduced_redundancy = kwargs.get( 'os_use_reduced_redundancy', False ) self.object_store_cache_size = float(kwargs.get( 'object_store_cache_size', -1 )) self.distributed_object_store_config_file = kwargs.get( 'distributed_object_store_config_file', None ) # Parse global_conf and save the parser diff --git a/lib/galaxy/objectstore/__init__.py b/lib/galaxy/objectstore/__init__.py index 6af4414ac8c..ab2fe8f8aae 100644 --- a/lib/galaxy/objectstore/__init__.py +++ b/lib/galaxy/objectstore/__init__.py @@ -25,6 +25,7 @@ from sqlalchemy.orm import object_session if sys.version_info >= (2, 6): import multiprocessing from galaxy.objectstore.s3_multipart_upload import multipart_upload + import boto from boto.s3.key import Key from boto.s3.connection import S3Connection from boto.exception import S3ResponseError @@ -377,9 +378,9 @@ class S3ObjectStore(ObjectStore): super(S3ObjectStore, self).__init__() self.config = config self.staging_path = self.config.file_path - self.s3_conn = S3Connection() - self.bucket = self._get_bucket(self.config.s3_bucket) - self.use_rr = self.config.use_reduced_redundancy + self.s3_conn = get_OS_connection(self.config) + self.bucket = self._get_bucket(self.config.os_bucket_name) + self.use_rr = self.config.os_use_reduced_redundancy self.cache_size = self.config.object_store_cache_size self.transfer_progress = 0 # Clean cache only if value is set in universe_wsgi.ini @@ -468,7 +469,7 @@ class S3ObjectStore(ObjectStore): for i in range(5): try: bucket = self.s3_conn.get_bucket(bucket_name) - log.debug("Using S3 object store; got bucket '%s'" % bucket.name) + log.debug("Using cloud object store with bucket '%s'" % bucket.name) return bucket except S3ResponseError: log.debug("Could not get bucket '%s', attempt %s/5" % (bucket_name, i+1)) @@ -843,7 +844,6 @@ class S3ObjectStore(ObjectStore): def get_store_usage_percent(self): return 0.0 - class DistributedObjectStore(ObjectStore): """ ObjectStore that defers to a list of backends, for getting objects the @@ -1009,14 +1009,14 @@ def build_object_store_from_config(config): store = config.object_store if store == 'disk': return DiskObjectStore(config=config) - elif store == 's3': - os.environ['AWS_ACCESS_KEY_ID'] = config.aws_access_key - os.environ['AWS_SECRET_ACCESS_KEY'] = config.aws_secret_key + elif store == 's3' or store == 'swift': return S3ObjectStore(config=config) elif store == 'distributed': return DistributedObjectStore(config=config) elif store == 'hierarchical': return HierarchicalObjectStore() + else: + log.error("Unrecognized object store definition: {0}".format(store)) def convert_bytes(bytes): """ A helper function used for pretty printing disk usage """ @@ -1039,3 +1039,26 @@ def convert_bytes(bytes): else: size = '%.2fb' % bytes return size + +def get_OS_connection(config): + """ + Get a connection object for a cloud Object Store specified in the config. + Currently, this is a ``boto`` connection object. + """ + log.debug("Getting a connection object for '{0}' object store".format(config.object_store)) + a_key = config.os_access_key + s_key = config.os_secret_key + if config.object_store == 's3': + return S3Connection(a_key, s_key) + else: + # Establish the connection now + calling_format = boto.s3.connection.OrdinaryCallingFormat() + s3_conn = boto.connect_s3(aws_access_key_id=a_key, + aws_secret_access_key=s_key, + is_secure=config.os_is_secure, + host=config.os_host, + port=int(config.os_port), + calling_format=calling_format, + path=config.os_conn_path) + return s3_conn + diff --git a/universe_wsgi.ini.sample b/universe_wsgi.ini.sample index abb17224c1e..4e627e5ac4c 100644 --- a/universe_wsgi.ini.sample +++ b/universe_wsgi.ini.sample @@ -481,16 +481,22 @@ use_interactive = True # -- Beta features -# Object store mode (valid options are: disk, s3, distributed, hierarchical) +# Object store mode (valid options are: disk, s3, swift, distributed, hierarchical) #object_store = disk -#aws_access_key = -#aws_secret_key = -#s3_bucket = -#use_reduced_redundancy = True - +#os_access_key = +#os_secret_key = +#os_bucket_name = +# If using 'swift' object store, you must specify the following connection properties +#os_host = swift.rc.nectar.org.au +#os_port = 8888 +#os_is_secure = False +#os_conn_path = / +# Reduced redundancy can be used only with the 's3' object store +#os_use_reduced_redundancy = False # Size (in GB) that the cache used by object store should be limited to. -# If the value is not specified, the cache size will be limited only by the file -# system size. +# If the value is not specified, the cache size will be limited only by the +# file system size. The file system location of the cache is considered the +# configuration of the ``file_path`` directive defined above. #object_store_cache_size = 100 # Configuration file for the distributed object store, if object_store = From cdd15be456422584f230e247d96e0ce4bff239c3 Mon Sep 17 00:00:00 2001 From: Greg Von Kuster Date: Wed, 15 Aug 2012 11:55:17 -0400 Subject: [PATCH 10/87] Fix for setting tooll dependency metadata where at least one tool in the repository does not include a tag set. --- lib/galaxy/util/shed_util.py | 2 +- lib/galaxy/webapps/community/config.py | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/util/shed_util.py b/lib/galaxy/util/shed_util.py index 23512eed871..cf70499c2e7 100644 --- a/lib/galaxy/util/shed_util.py +++ b/lib/galaxy/util/shed_util.py @@ -533,7 +533,7 @@ def can_generate_tool_dependency_metadata( root, metadata_dict ): if req_name==tool_dependency_name and req_version==tool_dependency_version and req_type==tool_dependency_type: can_generate_dependency_metadata = True break - if not can_generate_dependency_metadata: + if requirements and not can_generate_dependency_metadata: # We've discovered at least 1 combination of name, version and type that is not defined in the # tag for any tool in the repository. break diff --git a/lib/galaxy/webapps/community/config.py b/lib/galaxy/webapps/community/config.py index 3f07abaa336..c7a768eff9d 100644 --- a/lib/galaxy/webapps/community/config.py +++ b/lib/galaxy/webapps/community/config.py @@ -87,6 +87,7 @@ class Configuration( object ): self.server_name = '' self.job_manager = '' self.default_job_handlers = [] + self.default_cluster_job_runner = 'local:///' self.job_handlers = [] self.tool_handlers = [] self.tool_runners = [] From c0898c6882d87efb01eac894476463c7cc0dc405 Mon Sep 17 00:00:00 2001 From: Greg Von Kuster Date: Wed, 15 Aug 2012 14:50:38 -0400 Subject: [PATCH 11/87] Improve error message handling when setting metadata on tool shed repositories. Display the repository owner in the table grid when displaying invalid tools in the tool shed. --- .../webapps/community/controllers/admin.py | 11 +++- .../webapps/community/controllers/common.py | 58 ++++++++++++------- .../community/controllers/repository.py | 22 +++++-- .../repository/browse_invalid_tools.mako | 4 +- 4 files changed, 65 insertions(+), 30 deletions(-) diff --git a/lib/galaxy/webapps/community/controllers/admin.py b/lib/galaxy/webapps/community/controllers/admin.py index ad6d018635f..37028ce7eb1 100644 --- a/lib/galaxy/webapps/community/controllers/admin.py +++ b/lib/galaxy/webapps/community/controllers/admin.py @@ -696,9 +696,14 @@ class AdminController( BaseUIController, Admin ): owner = repository_name_owner_list[ 1 ] repository = get_repository_by_name_and_owner( trans, name, owner ) try: - reset_all_metadata_on_repository( trans, trans.security.encode_id( repository.id ) ) - log.debug( "Successfully reset metadata on repository %s" % repository.name ) - successful_count += 1 + invalid_file_tups = reset_all_metadata_on_repository( trans, trans.security.encode_id( repository.id ) ) + if invalid_file_tups: + message = generate_message_for_invalid_tools( invalid_file_tups, repository, None, as_html=False ) + log.debug( message ) + unsuccessful_count += 1 + else: + log.debug( "Successfully reset metadata on repository %s" % repository.name ) + successful_count += 1 except Exception, e: log.debug( "Error attempting to reset metadata on repository '%s': %s" % ( repository.name, str( e ) ) ) unsuccessful_count += 1 diff --git a/lib/galaxy/webapps/community/controllers/common.py b/lib/galaxy/webapps/community/controllers/common.py index cb7d832a655..a102439a631 100644 --- a/lib/galaxy/webapps/community/controllers/common.py +++ b/lib/galaxy/webapps/community/controllers/common.py @@ -277,6 +277,41 @@ def generate_clone_url( trans, repository_id ): return '%s://%s%s/repos/%s/%s' % ( protocol, username, base, repository.user.username, repository.name ) else: return '%s/repos/%s/%s' % ( base_url, repository.user.username, repository.name ) +def generate_message_for_invalid_tools( invalid_file_tups, repository, metadata_dict, as_html=True ): + if as_html: + new_line = '
' + bold_start = '' + bold_end = '' + else: + new_line = '\n' + bold_start = '' + bold_end = '' + message = '' + if metadata_dict: + message += "Metadata was defined for some items in revision '%s'. " % str( repository.tip ) + message += "Correct the following problems if necessary and reset metadata.%s" % new_line + else: + message += "Metadata cannot be defined for revision '%s' so this revision cannot be automatically " % str( repository.tip ) + message += "installed into a local Galaxy instance. Correct the following problems and reset metadata.%s" % new_line + for itc_tup in invalid_file_tups: + tool_file, exception_msg = itc_tup + if exception_msg.find( 'No such file or directory' ) >= 0: + exception_items = exception_msg.split() + missing_file_items = exception_items[ 7 ].split( '/' ) + missing_file = missing_file_items[ -1 ].rstrip( '\'' ) + if missing_file.endswith( '.loc' ): + sample_ext = '%s.sample' % missing_file + else: + sample_ext = missing_file + correction_msg = "This file refers to a missing file %s%s%s. " % ( bold_start, str( missing_file ), bold_end ) + correction_msg += "Upload a file named %s%s%s to the repository to correct this error." % ( bold_start, sample_ext, bold_end ) + else: + if as_html: + correction_msg = exception_msg + else: + correction_msg = exception_msg.replace( '
', new_line ).replace( '', bold_start ).replace( '', bold_end ) + message += "%s%s%s - %s%s" % ( bold_start, tool_file, bold_end, correction_msg, new_line ) + return message def generate_tool_guid( trans, repository, tool ): """ Generate a guid for the received tool. The form of the guid is @@ -854,6 +889,7 @@ def reset_all_metadata_on_repository( trans, id, **kwd ): clean_repository_metadata( trans, id, changeset_revisions ) # Set tool version information for all downloadable changeset revisions. Get the list of changeset revisions from the changelog. reset_all_tool_versions( trans, id, repo ) + return invalid_file_tups def set_repository_metadata( trans, repository, content_alert_str='', **kwd ): """ Set metadata using the repository's current disk files, returning specific error messages (if any) to alert the repository owner that the changeset @@ -931,27 +967,7 @@ def set_repository_metadata( trans, repository, content_alert_str='', **kwd ): message += "be defined so this revision cannot be automatically installed into a local Galaxy instance." status = "error" if invalid_file_tups: - if metadata_dict: - message += "Metadata was defined for some items in revision '%s'. " % str( repository.tip ) - message += "Correct the following problems if necessary and reset metadata.
" - else: - message += "Metadata cannot be defined for revision '%s' so this revision cannot be automatically " % str( repository.tip ) - message += "installed into a local Galaxy instance. Correct the following problems and reset metadata.
" - for itc_tup in invalid_file_tups: - tool_file, exception_msg = itc_tup - if exception_msg.find( 'No such file or directory' ) >= 0: - exception_items = exception_msg.split() - missing_file_items = exception_items[ 7 ].split( '/' ) - missing_file = missing_file_items[ -1 ].rstrip( '\'' ) - if missing_file.endswith( '.loc' ): - sample_ext = '%s.sample' % missing_file - else: - sample_ext = missing_file - correction_msg = "This file refers to a missing file %s. " % str( missing_file ) - correction_msg += "Upload a file named %s to the repository to correct this error." % sample_ext - else: - correction_msg = exception_msg - message += "%s - %s
" % ( tool_file, correction_msg ) + message = generate_message_for_invalid_tools( invalid_file_tups, repository, metadata_dict ) status = 'error' return message, status def set_repository_metadata_due_to_new_tip( trans, repository, content_alert_str=None, **kwd ): diff --git a/lib/galaxy/webapps/community/controllers/repository.py b/lib/galaxy/webapps/community/controllers/repository.py index bfaa4553a88..5526d5427a3 100644 --- a/lib/galaxy/webapps/community/controllers/repository.py +++ b/lib/galaxy/webapps/community/controllers/repository.py @@ -458,7 +458,10 @@ class RepositoryController( BaseUIController, ItemRatings ): metadata = downloadable_revision.metadata invalid_tools = metadata.get( 'invalid_tools', [] ) for invalid_tool_config in invalid_tools: - invalid_tools_dict[ invalid_tool_config ] = ( repository.id, repository.name, downloadable_revision.changeset_revision ) + invalid_tools_dict[ invalid_tool_config ] = ( repository.id, + repository.name, + repository.user.username, + downloadable_revision.changeset_revision ) else: for repository in trans.sa_session.query( trans.model.Repository ) \ .filter( and_( trans.model.Repository.table.c.deleted == False, @@ -468,7 +471,10 @@ class RepositoryController( BaseUIController, ItemRatings ): metadata = downloadable_revision.metadata invalid_tools = metadata.get( 'invalid_tools', [] ) for invalid_tool_config in invalid_tools: - invalid_tools_dict[ invalid_tool_config ] = ( repository.id, repository.name, downloadable_revision.changeset_revision ) + invalid_tools_dict[ invalid_tool_config ] = ( repository.id, + repository.name, + repository.user.username, + downloadable_revision.changeset_revision ) return trans.fill_template( '/webapps/community/repository/browse_invalid_tools.mako', cntrller=cntrller, invalid_tools_dict=invalid_tools_dict, @@ -1373,6 +1379,7 @@ class RepositoryController( BaseUIController, ItemRatings ): return trans.response.send_redirect( url ) @web.expose def load_invalid_tool( self, trans, repository_id, tool_config, changeset_revision, **kwd ): + # FIXME: loading an invalid tool should display an appropriate message as to why the tool is invalid. This worked until recently. params = util.Params( kwd ) message = util.restore_text( params.get( 'message', '' ) ) status = params.get( 'status', 'error' ) @@ -1752,9 +1759,14 @@ class RepositoryController( BaseUIController, ItemRatings ): status=status ) @web.expose def reset_all_metadata( self, trans, id, **kwd ): - reset_all_metadata_on_repository( trans, id, **kwd ) - message = "All repository metadata has been reset." - status = 'done' + invalid_file_tups = reset_all_metadata_on_repository( trans, id, **kwd ) + if invalid_file_tups: + repository = get_repository( trans, id ) + message = generate_message_for_invalid_tools( invalid_file_tups, repository, None ) + status = 'error' + else: + message = "All repository metadata has been reset." + status = 'done' return trans.response.send_redirect( web.url_for( controller='repository', action='manage_repository', id=id, diff --git a/templates/webapps/community/repository/browse_invalid_tools.mako b/templates/webapps/community/repository/browse_invalid_tools.mako index b3a8ef29784..5e7d293df7c 100644 --- a/templates/webapps/community/repository/browse_invalid_tools.mako +++ b/templates/webapps/community/repository/browse_invalid_tools.mako @@ -13,10 +13,11 @@ Tool config Repository name + Repository owner Changeset revision %for invalid_tool_config, repository_tup in invalid_tools_dict.items(): - <% repository_id, repository_name, changeset_revision = repository_tup %> + <% repository_id, repository_name, repository_owner, changeset_revision = repository_tup %> @@ -24,6 +25,7 @@ ${repository_name} + ${repository_owner} ${changeset_revision} %endfor From 5dda060e33d047b23256c1dcf821dc5ca4d65bee Mon Sep 17 00:00:00 2001 From: Jeremy Goecks Date: Wed, 15 Aug 2012 17:49:39 -0400 Subject: [PATCH 12/87] Rewrite sampling code for BBI data provider to handle (a) boundary cases during base-level resolution and (b) remainder of region not sampled during first pass. --- .../visualization/tracks/data_providers.py | 102 ++++++++++-------- 1 file changed, 58 insertions(+), 44 deletions(-) diff --git a/lib/galaxy/visualization/tracks/data_providers.py b/lib/galaxy/visualization/tracks/data_providers.py index 19581e544c8..50f969239c4 100644 --- a/lib/galaxy/visualization/tracks/data_providers.py +++ b/lib/galaxy/visualization/tracks/data_providers.py @@ -947,55 +947,69 @@ class BBIDataProvider( TracksDataProvider ): return dict( data=dict( min=summary.min_val[0], max=summary.max_val[0], mean=mean, sd=sd ) ) - # The following seems not to work very well, for example it will only return one - # data point if the tile is 1280px wide. Not sure what the intent is. + # Sample from region using approximately this many samples. + N = 1000 - # The first zoom level for BBI files is 640. If too much is requested, it will look at each block instead - # of summaries. The calculation done is: zoom <> (end-start)/num_points/2. - # Thus, the optimal number of points is (end-start)/num_points/2 = 640 - # num_points = (end-start) / 1280 - #num_points = (end-start) / 1280 - #if num_points < 1: - # num_points = end - start - #else: - # num_points = min(num_points, 500) + def summarize_region( bbi, chrom, start, end, num_points ): + ''' + Returns results from summarizing a region using num_points. + NOTE: num_points cannot be greater than end - start or BBI + will return None for all positions.s + ''' + result = [] - # For now, we'll do 1000 data points by default. However, the summaries - # don't seem to work when a summary pixel corresponds to less than one - # datapoint, so we prevent that. - - # FIXME: need to choose the number of points to maximize coverage of the area. - # It appears that BBI calculates points using intervals of - # floor( num_points / end - start ) - # In some cases, this prevents sampling near the end of the interval, - # especially when (a) the total interval is small ( < 20-30Kb) and (b) the - # computed interval size has a large fraction, e.g. 14.7 or 35.8 - num_points = min( 1000, end - start ) - - # HACK to address the FIXME above; should generalize. - if end - start <= 2000: - num_points = end - start - - summary = bbi.summarize( chrom, start, end, num_points ) - f.close() - - result = [] - - if summary: - #mean = summary.sum_data / summary.valid_count + # Get summary; this samples at intervals of length + # (end - start)/num_points -- i.e. drops any fractional component + # of interval length. + summary = bbi.summarize( chrom, start, end, num_points ) + if summary: + #mean = summary.sum_data / summary.valid_count + + ## Standard deviation by bin, not yet used + ## var = summary.sum_squares - mean + ## var /= minimum( valid_count - 1, 1 ) + ## sd = sqrt( var ) - ## Standard deviation by bin, not yet used - ## var = summary.sum_squares - mean - ## var /= minimum( valid_count - 1, 1 ) - ## sd = sqrt( var ) - - pos = start - step_size = (end - start) / num_points + pos = start + step_size = (end - start) / num_points - for i in range( num_points ): - result.append( (pos, float_nan( summary.sum_data[i] / summary.valid_count[i] ) ) ) - pos += step_size + for i in range( num_points ): + result.append( (pos, float_nan( summary.sum_data[i] / summary.valid_count[i] ) ) ) + pos += step_size + + return result + + # Approach is different depending on region size. + if end - start < N: + # Get values for individual bases in region, including start and end. + # To do this, need to increase end to next base and request number of points. + num_points = end - start + 1 + end += 1 + + result = summarize_region( bbi, chrom, start, end, num_points ) + else: + # + # The goal is to sample the region between start and end uniformly + # using N data points. The challenge is that the size of sampled + # intervals rarely is full bases, so sampling using N points will + # leave the end of the region unsampled. To recitify this, samples + # beyond N are taken at the end of the interval. + # + + # Do initial summary. + num_points = N + result = summarize_region( bbi, chrom, start, end, num_points ) + + # Do summary of remaining part of region. + step_size = ( end - start ) / num_points + new_start = start + step_size * num_points + new_num_points = min( ( end - new_start ) / step_size, end - start ) + if new_num_points is not 0: + result.extend( summarize_region( bbi, chrom, new_start, end, new_num_points ) ) + #TODO: progressively reduce step_size to generate more datapoints. + # Cleanup and return. + f.close() return { 'data': result } class BigBedDataProvider( BBIDataProvider ): From 54b1f8097aff613478d674f39a8196ea5cf3bf5b Mon Sep 17 00:00:00 2001 From: Jeremy Goecks Date: Wed, 15 Aug 2012 18:09:48 -0400 Subject: [PATCH 13/87] Cleanup for previous commit, 565476ce4f03, mainly to further comment and simplify code and avoid going to index multiple times. --- .../visualization/tracks/data_providers.py | 31 ++++++++++--------- 1 file changed, 17 insertions(+), 14 deletions(-) diff --git a/lib/galaxy/visualization/tracks/data_providers.py b/lib/galaxy/visualization/tracks/data_providers.py index 50f969239c4..eddc8f57169 100644 --- a/lib/galaxy/visualization/tracks/data_providers.py +++ b/lib/galaxy/visualization/tracks/data_providers.py @@ -985,28 +985,31 @@ class BBIDataProvider( TracksDataProvider ): # To do this, need to increase end to next base and request number of points. num_points = end - start + 1 end += 1 - - result = summarize_region( bbi, chrom, start, end, num_points ) else: # # The goal is to sample the region between start and end uniformly - # using N data points. The challenge is that the size of sampled + # using ~N data points. The challenge is that the size of sampled # intervals rarely is full bases, so sampling using N points will - # leave the end of the region unsampled. To recitify this, samples - # beyond N are taken at the end of the interval. + # leave the end of the region unsampled due to remainders for each + # interval. To recitify this, a new N is calculated based on the + # step size that covers as much of the region as possible. + # + # However, this still leaves some of the region unsampled. This + # could be addressed by repeatedly sampling remainder using a + # smaller and smaller step_size, but that would require iteratively + # going to BBI, which could be time consuming. # - # Do initial summary. + # Start with N samples. num_points = N - result = summarize_region( bbi, chrom, start, end, num_points ) - - # Do summary of remaining part of region. step_size = ( end - start ) / num_points - new_start = start + step_size * num_points - new_num_points = min( ( end - new_start ) / step_size, end - start ) - if new_num_points is not 0: - result.extend( summarize_region( bbi, chrom, new_start, end, new_num_points ) ) - #TODO: progressively reduce step_size to generate more datapoints. + # Add additional points to sample in the remainder not covered by + # the initial N samples. + remainder_start = start + step_size * num_points + additional_points = ( end - remainder_start ) / step_size + num_points += additional_points + + result = summarize_region( bbi, chrom, start, end, num_points ) # Cleanup and return. f.close() From 9a47e09e21e0679ac60762756cecd062a4370736 Mon Sep 17 00:00:00 2001 From: Enis Afgan Date: Thu, 16 Aug 2012 09:50:00 +1000 Subject: [PATCH 14/87] Handle AWS-specific config options for backward compatibility --- lib/galaxy/config.py | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index c55516ec843..aaf2890b4d1 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -169,14 +169,21 @@ class Configuration( object ): if self.nginx_upload_store: self.nginx_upload_store = os.path.abspath( self.nginx_upload_store ) self.object_store = kwargs.get( 'object_store', 'disk' ) - self.os_access_key = kwargs.get( 'os_access_key', None ) - self.os_secret_key = kwargs.get( 'os_secret_key', None ) - self.os_bucket_name = kwargs.get( 'os_bucket_name', None ) + # Handle AWS-specific config options for backward compatibility + if kwargs.get( 'aws_access_key', None) is not None: + self.os_access_key= kwargs.get( 'aws_access_key', None ) + self.os_secret_key= kwargs.get( 'aws_secret_key', None ) + self.os_bucket_name= kwargs.get( 's3_bucket', None ) + self.os_use_reduced_redundancy = kwargs.get( 'use_reduced_redundancy', False ) + else: + self.os_access_key = kwargs.get( 'os_access_key', None ) + self.os_secret_key = kwargs.get( 'os_secret_key', None ) + self.os_bucket_name = kwargs.get( 'os_bucket_name', None ) + self.os_use_reduced_redundancy = kwargs.get( 'os_use_reduced_redundancy', False ) self.os_host = kwargs.get( 'os_host', None ) self.os_port = kwargs.get( 'os_port', None ) self.os_is_secure = string_as_bool( kwargs.get( 'os_is_secure', True ) ) self.os_conn_path = kwargs.get( 'os_conn_path', '/' ) - self.os_use_reduced_redundancy = kwargs.get( 'os_use_reduced_redundancy', False ) self.object_store_cache_size = float(kwargs.get( 'object_store_cache_size', -1 )) self.distributed_object_store_config_file = kwargs.get( 'distributed_object_store_config_file', None ) # Parse global_conf and save the parser From 65179d6c0b55a1cc0f5037a10830485f99411e4f Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 16 Aug 2012 10:58:01 -0500 Subject: [PATCH 15/87] Dynamic job runner bug fixes. --- lib/galaxy/jobs/handler.py | 2 +- lib/galaxy/jobs/mapper.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index 9fe115277e0..cd0fd94b1e8 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -360,7 +360,7 @@ class DefaultJobDispatcher( object ): def __init__( self, app ): self.app = app self.job_runners = {} - start_job_runners = ["local", "lwr", "dynamic"] + start_job_runners = ["local", "lwr"] if app.config.start_job_runners is not None: start_job_runners.extend( [ x.strip() for x in util.listify( app.config.start_job_runners ) ] ) if app.config.use_tasked_jobs: diff --git a/lib/galaxy/jobs/mapper.py b/lib/galaxy/jobs/mapper.py index 2c0d003d36d..efa5cbcffa5 100644 --- a/lib/galaxy/jobs/mapper.py +++ b/lib/galaxy/jobs/mapper.py @@ -111,7 +111,7 @@ class JobRunnerMapper( object ): expand_function = self.__get_expand_function( expand_function_name ) return self.__invoke_expand_function( expand_function ) else: - raise Exception( "Unhandled dynamic job runner type specified - %s" % calculation_type ) + raise Exception( "Unhandled dynamic job runner type specified - %s" % expand_type ) def __cache_job_runner_url( self, params ): raw_job_runner_url = self.job_wrapper.tool.get_job_runner_url( params ) From 48ad38ff3cf9e0d1526f9d88ab8ef55d866b9e83 Mon Sep 17 00:00:00 2001 From: Dave Bouvier Date: Thu, 16 Aug 2012 15:43:02 -0400 Subject: [PATCH 16/87] Improved error handling for failed indexing jobs. --- lib/galaxy/tools/genome_index/__init__.py | 21 ++++++++++++++----- lib/galaxy/tools/genome_index/index_genome.py | 13 +++++++----- 2 files changed, 24 insertions(+), 10 deletions(-) diff --git a/lib/galaxy/tools/genome_index/__init__.py b/lib/galaxy/tools/genome_index/__init__.py index 6318325f226..0eade5c4d8c 100644 --- a/lib/galaxy/tools/genome_index/__init__.py +++ b/lib/galaxy/tools/genome_index/__init__.py @@ -18,13 +18,16 @@ def load_genome_index_tools( toolbox ): - $__GENOME_INDEX_COMMAND__ $output_file $output_file.files_path $__app__.config.rsync_url "$__app__.config.tool_data_path" + $__GENOME_INDEX_COMMAND__ $output_file $output_file.files_path "$__app__.config.rsync_url" "$__app__.config.tool_data_path" + + + """ @@ -64,6 +67,18 @@ class GenomeIndexToolWrapper( object ): if gitd: + fp = open( gitd.dataset.get_file_name(), 'r' ) + deferred = sa_session.query( model.DeferredJob ).filter_by( id=gitd.deferred_job_id ).first() + try: + logloc = json.load( fp ) + except ValueError: + deferred.state = app.model.DeferredJob.states.ERROR + sa_session.add( deferred ) + sa_session.flush() + log.debug( 'Indexing job failed, setting deferred job state to error.' ) + return False + finally: + fp.close() destination = None tdtman = ToolDataTableManager( app.config.tool_data_path ) xmltree = tdtman.load_from_config_file( app.config.tool_data_table_config_path, app.config.tool_data_path ) @@ -72,16 +87,12 @@ class GenomeIndexToolWrapper( object ): location = node.findall('file')[0].get('path') self.locations[table] = os.path.abspath( location ) locbase = os.path.abspath( os.path.split( self.locations['all_fasta'] )[0] ) - deferred = sa_session.query( model.DeferredJob ).filter_by( id=gitd.deferred_job_id ).first() params = deferred.params dbkey = params[ 'dbkey' ] basepath = os.path.join( os.path.abspath( app.config.genome_data_path ), dbkey ) intname = params[ 'intname' ] indexer = gitd.indexer workingdir = os.path.abspath( gitd.dataset.extra_files_path ) - fp = open( gitd.dataset.get_file_name(), 'r' ) - logloc = json.load( fp ) - fp.close() location = [] indexdata = gitd.dataset.extra_files_path if indexer == '2bit': diff --git a/lib/galaxy/tools/genome_index/index_genome.py b/lib/galaxy/tools/genome_index/index_genome.py index 26aa918fe9f..553733eb877 100644 --- a/lib/galaxy/tools/genome_index/index_genome.py +++ b/lib/galaxy/tools/genome_index/index_genome.py @@ -40,15 +40,16 @@ class ManagedIndexer(): self.genome = os.path.splitext( self.fafile )[0] with WithChDir( self.basedir ): if indexer not in self.indexers: - raise KeyError, 'The requested indexing function does not exist' + sys.stderr.write( 'The requested indexing function does not exist' ) + exit(127) else: with WithChDir( self.workingdir ): self._log( 'Running indexer %s.' % indexer ) result = getattr( self, self.indexers[ indexer ] )() if result in [ None, False ]: - self._log( 'Error running indexer %s, %s' % ( indexer, result ) ) + sys.stderr.write( 'Error running indexer %s, %s' % ( indexer, result ) ) self._flush_files() - return True + exit(1) else: self._log( self.locations ) self._log( 'Indexer %s completed successfully.' % indexer ) @@ -309,5 +310,7 @@ if __name__ == "__main__": # Create archive. idxobj = ManagedIndexer( outfile, infile, working_dir, rsync_url, tooldata ) - idxobj.run_indexer( indexer ) - \ No newline at end of file + returncode = idxobj.run_indexer( indexer ) + if not returncode: + exit(1) + exit(0) \ No newline at end of file From 9fabf5dfbe07844c0e6fc2a20796f9b36063427c Mon Sep 17 00:00:00 2001 From: Jeremy Goecks Date: Thu, 16 Aug 2012 16:32:00 -0400 Subject: [PATCH 17/87] Add feature/attribute name indexing framework to converters. Provide full text indexing for GFF attributes. --- datatypes_conf.xml.sample | 2 + lib/galaxy/datatypes/converters/gff_to_fli.py | 53 +++++++++++++++++++ .../converters/gff_to_fli_converter.xml | 13 +++++ lib/galaxy/datatypes/tabular.py | 7 +++ .../visualization/tracks/data_providers.py | 47 +++++++++++++++- lib/galaxy/web/controllers/tracks.py | 14 +++++ tools/filters/gff/sort_gtf.py | 1 + 7 files changed, 136 insertions(+), 1 deletion(-) create mode 100644 lib/galaxy/datatypes/converters/gff_to_fli.py create mode 100644 lib/galaxy/datatypes/converters/gff_to_fli_converter.xml diff --git a/datatypes_conf.xml.sample b/datatypes_conf.xml.sample index cb2f1b3af88..099addfa73a 100644 --- a/datatypes_conf.xml.sample +++ b/datatypes_conf.xml.sample @@ -4,6 +4,7 @@ + @@ -79,6 +80,7 @@ + diff --git a/lib/galaxy/datatypes/converters/gff_to_fli.py b/lib/galaxy/datatypes/converters/gff_to_fli.py new file mode 100644 index 00000000000..cfa056e573f --- /dev/null +++ b/lib/galaxy/datatypes/converters/gff_to_fli.py @@ -0,0 +1,53 @@ +''' +Creates a feature location index for a given GFF file. +''' + +import sys +from galaxy import eggs +from galaxy.datatypes.util.gff_util import read_unordered_gtf, convert_gff_coords_to_bed + +# Process arguments. +in_fname = sys.argv[1] +out_fname = sys.argv[2] + +# Create dict of name-location pairings. +name_loc_dict = {} +for feature in read_unordered_gtf( open( in_fname, 'r' ) ): + for name in feature.attributes: + val = feature.attributes[ name ] + try: + float( val ) + continue + except: + convert_gff_coords_to_bed( feature ) + # Value is not a number, so it can be indexed. + if val not in name_loc_dict: + # Value is not in dictionary. + name_loc_dict[ val ] = { + 'contig': feature.chrom, + 'start': feature.start, + 'end': feature.end + } + else: + # Value already in dictionary, so update dictionary. + loc = name_loc_dict[ val ] + if feature.start < loc[ 'start' ]: + loc[ 'start' ] = feature.start + if feature.end > loc[ 'end' ]: + loc[ 'end' ] = feature.end + +# Print name, loc in sorted order. +out = open( out_fname, 'w' ) +max_len = 0 +entries = [] +for name in sorted( name_loc_dict.iterkeys() ): + loc = name_loc_dict[ name ] + entry = '%s\t%s' % ( name, '%s:%i-%i' % ( loc[ 'contig' ], loc[ 'start' ], loc[ 'end' ] ) ) + if len( entry ) > max_len: + max_len = len( entry ) + entries.append( entry ) + +out.write( str( max_len + 1 ).ljust( max_len ) + '\n' ) +for entry in entries: + out.write( entry.ljust( max_len ) + '\n' ) +out.close() \ No newline at end of file diff --git a/lib/galaxy/datatypes/converters/gff_to_fli_converter.xml b/lib/galaxy/datatypes/converters/gff_to_fli_converter.xml new file mode 100644 index 00000000000..22f739438a0 --- /dev/null +++ b/lib/galaxy/datatypes/converters/gff_to_fli_converter.xml @@ -0,0 +1,13 @@ + + + + gff_to_fli.py $input1 $output1 + + + + + + + + + diff --git a/lib/galaxy/datatypes/tabular.py b/lib/galaxy/datatypes/tabular.py index 0f5c5b5d348..ce3b043b658 100644 --- a/lib/galaxy/datatypes/tabular.py +++ b/lib/galaxy/datatypes/tabular.py @@ -638,3 +638,10 @@ class Eland( Tabular ): dataset.metadata.reads = reads.keys() +class FeatureLocationIndex( Tabular ): + """ + An index that stores feature locations in tabular format. + """ + file_ext='fli' + MetadataElement( name="columns", default=2, desc="Number of columns", readonly=True, visible=False ) + MetadataElement( name="column_types", default=['str', 'str'], param=metadata.ColumnTypesParameter, desc="Column types", readonly=True, visible=False, no_value=[] ) \ No newline at end of file diff --git a/lib/galaxy/visualization/tracks/data_providers.py b/lib/galaxy/visualization/tracks/data_providers.py index eddc8f57169..0999f9e3932 100644 --- a/lib/galaxy/visualization/tracks/data_providers.py +++ b/lib/galaxy/visualization/tracks/data_providers.py @@ -2,7 +2,7 @@ Data providers for tracks visualizations. """ -import sys +import os, sys from math import ceil, log import pkg_resources pkg_resources.require( "bx-python" ) @@ -59,6 +59,51 @@ def _convert_between_ucsc_and_ensemble_naming( chrom ): def _chrom_naming_matches( chrom1, chrom2 ): return ( chrom1.startswith( 'chr' ) and chrom2.startswith( 'chr' ) ) or ( not chrom1.startswith( 'chr' ) and not chrom2.startswith( 'chr' ) ) + +class FeatureLocationIndexDataProvider( object ): + ''' + + ''' + + def __init__( self, converted_dataset ): + self.converted_dataset = converted_dataset + + def get_data( self, query ): + # Init. + textloc_file = open( self.converted_dataset.file_name, 'r' ) + line_len = int( textloc_file.readline() ) + file_len = os.path.getsize( self.converted_dataset.file_name ) + + # Find query in file using binary search. + low = 0 + high = file_len / line_len + while low < high: + mid = ( low + high ) // 2 + position = mid * line_len + textloc_file.seek( position ) + + # Compare line with query and update low, high. + line = textloc_file.readline() + print '--', mid, line + if line < query: + low = mid + 1 + else: + high = mid + + position = low * line_len + + # At right point in file, generate hits. + result = [ ] + while True: + line = textloc_file.readline() + if not line.startswith( query ): + break + if line[ -1: ] == '\n': + line = line[ :-1 ] + result.append( line.split() ) + + textloc_file.close() + return result class TracksDataProvider( object ): """ Base class for tracks data providers. """ diff --git a/lib/galaxy/web/controllers/tracks.py b/lib/galaxy/web/controllers/tracks.py index 10ae947d9a0..d468b0dfb35 100644 --- a/lib/galaxy/web/controllers/tracks.py +++ b/lib/galaxy/web/controllers/tracks.py @@ -345,6 +345,20 @@ class TracksController( BaseUIController, UsesVisualizationMixin, UsesHistoryDat # Have data if we get here return { "status": messages.DATA, "valid_chroms": valid_chroms } + + @web.json + def feature_loc( self, trans, hda_ldda, dataset_id, query ): + """ + Returns features, locations in dataset that match query. Format is a + list of features; each feature is a list itself: [name, location] + """ + dataset = self.get_hda_or_ldda( trans, hda_ldda, dataset_id ) + converted_dataset = dataset.get_converted_dataset( trans, "fli" ) + data_provider = FeatureLocationIndexDataProvider( converted_dataset=converted_dataset ) + if data_provider: + return data_provider.get_data( query ) + else: + return 'None' @web.json def data( self, trans, hda_ldda, dataset_id, chrom, low, high, start_val=0, max_vals=None, **kwargs ): diff --git a/tools/filters/gff/sort_gtf.py b/tools/filters/gff/sort_gtf.py index 2f757f6793c..fa4904a0612 100644 --- a/tools/filters/gff/sort_gtf.py +++ b/tools/filters/gff/sort_gtf.py @@ -24,5 +24,6 @@ for feature in read_unordered_gtf( open( in_fname, 'r' ) ): # Print feature. for interval in feature.intervals: out.write( "\t".join(interval.fields) ) +out.close() # TODO: print status information: how many lines processed and features found. \ No newline at end of file From 3d5fd7b5e27d33df02ea929f73250163f2e4e007 Mon Sep 17 00:00:00 2001 From: Dannon Baker Date: Fri, 17 Aug 2012 11:35:31 -0400 Subject: [PATCH 18/87] Migrate NCBI Blast+ (and blastxml) to the toolshed. --- lib/galaxy/datatypes/registry.py | 3 - lib/galaxy/datatypes/xml.py | 120 ------------------ .../tool_shed/migrate/versions/0004_tools.py | 12 ++ scripts/migrate_tools/0004_tools.sh | 4 + scripts/migrate_tools/0004_tools.xml | 12 ++ 5 files changed, 28 insertions(+), 123 deletions(-) create mode 100644 lib/galaxy/tool_shed/migrate/versions/0004_tools.py create mode 100644 scripts/migrate_tools/0004_tools.sh create mode 100644 scripts/migrate_tools/0004_tools.xml diff --git a/lib/galaxy/datatypes/registry.py b/lib/galaxy/datatypes/registry.py index 50315904e00..b054b423d11 100644 --- a/lib/galaxy/datatypes/registry.py +++ b/lib/galaxy/datatypes/registry.py @@ -276,7 +276,6 @@ class Registry( object ): 'axt' : sequence.Axt(), 'bam' : binary.Bam(), 'bed' : interval.Bed(), - 'blastxml' : xml.BlastXml(), 'coverage' : coverage.LastzCoverage(), 'customtrack' : interval.CustomTrack(), 'csfasta' : sequence.csFasta(), @@ -310,7 +309,6 @@ class Registry( object ): 'axt' : 'text/plain', 'bam' : 'application/octet-stream', 'bed' : 'text/plain', - 'blastxml' : 'application/xml', 'customtrack' : 'text/plain', 'csfasta' : 'text/plain', 'eland' : 'application/octet-stream', @@ -348,7 +346,6 @@ class Registry( object ): self.sniff_order = [ binary.Bam(), binary.Sff(), - xml.BlastXml(), xml.GenericXml(), sequence.Maf(), sequence.Lav(), diff --git a/lib/galaxy/datatypes/xml.py b/lib/galaxy/datatypes/xml.py index ef1f1623e28..37f34f55169 100644 --- a/lib/galaxy/datatypes/xml.py +++ b/lib/galaxy/datatypes/xml.py @@ -27,9 +27,6 @@ class GenericXml( data.Text ): >>> fname = get_test_fname( 'megablast_xml_parser_test1.blastxml' ) >>> GenericXml().sniff( fname ) True - >>> fname = get_test_fname( 'tblastn_four_human_vs_rhodopsin.xml' ) - >>> BlastXml().sniff( fname ) - True >>> fname = get_test_fname( 'interval.interval' ) >>> GenericXml().sniff( fname ) False @@ -50,123 +47,6 @@ class GenericXml( data.Text ): data.Text.merge(split_files, output_file) merge = staticmethod(merge) -class BlastXml( GenericXml ): - """NCBI Blast XML Output data""" - file_ext = "blastxml" - - def set_peek( self, dataset, is_multi_byte=False ): - """Set the peek and blurb text""" - if not dataset.dataset.purged: - dataset.peek = data.get_file_peek( dataset.file_name, is_multi_byte=is_multi_byte ) - dataset.blurb = 'NCBI Blast XML data' - else: - dataset.peek = 'file does not exist' - dataset.blurb = 'file purged from disk' - def sniff( self, filename ): - """ - Determines whether the file is blastxml - - >>> fname = get_test_fname( 'megablast_xml_parser_test1.blastxml' ) - >>> BlastXml().sniff( fname ) - True - >>> fname = get_test_fname( 'tblastn_four_human_vs_rhodopsin.xml' ) - >>> BlastXml().sniff( fname ) - True - >>> fname = get_test_fname( 'interval.interval' ) - >>> BlastXml().sniff( fname ) - False - """ - #TODO - Use a context manager on Python 2.5+ to close handle - handle = open(filename) - line = handle.readline() - if line.strip() != '': - handle.close() - return False - line = handle.readline() - if line.strip() not in ['', - '']: - handle.close() - return False - line = handle.readline() - if line.strip() != '': - handle.close() - return False - handle.close() - return True - - def merge(split_files, output_file): - """Merging multiple XML files is non-trivial and must be done in subclasses.""" - if len(split_files) == 1: - #For one file only, use base class method (move/copy) - return data.Text.merge(split_files, output_file) - out = open(output_file, "w") - h = None - for f in split_files: - h = open(f) - body = False - header = h.readline() - if not header: - out.close() - h.close() - raise ValueError("BLAST XML file %s was empty" % f) - if header.strip() != '': - out.write(header) #for diagnosis - out.close() - h.close() - raise ValueError("%s is not an XML file!" % f) - line = h.readline() - header += line - if line.strip() not in ['', - '']: - out.write(header) #for diagnosis - out.close() - h.close() - raise ValueError("%s is not a BLAST XML file!" % f) - while True: - line = h.readline() - if not line: - out.write(header) #for diagnosis - out.close() - h.close() - raise ValueError("BLAST XML file %s ended prematurely" % f) - header += line - if "" in line: - break - if len(header) > 10000: - #Something has gone wrong, don't load too much into memory! - #Write what we have to the merged file for diagnostics - out.write(header) - out.close() - h.close() - raise ValueError("BLAST XML file %s has too long a header!" % f) - if "" not in header: - out.close() - h.close() - raise ValueError("%s is not a BLAST XML file:\n%s\n..." % (f, header)) - if f == split_files[0]: - out.write(header) - old_header = header - elif old_header[:300] != header[:300]: - #Enough to check and match - out.close() - h.close() - raise ValueError("BLAST XML headers don't match for %s and %s - have:\n%s\n...\n\nAnd:\n%s\n...\n" \ - % (split_files[0], f, old_header[:300], header[:300])) - else: - out.write(" \n") - for line in h: - if "" in line: - break - #TODO - Increment and if required automatic query names - #like Query_3 to be increasing? - out.write(line) - h.close() - out.write(" \n") - out.write("\n") - out.close() - merge = staticmethod(merge) - - class MEMEXml( GenericXml ): """MEME XML Output data""" file_ext = "memexml" diff --git a/lib/galaxy/tool_shed/migrate/versions/0004_tools.py b/lib/galaxy/tool_shed/migrate/versions/0004_tools.py new file mode 100644 index 00000000000..f20bd83b251 --- /dev/null +++ b/lib/galaxy/tool_shed/migrate/versions/0004_tools.py @@ -0,0 +1,12 @@ +""" +The NCBI BLAST+ tools have been eliminated from the distribution. The tools and datatypes are are now available in repositories named ncbi_blast_plus +and blast_datatypes, respectively, from the main Galaxy tool shed at http://toolshed.g2.bx.psu.edu will be installed into your local Galaxy instance +at the location discussed above by running the following command. +""" + +import sys + +def upgrade(): + print __doc__ +def downgrade(): + pass diff --git a/scripts/migrate_tools/0004_tools.sh b/scripts/migrate_tools/0004_tools.sh new file mode 100644 index 00000000000..40b76956fa2 --- /dev/null +++ b/scripts/migrate_tools/0004_tools.sh @@ -0,0 +1,4 @@ +#!/bin/sh + +cd `dirname $0`/../.. +python ./scripts/migrate_tools/migrate_tools.py 0004_tools.xml $@ diff --git a/scripts/migrate_tools/0004_tools.xml b/scripts/migrate_tools/0004_tools.xml new file mode 100644 index 00000000000..53664f9823c --- /dev/null +++ b/scripts/migrate_tools/0004_tools.xml @@ -0,0 +1,12 @@ + + + + + + + + + + + + From 97e69fbbdfd820b3044ef3289484850bf3faa566 Mon Sep 17 00:00:00 2001 From: Dannon Baker Date: Fri, 17 Aug 2012 13:55:37 -0400 Subject: [PATCH 19/87] Fix chunk-serving logic for tabular files. --- lib/galaxy/datatypes/tabular.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/datatypes/tabular.py b/lib/galaxy/datatypes/tabular.py index ce3b043b658..0939edae159 100644 --- a/lib/galaxy/datatypes/tabular.py +++ b/lib/galaxy/datatypes/tabular.py @@ -264,10 +264,10 @@ class Tabular( data.Text ): def display_data(self, trans, dataset, preview=False, filename=None, to_ext=None, chunk=None): #TODO Prevent failure when displaying extremely long > 50kb lines. - if to_ext or not preview: - return self._serve_raw(trans, dataset, to_ext) if chunk: return self.get_chunk(trans, dataset, chunk) + if to_ext or not preview: + return self._serve_raw(trans, dataset, to_ext) else: column_names = 'null' if dataset.metadata.column_names: @@ -644,4 +644,5 @@ class FeatureLocationIndex( Tabular ): """ file_ext='fli' MetadataElement( name="columns", default=2, desc="Number of columns", readonly=True, visible=False ) - MetadataElement( name="column_types", default=['str', 'str'], param=metadata.ColumnTypesParameter, desc="Column types", readonly=True, visible=False, no_value=[] ) \ No newline at end of file + MetadataElement( name="column_types", default=['str', 'str'], param=metadata.ColumnTypesParameter, desc="Column types", readonly=True, visible=False, no_value=[] ) + From 3360fe30725e8aeb45595a44936470753d264811 Mon Sep 17 00:00:00 2001 From: Carl Eberhard Date: Fri, 17 Aug 2012 14:39:32 -0400 Subject: [PATCH 20/87] new data for bwa_wrapper/bwa_color_wrapper functional tests From dbc2c27a82c11d6693b92c9dd8358879390cfca8 Mon Sep 17 00:00:00 2001 From: Scott McManus Date: Fri, 17 Aug 2012 15:15:48 -0400 Subject: [PATCH 21/87] Minor tweak to exit code handling --- lib/galaxy/jobs/__init__.py | 18 ++++++------------ 1 file changed, 6 insertions(+), 12 deletions(-) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 5ac63f9957b..50c185834c2 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -490,7 +490,6 @@ class JobWrapper( object ): if stderr contains anything, then False is returned. Note that the job id is just for messages. """ - err_msg = "" # By default, the tool succeeded. This covers the case where the code # has a bug but the tool was ok, and it lets a workflow continue. success = True @@ -507,7 +506,7 @@ class JobWrapper( object ): # Check the exit code ranges in the order in which # they were specified. Each exit_code is a StdioExitCode # that includes an applicable range. If the exit code was in - # that range, then apply the error level and add in a message. + # that range, then apply the error level and add a message. # If we've reached a fatal error rule, then stop. max_error_level = galaxy.tools.StdioErrorLevel.NO_ERROR for stdio_exit_code in self.tool.stdio_exit_codes: @@ -515,20 +514,16 @@ class JobWrapper( object ): tool_exit_code <= stdio_exit_code.range_end ): # Tack on a generic description of the code # plus a specific code description. For example, - # this might append "Job 42: Warning: Out of Memory\n". - # TODO: Find somewhere to stick the err_msg - - # possibly to the source (stderr/stdout), possibly - # in a new db column. + # this might prepend "Job 42: Warning: Out of Memory\n". code_desc = stdio_exit_code.desc if ( None == code_desc ): code_desc = "" - tool_msg = ( "Job %s: %s: Exit code %d: %s" % ( - job.get_id_tag(), - galaxy.tools.StdioErrorLevel.desc( tool_exit_code ), + tool_msg = ( "%s: Exit code %d: %s" % ( + galaxy.tools.StdioErrorLevel.desc( stdio_exit_code.error_level ), tool_exit_code, code_desc ) ) - log.info( tool_msg ) - stderr = err_msg + stderr + log.info( "Job %s: %s" % (job.get_id_tag(), tool_msg) ) + stderr = tool_msg + "\n" + stderr max_error_level = max( max_error_level, stdio_exit_code.error_level ) if ( max_error_level >= @@ -571,7 +566,6 @@ class JobWrapper( object ): re.IGNORECASE ) if ( regex_match ): rexmsg = self.regex_err_msg( regex_match, regex) - # DELETEME log.info( "Job %s: %s" % ( job.get_id_tag(), rexmsg ) ) stderr = rexmsg + "\n" + stderr From 860bb2673c6a0324a6d0dc4738fe632d2e1a458b Mon Sep 17 00:00:00 2001 From: Jeremy Goecks Date: Fri, 17 Aug 2012 17:49:41 -0400 Subject: [PATCH 22/87] Converters enhancements: (a) ignore track lines in wiggle to bigwig converter and add clip option and (b) enable visualization of bedgraph datasets via a bedgraph to bigwig converter. --- .../converters/bedgraph_to_bigwig_converter.xml | 14 ++++++++++++++ .../converters/wig_to_bigwig_converter.xml | 2 +- lib/galaxy/datatypes/interval.py | 5 +++-- 3 files changed, 18 insertions(+), 3 deletions(-) create mode 100644 lib/galaxy/datatypes/converters/bedgraph_to_bigwig_converter.xml diff --git a/lib/galaxy/datatypes/converters/bedgraph_to_bigwig_converter.xml b/lib/galaxy/datatypes/converters/bedgraph_to_bigwig_converter.xml new file mode 100644 index 00000000000..77d2c2a42e6 --- /dev/null +++ b/lib/galaxy/datatypes/converters/bedgraph_to_bigwig_converter.xml @@ -0,0 +1,14 @@ + \ No newline at end of file diff --git a/lib/galaxy/datatypes/converters/wig_to_bigwig_converter.xml b/lib/galaxy/datatypes/converters/wig_to_bigwig_converter.xml index e85c3a7b138..d90702efe78 100644 --- a/lib/galaxy/datatypes/converters/wig_to_bigwig_converter.xml +++ b/lib/galaxy/datatypes/converters/wig_to_bigwig_converter.xml @@ -1,6 +1,6 @@