From 4dccf1227eb7648882bb485354dcf71c762f0201 Mon Sep 17 00:00:00 2001 From: Dannon Baker Date: Wed, 29 Aug 2012 10:17:53 -0400 Subject: [PATCH] Merge of rpark's workflow API changes to support parameter execution and workflow creation. TODO: Refactor copied workflow methods back out. --- lib/galaxy/web/api/workflows.py | 151 +++++++++++--------------------- lib/galaxy/web/buildapp.py | 19 ++-- 2 files changed, 54 insertions(+), 116 deletions(-) diff --git a/lib/galaxy/web/api/workflows.py b/lib/galaxy/web/api/workflows.py index 4d025d93fdf..418acf86bcf 100644 --- a/lib/galaxy/web/api/workflows.py +++ b/lib/galaxy/web/api/workflows.py @@ -10,25 +10,12 @@ from galaxy.tools.parameters import visit_input_values, DataToolParameter 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.model.item_attrs import UsesAnnotations 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): +class WorkflowsAPIController(BaseAPIController, UsesAnnotations): @web.expose_api def index(self, trans, **kwd): """ @@ -100,16 +87,16 @@ 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') ): param_map = payload['parameters']; - # ------------------------------------------------------------------------------- # - + # ------------------------------------------------------------------------------- # + + + - - if 'workflow_id' not in payload: # create new if 'installed_repository_file' in payload: @@ -194,7 +181,7 @@ 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 - + #################################################### #################################################### # RPARK: IF TOOL_NAME IN PARAMETER MAP # @@ -204,7 +191,7 @@ class WorkflowsAPIController(BaseAPIController): step.state.inputs[change_param] = change_value; #################################################### #################################################### - + if step.tool_errors: trans.response.status = 400 return "Workflow cannot be run because of validation errors in some steps: %s" % step_errors @@ -267,40 +254,40 @@ class WorkflowsAPIController(BaseAPIController): def workflow_dict( self, trans, workflow_id, **kwd ): """ GET /api/workflows/{encoded_workflow_id}/download - Returns a selected workflow as a json dictionary. + Returns a selected workflow as a json dictionary. """ - + 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 + + # 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") - + ret_dict = self._workflow_to_dict( trans, stored_workflow ); return ret_dict - + @web.expose_api - def delete( self, trans, id, **kwd ): + 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 + + # 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 @@ -309,42 +296,37 @@ class WorkflowsAPIController(BaseAPIController): #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/upload Importing dynamic workflows from the api. Return newly generated workflow id. - Author: rpark - + 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 + + # galaxy workflow newly created id workflow_id = workflow.id; - # api encoded, 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); - + + rval.append(item); + return item; def _workflow_from_dict( self, trans, data, source=None ): @@ -392,21 +374,19 @@ class WorkflowsAPIController(BaseAPIController): 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' ) + #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 = 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: @@ -431,14 +411,14 @@ class WorkflowsAPIController(BaseAPIController): 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 ) @@ -446,8 +426,8 @@ class WorkflowsAPIController(BaseAPIController): 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 @@ -457,22 +437,22 @@ class WorkflowsAPIController(BaseAPIController): ## 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 + annotation_str = step_annotation.annotation ### ----------------------------------- ### - + # Step info step_dict = { 'id': step.order_index, @@ -484,18 +464,18 @@ class WorkflowsAPIController(BaseAPIController): '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, + 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 @@ -559,37 +539,4 @@ class WorkflowsAPIController(BaseAPIController): # 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 92038157877..2c4a3e7e0ac 100644 --- a/lib/galaxy/web/buildapp.py +++ b/lib/galaxy/web/buildapp.py @@ -151,20 +151,11 @@ 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() - # ---------------------------------------------- # - # ---------------------------------------------- # - + # "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'])) + webapp.finalize_config() # Wrap the webapp in some useful middleware if kwargs.get( 'middleware', True ):