mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge of rpark's workflow API changes to support parameter execution and workflow creation.
TODO: Refactor copied workflow methods back out.
This commit is contained in:
+49
-102
@@ -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 )
|
||||
|
||||
|
||||
@@ -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 ):
|
||||
|
||||
Reference in New Issue
Block a user