mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge pull request #4830 from phac-nml/dev
Workflow-to-Job Scheduling Parameters
This commit is contained in:
@@ -192,6 +192,7 @@ var View = Backbone.View.extend({
|
||||
this._renderParameters();
|
||||
this._renderHistory();
|
||||
this._renderUseCachedJob();
|
||||
this._renderResourceParameters();
|
||||
_.each(this.steps, step => {
|
||||
self._renderStep(step);
|
||||
});
|
||||
@@ -302,6 +303,19 @@ var View = Backbone.View.extend({
|
||||
});
|
||||
this._append(this.$steps, this.history_form.$el);
|
||||
},
|
||||
|
||||
/** Render Workflow Options */
|
||||
_renderResourceParameters: function() {
|
||||
this.workflow_resource_parameters_form = null;
|
||||
if(!_.isEmpty(this.model.get('workflow_resource_parameters'))){
|
||||
this.workflow_resource_parameters_form = new Form({
|
||||
cls : 'ui-portlet-narrow',
|
||||
title : '<b>Workflow Resource Options</b>',
|
||||
inputs : this.model.get('workflow_resource_parameters')
|
||||
});
|
||||
this._append( this.$steps, this.workflow_resource_parameters_form.$el );
|
||||
}
|
||||
},
|
||||
|
||||
/** Render job caching option */
|
||||
_renderUseCachedJob: function() {
|
||||
@@ -524,6 +538,7 @@ var View = Backbone.View.extend({
|
||||
var job_def = {
|
||||
new_history_name: history_form_data["new_history|name"] ? history_form_data["new_history|name"] : null,
|
||||
history_id: !history_form_data["new_history|name"] ? this.model.get("history_id") : null,
|
||||
resource_params: this.workflow_resource_parameters_form ? this.workflow_resource_parameters_form.data.create() : {},
|
||||
replacement_params: this.wp_form ? this.wp_form.data.create() : {},
|
||||
parameters: {},
|
||||
// Tool form will submit flat maps for each parameter
|
||||
|
||||
@@ -1635,6 +1635,23 @@ galaxy:
|
||||
# processors, memory and walltime.
|
||||
#job_resource_params_file: config/job_resource_params_conf.xml
|
||||
|
||||
# Similar to the above parameter, workflows can describe parameters
|
||||
# used to influence scheduling of jobs within the workflow. This
|
||||
# requires both a description of the fields available (which defaults
|
||||
# to the definitions in job_resource_params_file if not set).
|
||||
#workflow_resource_params_file: config/workflow_resource_params_conf.xml
|
||||
|
||||
# This parameter describes how to map users and workflows to a set of
|
||||
# workflow resource parameter to present (typically input IDs from
|
||||
# workflow_resource_params_file). If this this is a function reference
|
||||
# it will be passed various inputs (workflow model object and user)
|
||||
# and it should produce a list of input IDs. If it is a path it is
|
||||
# expected to an XML or YAML file describing how to map group names to
|
||||
# parameter descriptions (additional types of mappings via these files
|
||||
# could be implemented but haven't yet - for instance using workflow
|
||||
# tags to do the mapping).
|
||||
#workflow_resource_params_mapper: config/workflow_resource_mapper_conf.yml
|
||||
|
||||
# If using job concurrency limits (configured in job_config_file),
|
||||
# several extra database queries must be performed to determine the
|
||||
# number of jobs a user has dispatched to a given destination. By
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
by_group:
|
||||
default: default
|
||||
groups:
|
||||
default: [project, priority]
|
||||
prio_basic: [{name: priority, options: ["low", "med"]}]
|
||||
prio_advanced: [{name: priority, options: ["low", "med", "high"]}]
|
||||
prio_super:
|
||||
- time
|
||||
- memory
|
||||
- processors
|
||||
- name: priority
|
||||
options:
|
||||
- ultra
|
||||
- plus_ultra
|
||||
@@ -0,0 +1,13 @@
|
||||
<parameters>
|
||||
<param label="Processors" name="processors" type="integer" size="2" min="1" max="64" value="" help="Number of processing cores, 'ppn' value (1-64). Leave blank to use default value." />
|
||||
<param label="Memory" name="memory" type="integer" size="3" min="1" max="256" value="" help="Memory size in gigabytes, 'pmem' value (1-256). Leave blank to use default value." />
|
||||
<param label="Time" name="time" type="integer" size="3" min="1" max="744" value="" help="Maximum job time in hours, 'walltime' value (1-744). Leave blank to use default value." />
|
||||
<param label="Project" name="project" type="text" value="" help="Project to assign resource allocation to. Leave blank to use default value." />
|
||||
<param label="Workflow Job Priority" name="priority" type="select" value="med" help="What priority should the jobs in this workflow run at? (Overrides any declared job priority)">
|
||||
<option value="low" label="Low"/>
|
||||
<option value="med" label="Medium"/>
|
||||
<option value="high" label="High"/>
|
||||
<option value="ultra" label="Ultra"/>
|
||||
<option value="plus_ultra" label="Plus Ultra"/>
|
||||
</param>
|
||||
</parameters>
|
||||
@@ -47,6 +47,7 @@ PATH_DEFAULTS = dict(
|
||||
error_report_file=['config/error_report.yml', 'config/error_report.yml.sample'],
|
||||
dependency_resolvers_config_file=['config/dependency_resolvers_conf.xml', 'dependency_resolvers_conf.xml'],
|
||||
job_resource_params_file=['config/job_resource_params_conf.xml', 'job_resource_params_conf.xml'],
|
||||
workflow_resource_params_file=['config/workflow_resource_params_conf.xml', 'workflow_resource_params_conf.xml'],
|
||||
migrated_tools_config=['migrated_tools_conf.xml', 'config/migrated_tools_conf.xml'],
|
||||
object_store_config_file=['config/object_store_conf.xml', 'object_store_conf.xml'],
|
||||
openid_config_file=['config/openid_conf.xml', 'openid_conf.xml', 'config/openid_conf.xml.sample'],
|
||||
@@ -391,6 +392,15 @@ class Configuration(object):
|
||||
self.maximum_workflow_invocation_duration = int(kwargs.get("maximum_workflow_invocation_duration", 2678400))
|
||||
self.maximum_workflow_jobs_per_scheduling_iteration = int(kwargs.get("maximum_workflow_jobs_per_scheduling_iteration", -1))
|
||||
|
||||
workflow_resource_params_mapper = kwargs.get("workflow_resource_params_mapper", None)
|
||||
if not workflow_resource_params_mapper:
|
||||
workflow_resource_params_mapper = None
|
||||
elif ":" not in workflow_resource_params_mapper:
|
||||
# Assume it is not a Python function, so a file
|
||||
workflow_resource_params_mapper = self.resolve_path(workflow_resource_params_mapper)
|
||||
# else: a Python a function!
|
||||
self.workflow_resource_params_mapper = workflow_resource_params_mapper
|
||||
|
||||
self.cache_user_job_count = string_as_bool(kwargs.get('cache_user_job_count', False))
|
||||
self.pbs_application_server = kwargs.get('pbs_application_server', "")
|
||||
self.pbs_dataset_server = kwargs.get('pbs_dataset_server', "")
|
||||
|
||||
@@ -398,17 +398,7 @@ class JobConfiguration(ConfiguresHandlers):
|
||||
return conditional_element
|
||||
|
||||
def __parse_resource_parameters(self):
|
||||
if os.path.exists(self.app.config.job_resource_params_file):
|
||||
resource_param_file = self.app.config.job_resource_params_file
|
||||
try:
|
||||
resource_definitions = util.parse_xml(resource_param_file)
|
||||
except Exception as e:
|
||||
raise config_exception(e, resource_param_file)
|
||||
resource_definitions_root = resource_definitions.getroot()
|
||||
# TODO: Also handling conditionals would be awesome!
|
||||
for parameter_elem in resource_definitions_root.findall("param"):
|
||||
name = parameter_elem.get("name")
|
||||
self.resource_parameters[name] = parameter_elem
|
||||
self.resource_parameters = util.parse_resource_parameters(self.app.config.job_resource_params_file)
|
||||
|
||||
def __get_params(self, parent):
|
||||
"""Parses any child <param> tags in to a dictionary suitable for persistence.
|
||||
|
||||
@@ -3,6 +3,7 @@ from __future__ import print_function
|
||||
import argparse
|
||||
import collections
|
||||
import copy
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
@@ -1229,7 +1230,35 @@ def map_tool_to_destination(
|
||||
default_priority = 'med'
|
||||
priority = default_priority
|
||||
|
||||
if config is not None:
|
||||
# fetch priority information from workflow/job parameters
|
||||
job_parameter_list = job.get_parameters()
|
||||
workflow_params = None
|
||||
job_params = None
|
||||
if job_parameter_list is not None:
|
||||
for param in job_parameter_list:
|
||||
if param.name == "__workflow_resource_params__":
|
||||
workflow_params = param.value
|
||||
if param.name == "__job_resource":
|
||||
job_params = param.value
|
||||
|
||||
# Priority coming from workflow invocation takes precedence over job specific priorities
|
||||
if workflow_params is not None:
|
||||
resource_params = json.loads(workflow_params)
|
||||
if 'priority' in resource_params:
|
||||
# For by_group mapping, this priority has already been validated when the
|
||||
# request was created.
|
||||
if resource_params['priority'] is not None:
|
||||
priority = resource_params['priority']
|
||||
|
||||
elif job_params is not None:
|
||||
resource_params = json.loads(job_params)
|
||||
if 'priority' in resource_params:
|
||||
if resource_params['priority'] is not None:
|
||||
priority = resource_params['priority']
|
||||
|
||||
if fail_message is not None:
|
||||
destination = "fail"
|
||||
elif config is not None:
|
||||
# get the users priority
|
||||
if "users" in config:
|
||||
if user_email in config["users"]:
|
||||
@@ -1363,7 +1392,8 @@ def map_tool_to_destination(
|
||||
fail_message = "Job '" + str(tool.old_id) + "' failed; "
|
||||
fail_message += "no global default destination specified in config!"
|
||||
|
||||
# if config is not None
|
||||
# if fail_message is not None
|
||||
# elif config is not None
|
||||
else:
|
||||
destination = "fail"
|
||||
fail_message = "No config file supplied!"
|
||||
@@ -1377,7 +1407,7 @@ def map_tool_to_destination(
|
||||
if config is not None:
|
||||
if destination == "fail":
|
||||
output = "An error occurred: " + fail_message
|
||||
|
||||
log.debug(output)
|
||||
else:
|
||||
output = "Running '" + str(tool.old_id) + "' with '"
|
||||
output += destination + "'."
|
||||
|
||||
@@ -137,6 +137,11 @@ class JobRunnerMapper(object):
|
||||
workflow_invocation_uuid = param_values.get("__workflow_invocation_uuid__", None)
|
||||
actual_args["workflow_invocation_uuid"] = workflow_invocation_uuid
|
||||
|
||||
if "workflow_resource_params" in function_arg_names:
|
||||
param_values = job.raw_param_dict()
|
||||
workflow_resource_params = param_values.get("__workflow_resource_params__", None)
|
||||
actual_args["workflow_resource_params"] = workflow_resource_params
|
||||
|
||||
return expand_function(**actual_args)
|
||||
|
||||
def __job_params(self, job):
|
||||
|
||||
@@ -36,6 +36,7 @@ from galaxy.workflow.modules import (
|
||||
ToolModule,
|
||||
WorkflowModuleInjector
|
||||
)
|
||||
from galaxy.workflow.resources import get_resource_mapper_function
|
||||
from galaxy.workflow.steps import attach_ordered_steps
|
||||
from .base import decode_id
|
||||
|
||||
@@ -205,6 +206,7 @@ class WorkflowContentsManager(UsesAnnotations):
|
||||
|
||||
def __init__(self, app):
|
||||
self.app = app
|
||||
self._resource_mapper_function = get_resource_mapper_function(app)
|
||||
|
||||
def build_workflow_from_dict(
|
||||
self,
|
||||
@@ -432,14 +434,20 @@ class WorkflowContentsManager(UsesAnnotations):
|
||||
step_model['messages'] = step.upgrade_messages
|
||||
step_models.append(step_model)
|
||||
return {
|
||||
'id' : trans.app.security.encode_id(stored.id),
|
||||
'history_id' : trans.app.security.encode_id(trans.history.id) if trans.history else None,
|
||||
'name' : stored.name,
|
||||
'steps' : step_models,
|
||||
'step_version_changes' : step_version_changes,
|
||||
'has_upgrade_messages' : has_upgrade_messages
|
||||
'id': trans.app.security.encode_id(stored.id),
|
||||
'history_id': trans.app.security.encode_id(trans.history.id) if trans.history else None,
|
||||
'name': stored.name,
|
||||
'steps': step_models,
|
||||
'step_version_changes': step_version_changes,
|
||||
'has_upgrade_messages': has_upgrade_messages,
|
||||
'workflow_resource_parameters': self._workflow_resource_parameters(trans, stored, workflow),
|
||||
}
|
||||
|
||||
def _workflow_resource_parameters(self, trans, stored, workflow):
|
||||
"""Get workflow scheduling resource parameters for this user and workflow or None if unconfigured.
|
||||
"""
|
||||
return self._resource_mapper_function(trans=trans, stored_workflow=stored, workflow=workflow)
|
||||
|
||||
def _workflow_to_dict_editor(self, trans, stored):
|
||||
workflow = stored.latest_workflow
|
||||
# Pack workflow data into a dictionary and return
|
||||
|
||||
@@ -4150,6 +4150,16 @@ class WorkflowInvocation(UsesCreateAndUpdateTime, Dictifiable):
|
||||
request_to_content.workflow_step_id = step_id
|
||||
self.input_step_parameters.append(request_to_content)
|
||||
|
||||
@property
|
||||
def resource_parameters(self):
|
||||
resource_type = WorkflowRequestInputParameter.types.RESOURCE_PARAMETERS
|
||||
_resource_parameters = {}
|
||||
for input_parameter in self.input_parameters:
|
||||
if input_parameter.type == resource_type:
|
||||
_resource_parameters[input_parameter.name] = input_parameter.value
|
||||
|
||||
return _resource_parameters
|
||||
|
||||
def has_input_for_step(self, step_id):
|
||||
for content in self.input_datasets:
|
||||
if content.workflow_step_id == step_id:
|
||||
@@ -4258,7 +4268,8 @@ class WorkflowRequestInputParameter(Dictifiable):
|
||||
dict_collection_visible_keys = ['id', 'name', 'value', 'type']
|
||||
types = Bunch(
|
||||
REPLACEMENT_PARAMETERS='replacements',
|
||||
META_PARAMETERS='meta', #
|
||||
META_PARAMETERS='meta',
|
||||
RESOURCE_PARAMETERS='resource',
|
||||
)
|
||||
|
||||
def __init__(self, name=None, value=None, type=None):
|
||||
|
||||
@@ -31,7 +31,7 @@ class PartialJobExecution(Exception):
|
||||
MappingParameters = collections.namedtuple("MappingParameters", ["param_template", "param_combinations"])
|
||||
|
||||
|
||||
def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, collection_info=None, workflow_invocation_uuid=None, invocation_step=None, max_num_jobs=None, job_callback=None, completed_jobs=None):
|
||||
def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, collection_info=None, workflow_invocation_uuid=None, invocation_step=None, max_num_jobs=None, job_callback=None, completed_jobs=None, workflow_resource_parameters=None):
|
||||
"""
|
||||
Execute a tool and return object containing summary (output data, number of
|
||||
failures, etc...).
|
||||
@@ -58,7 +58,12 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle
|
||||
# Only workflow invocation code gets to set this, ignore user supplied
|
||||
# values or rerun parameters.
|
||||
del params['__workflow_invocation_uuid__']
|
||||
|
||||
if workflow_resource_parameters:
|
||||
params['__workflow_resource_params__'] = workflow_resource_parameters
|
||||
elif '__workflow_resource_params__' in params:
|
||||
# Only workflow invocation code gets to set this, ignore user supplied
|
||||
# values or rerun parameters.
|
||||
del params['__workflow_resource_params__']
|
||||
job, result = tool.handle_single_execution(trans, rerun_remap_job_id, execution_slice, history, execution_cache, completed_job)
|
||||
if job:
|
||||
message = EXECUTION_SUCCESS_MESSAGE % (tool.id, job.id, job_timer)
|
||||
|
||||
@@ -817,6 +817,22 @@ def xml_text(root, name=None):
|
||||
return ''
|
||||
|
||||
|
||||
def parse_resource_parameters(resource_param_file):
|
||||
"""Code shared between jobs and workflows for reading resource parameter configuration files.
|
||||
|
||||
TODO: Allow YAML in addition to XML.
|
||||
"""
|
||||
resource_parameters = {}
|
||||
if os.path.exists(resource_param_file):
|
||||
resource_definitions = parse_xml(resource_param_file)
|
||||
resource_definitions_root = resource_definitions.getroot()
|
||||
for parameter_elem in resource_definitions_root.findall("param"):
|
||||
name = parameter_elem.get("name")
|
||||
resource_parameters[name] = parameter_elem
|
||||
|
||||
return resource_parameters
|
||||
|
||||
|
||||
# asbool implementation pulled from PasteDeploy
|
||||
truthy = frozenset(['true', 'yes', 'on', 'y', 't', '1'])
|
||||
falsy = frozenset(['false', 'no', 'off', 'n', 'f', '0'])
|
||||
|
||||
@@ -2539,6 +2539,29 @@ mapping:
|
||||
overwrite default job resources such as number of processors, memory and
|
||||
walltime.
|
||||
|
||||
workflow_resource_params_file:
|
||||
type: str
|
||||
default: config/workflow_resource_params_conf.xml
|
||||
required: false
|
||||
desc: |
|
||||
Similar to the above parameter, workflows can describe parameters used to
|
||||
influence scheduling of jobs within the workflow. This requires both a description
|
||||
of the fields available (which defaults to the definitions in
|
||||
job_resource_params_file if not set).
|
||||
|
||||
workflow_resource_params_mapper:
|
||||
type: str
|
||||
default: config/workflow_resource_mapper_conf.yml
|
||||
required: false
|
||||
desc: |
|
||||
This parameter describes how to map users and workflows to a set of workflow
|
||||
resource parameter to present (typically input IDs from workflow_resource_params_file).
|
||||
If this this is a function reference it will be passed various inputs (workflow model
|
||||
object and user) and it should produce a list of input IDs. If it is a path
|
||||
it is expected to an XML or YAML file describing how to map group names to parameter
|
||||
descriptions (additional types of mappings via these files could be implemented but
|
||||
haven't yet - for instance using workflow tags to do the mapping).
|
||||
|
||||
cache_user_job_count:
|
||||
type: bool
|
||||
default: false
|
||||
|
||||
@@ -846,6 +846,7 @@ class ToolModule(WorkflowModule):
|
||||
else:
|
||||
iteration_elements_iter = [None]
|
||||
|
||||
resource_parameters = invocation.resource_parameters
|
||||
for iteration_elements in iteration_elements_iter:
|
||||
execution_state = tool_state.copy()
|
||||
# TODO: Move next step into copy()
|
||||
@@ -918,7 +919,8 @@ class ToolModule(WorkflowModule):
|
||||
invocation_step=invocation_step,
|
||||
max_num_jobs=max_num_jobs,
|
||||
job_callback=lambda job: self._handle_post_job_actions(step, job, invocation.replacement_dict),
|
||||
completed_jobs=completed_jobs
|
||||
completed_jobs=completed_jobs,
|
||||
workflow_resource_parameters=resource_parameters
|
||||
)
|
||||
complete = True
|
||||
except PartialJobExecution as pje:
|
||||
|
||||
@@ -0,0 +1,176 @@
|
||||
"""This package is something a placeholder for workflow resource parameters.
|
||||
|
||||
This file defines the baked in resource mapper types, and this package contains an
|
||||
example of a more open, pluggable approach with greater control.
|
||||
"""
|
||||
import functools
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
from copy import deepcopy
|
||||
|
||||
import yaml
|
||||
|
||||
import galaxy.util
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def get_resource_mapper_function(app):
|
||||
config = app.config
|
||||
mapper = getattr(config, "workflow_resource_params_mapper", None)
|
||||
|
||||
if mapper is None:
|
||||
return _null_mapper_function
|
||||
elif ":" in mapper:
|
||||
raw_function = _import_resource_mapping_function(mapper)
|
||||
# Bind resource parameters here just to not re-parse over and over.
|
||||
workflow_resource_params = _read_defined_parameter_definitions(config)
|
||||
return functools.partial(raw_function, workflow_resource_params=workflow_resource_params)
|
||||
else:
|
||||
workflow_resource_params = _read_defined_parameter_definitions(config)
|
||||
with open(mapper, "r") as f:
|
||||
mapper_definition = yaml.load(f)
|
||||
|
||||
if "by_group" in mapper_definition:
|
||||
by_group = mapper_definition["by_group"]
|
||||
return functools.partial(_resource_parameters_by_group, by_group=by_group, workflow_resource_params=workflow_resource_params)
|
||||
else:
|
||||
raise Exception("Currently workflow parameter mapper definitions require a by_group definition.")
|
||||
|
||||
|
||||
def _read_defined_parameter_definitions(config):
|
||||
params_file = getattr(config, "workflow_resource_params_file", None)
|
||||
if not params_file or not os.path.exists(params_file):
|
||||
# Just re-use job resource parameters.
|
||||
params_file = getattr(config, "job_resource_params_file", None)
|
||||
if not params_file or not os.path.exists(params_file):
|
||||
params_file = None
|
||||
log.debug("Loading workflow resource parameter definitions from %s" % params_file)
|
||||
if params_file:
|
||||
return galaxy.util.parse_resource_parameters(params_file)
|
||||
else:
|
||||
return {}
|
||||
|
||||
|
||||
def _resource_parameters_by_group(trans, **kwds):
|
||||
user = trans.user
|
||||
by_group = kwds["by_group"]
|
||||
workflow_resource_params = kwds["workflow_resource_params"]
|
||||
|
||||
params = []
|
||||
if validate_by_group_workflow_parameters_mapper(by_group, workflow_resource_params):
|
||||
user_permissions = {}
|
||||
user_groups = []
|
||||
for g in user.groups:
|
||||
user_groups.append(g.group.name)
|
||||
default_group = by_group.get('default', None)
|
||||
for group_name, group_def in by_group.get("groups", {}).items():
|
||||
if group_name == default_group or group_name in user_groups:
|
||||
for tag in group_def:
|
||||
if type(tag) is dict:
|
||||
if tag.get('name') not in user_permissions:
|
||||
user_permissions[tag.get('name')] = {}
|
||||
for option in tag.get('options'):
|
||||
user_permissions[tag.get('name')][option] = {}
|
||||
else:
|
||||
if tag not in user_permissions:
|
||||
user_permissions[tag] = {}
|
||||
|
||||
# user_permissions is now set.
|
||||
params = get_workflow_parameter_list(workflow_resource_params, user_permissions)
|
||||
return params
|
||||
|
||||
|
||||
# returns an array of parameters that a users set of permissions can access.
|
||||
def get_workflow_parameter_list(params, user_permissions):
|
||||
param_list = []
|
||||
for param_name, param_elem in params.items():
|
||||
attr = deepcopy(param_elem.attrib)
|
||||
if attr['name'] in user_permissions:
|
||||
# Allow 'select' type parameters to be used
|
||||
if attr['type'] == 'select':
|
||||
option_data = []
|
||||
reject_list = []
|
||||
for option_elem in param_elem.findall("option"):
|
||||
if option_elem.attrib['value'] in user_permissions[attr['name']]:
|
||||
option_data.append({
|
||||
'label': option_elem.attrib['label'],
|
||||
'value': option_elem.attrib['value']
|
||||
})
|
||||
else:
|
||||
reject_list.append(option_elem.attrib['label'])
|
||||
attr['data'] = option_data
|
||||
attr_help = ""
|
||||
if 'help' in attr:
|
||||
attr_help = attr['help']
|
||||
if reject_list:
|
||||
attr_help += "<br/><br/>The following options are available but disabled.<br/>" + \
|
||||
str(reject_list) + \
|
||||
"<br/>If you believe this is a mistake, please contact your Galaxy admin."
|
||||
attr['help'] = attr_help
|
||||
|
||||
param_list.append(attr)
|
||||
return param_list
|
||||
|
||||
|
||||
def validate_by_group_workflow_parameters_mapper(by_group, workflow_resource_params):
|
||||
valid = True
|
||||
try:
|
||||
if 'default' not in by_group:
|
||||
raise Exception("'workflow_resource_params_mapper' YAML file is malformed, 'default' attribute not found!")
|
||||
default_group = by_group['default']
|
||||
if 'groups' not in by_group:
|
||||
raise Exception("'workflow_resource_params_mapper' YAML file is malformed, 'groups' attribute not found!")
|
||||
if default_group not in by_group['groups']:
|
||||
raise Exception("'workflow_resource_params_mapper' YAML file is malformed, default group with title '" +
|
||||
default_group + "' not found in 'groups'!")
|
||||
for group in by_group['groups']:
|
||||
for attrib in by_group['groups'][group]:
|
||||
if type(attrib) is dict:
|
||||
if 'name' not in attrib:
|
||||
raise Exception("'workflow_resource_params_mapper' YAML file is malformed, "
|
||||
"'name' attribute not found in attribute of group '" + group + "'!")
|
||||
if attrib['name'] not in workflow_resource_params:
|
||||
raise Exception("'workflow_resource_params_mapper' YAML file is malformed, group with name '" +
|
||||
attrib['name'] + "' not found in 'workflow_resource_params'!")
|
||||
if 'options' not in attrib:
|
||||
raise Exception("'workflow_resource_params_mapper' YAML file is malformed, "
|
||||
"'options' attribute not found in attribute of group '" + group + "'!")
|
||||
|
||||
valid_options = []
|
||||
for param_option in workflow_resource_params[attrib['name']]:
|
||||
valid_options.append(param_option.attrib['value'])
|
||||
for option in attrib['options']:
|
||||
if option not in valid_options:
|
||||
raise Exception("'workflow_resource_params_mapper' YAML file is malformed, '" + option +
|
||||
"' in 'options' of '" + attrib['name'] + "' not found in attribute of group '" + group + "'!")
|
||||
else:
|
||||
if attrib not in workflow_resource_params:
|
||||
raise Exception("'workflow_resource_params_mapper' YAML file is malformed, attribute with name "
|
||||
"'" + attrib + "' not found in 'workflow_resource_params'!")
|
||||
|
||||
except Exception as e:
|
||||
log.exception(e)
|
||||
valid = False
|
||||
pass
|
||||
|
||||
return valid
|
||||
|
||||
|
||||
def _import_resource_mapping_function(qualified_function_path):
|
||||
full_module_name, function_name = qualified_function_path.split(":", 1)
|
||||
try:
|
||||
__import__(full_module_name)
|
||||
except ImportError:
|
||||
raise Exception("Failed to find workflow resource mapper module %s" % full_module_name)
|
||||
|
||||
module = sys.modules[full_module_name]
|
||||
if hasattr(module, function_name):
|
||||
return getattr(module, function_name)
|
||||
else:
|
||||
raise Exception("Failed to find workflow resource mapper function %s.%s" % (full_module_name, function_name))
|
||||
|
||||
|
||||
def _null_mapper_function(*args, **kwds):
|
||||
return None
|
||||
@@ -0,0 +1,22 @@
|
||||
import logging
|
||||
log = logging.getLogger( __name__ )
|
||||
|
||||
|
||||
def admin_mapping(trans, stored_workflow, **kwds):
|
||||
"""
|
||||
This example workflow resource parameter mapping simply provides admins the ability to
|
||||
specify priorities for workflows. To enable this setup ``workflow_resource_params_file``
|
||||
in the Galaxy configuration with a priority definition input called "priority" (such
|
||||
as in the example), copy this file without the .sample extension, and set
|
||||
``workflow_resource_params_mapper`` to ``galaxy.workflow.resources.example:admin_mapping``.
|
||||
"""
|
||||
workflow_resource_params = kwds["workflow_resource_params"]
|
||||
if trans.user_is_admin():
|
||||
priority_attrib = workflow_resource_params.get("priority").attrib
|
||||
priority_attrib['data'] = []
|
||||
for child in workflow_resource_params.get('priority').getchildren():
|
||||
priority_attrib['data'].append(child.attrib)
|
||||
time_attrib = workflow_resource_params.get("time").attrib
|
||||
return [priority_attrib, time_attrib]
|
||||
|
||||
return None
|
||||
@@ -7,6 +7,7 @@ from galaxy import (
|
||||
)
|
||||
from galaxy.managers import histories
|
||||
from galaxy.tools.parameters.meta import expand_workflow_inputs
|
||||
from galaxy.workflow.resources import get_resource_mapper_function
|
||||
|
||||
INPUT_STEP_TYPES = ['data_input', 'data_collection_input', 'parameter_input']
|
||||
|
||||
@@ -47,12 +48,14 @@ class WorkflowRunConfig(object):
|
||||
inputs=None,
|
||||
param_map=None,
|
||||
allow_tool_state_corrections=False,
|
||||
use_cached_job=False):
|
||||
use_cached_job=False,
|
||||
resource_params=None):
|
||||
self.target_history = target_history
|
||||
self.replacement_dict = replacement_dict
|
||||
self.copy_inputs_to_history = copy_inputs_to_history
|
||||
self.inputs = inputs or {}
|
||||
self.param_map = param_map or {}
|
||||
self.resource_params = resource_params or {}
|
||||
self.allow_tool_state_corrections = allow_tool_state_corrections
|
||||
self.use_cached_job = use_cached_job
|
||||
|
||||
@@ -305,6 +308,35 @@ def build_workflow_run_configs(trans, workflow, payload):
|
||||
normalized_inputs[key] = value['content']
|
||||
else:
|
||||
normalized_inputs[key] = value
|
||||
resource_params = payload.get('resource_params', {})
|
||||
if resource_params:
|
||||
# quick attempt to validate parameters, just handle select options now since is what
|
||||
# is needed for DTD - arbitrary plugins can define arbitrary logic at runtime in the
|
||||
# destination function. In the future this should be extended to allow arbitrary
|
||||
# pluggable validation.
|
||||
resource_mapper_function = get_resource_mapper_function(trans.app)
|
||||
# TODO: Do we need to do anything with the stored_workflow or can this be removed.
|
||||
resource_parameters = resource_mapper_function(trans=trans, stored_workflow=None, workflow=workflow)
|
||||
for resource_parameter in resource_parameters:
|
||||
if resource_parameter.get("type") == "select":
|
||||
name = resource_parameter.get("name")
|
||||
if name in resource_params:
|
||||
value = resource_params[name]
|
||||
valid_option = False
|
||||
# TODO: How should be handle the case where no selection is made by the user
|
||||
# This can happen when there is a select on the page but the user has no options to select
|
||||
# Here I have the validation pass it through. An alternative may be to remove the parameter if
|
||||
# it is None.
|
||||
if value is None:
|
||||
valid_option = True
|
||||
else:
|
||||
for option_elem in resource_parameter.get('data'):
|
||||
option_value = option_elem.get("value")
|
||||
if value == option_value:
|
||||
valid_option = True
|
||||
if not valid_option:
|
||||
raise exceptions.RequestParameterInvalidException("Invalid value for parameter '%s' found." % name)
|
||||
|
||||
run_configs.append(WorkflowRunConfig(
|
||||
target_history=history,
|
||||
replacement_dict=payload.get('replacement_params', {}),
|
||||
@@ -312,6 +344,7 @@ def build_workflow_run_configs(trans, workflow, payload):
|
||||
param_map=param_map,
|
||||
allow_tool_state_corrections=allow_tool_state_corrections,
|
||||
use_cached_job=use_cached_job,
|
||||
resource_params=resource_params,
|
||||
))
|
||||
|
||||
return run_configs
|
||||
@@ -351,7 +384,8 @@ def workflow_run_config_to_request(trans, run_config, workflow):
|
||||
use_cached_job=run_config.use_cached_job,
|
||||
inputs={},
|
||||
param_map={},
|
||||
allow_tool_state_corrections=run_config.allow_tool_state_corrections
|
||||
allow_tool_state_corrections=run_config.allow_tool_state_corrections,
|
||||
resource_params=run_config.resource_params
|
||||
)
|
||||
subworkflow_invocation = workflow_run_config_to_request(
|
||||
trans,
|
||||
@@ -373,6 +407,9 @@ def workflow_run_config_to_request(trans, run_config, workflow):
|
||||
for step_id, content in run_config.inputs.items():
|
||||
workflow_invocation.add_input(content, step_id)
|
||||
|
||||
resource_parameters = run_config.resource_params
|
||||
for key, value in resource_parameters.items():
|
||||
add_parameter(key, value, param_types.RESOURCE_PARAMETERS)
|
||||
add_parameter("copy_inputs_to_history", "true" if run_config.copy_inputs_to_history else "false", param_types.META_PARAMETERS)
|
||||
add_parameter("use_cached_job", "true" if run_config.use_cached_job else "false", param_types.META_PARAMETERS)
|
||||
return workflow_invocation
|
||||
@@ -384,6 +421,7 @@ def workflow_request_to_run_config(work_request_context, workflow_invocation):
|
||||
replacement_dict = {}
|
||||
inputs = {}
|
||||
param_map = {}
|
||||
resource_params = {}
|
||||
copy_inputs_to_history = None
|
||||
use_cached_job = False
|
||||
for parameter in workflow_invocation.input_parameters:
|
||||
@@ -396,6 +434,8 @@ def workflow_request_to_run_config(work_request_context, workflow_invocation):
|
||||
copy_inputs_to_history = (parameter.value == "true")
|
||||
if parameter.name == 'use_cached_job':
|
||||
use_cached_job = (parameter.value == 'true')
|
||||
elif parameter_type == param_types.RESOURCE_PARAMETERS:
|
||||
resource_params[parameter.name] = parameter.value
|
||||
for input_association in workflow_invocation.input_datasets:
|
||||
inputs[input_association.workflow_step_id] = input_association.dataset
|
||||
for input_association in workflow_invocation.input_dataset_collections:
|
||||
@@ -411,6 +451,7 @@ def workflow_request_to_run_config(work_request_context, workflow_invocation):
|
||||
param_map=param_map,
|
||||
copy_inputs_to_history=copy_inputs_to_history,
|
||||
use_cached_job=use_cached_job,
|
||||
resource_params=resource_params,
|
||||
)
|
||||
return workflow_run_config
|
||||
|
||||
|
||||
@@ -7,6 +7,7 @@ class Job(object):
|
||||
self.input_datasets = []
|
||||
self.input_library_datasets = []
|
||||
self.param_values = dict()
|
||||
self.parameters = []
|
||||
|
||||
def get_param_values(self, app, ignore_errors=False):
|
||||
return self.param_values
|
||||
@@ -17,6 +18,9 @@ class Job(object):
|
||||
def add_input_dataset(self, dataset):
|
||||
self.input_datasets.append(dataset)
|
||||
|
||||
def get_parameters(self):
|
||||
return self.parameters
|
||||
|
||||
|
||||
class InputDataset(object):
|
||||
def __init__(self, name, dataset):
|
||||
|
||||
Reference in New Issue
Block a user