Remove deprecated job runner URL config options

This commit is contained in:
Nate Coraor
2017-11-17 17:04:14 -05:00
parent b4863b2df8
commit eda29c81c1
5 changed files with 33 additions and 140 deletions
-52
View File
@@ -288,8 +288,6 @@ class Configuration(object):
self.collect_outputs_from = [x.strip() for x in kwargs.get('collect_outputs_from', 'new_file_path,job_working_directory').lower().split(',')]
self.template_path = resolve_path(kwargs.get("template_path", "templates"), self.root)
self.template_cache = resolve_path(kwargs.get("template_cache_path", "database/compiled_templates"), self.root)
self.local_job_queue_workers = int(kwargs.get("local_job_queue_workers", "5"))
self.cluster_job_queue_workers = int(kwargs.get("cluster_job_queue_workers", "3"))
self.job_queue_cleanup_interval = int(kwargs.get("job_queue_cleanup_interval", "5"))
self.cluster_files_directory = os.path.abspath(kwargs.get("cluster_files_directory", "database/pbs"))
@@ -312,11 +310,6 @@ class Configuration(object):
self.output_size_limit = int(kwargs.get('output_size_limit', 0))
self.retry_job_output_collection = int(kwargs.get('retry_job_output_collection', 0))
self.check_job_script_integrity = string_as_bool(kwargs.get("check_job_script_integrity", True))
self.job_walltime = kwargs.get('job_walltime', None)
self.job_walltime_delta = None
if self.job_walltime is not None:
h, m, s = [int(v) for v in self.job_walltime.split(':')]
self.job_walltime_delta = timedelta(0, s, 0, 0, m, h)
self.admin_users = kwargs.get("admin_users", "")
self.admin_users_list = [u.strip() for u in self.admin_users.split(',') if u]
self.mailing_join_addr = kwargs.get('mailing_join_addr', 'galaxy-announce-join@bx.psu.edu')
@@ -349,7 +342,6 @@ class Configuration(object):
self.smtp_password = kwargs.get('smtp_password', None)
self.smtp_ssl = kwargs.get('smtp_ssl', None)
self.track_jobs_in_database = string_as_bool(kwargs.get('track_jobs_in_database', 'True'))
self.start_job_runners = listify(kwargs.get('start_job_runners', ''))
self.expose_dataset_path = string_as_bool(kwargs.get('expose_dataset_path', 'False'))
self.expose_potentially_sensitive_job_metrics = string_as_bool(kwargs.get('expose_potentially_sensitive_job_metrics', 'False'))
self.enable_communication_server = string_as_bool(kwargs.get('enable_communication_server', 'False'))
@@ -397,12 +389,7 @@ class Configuration(object):
self.parallelize_workflow_scheduling_within_histories = string_as_bool(kwargs.get('parallelize_workflow_scheduling_within_histories', 'False'))
self.maximum_workflow_invocation_duration = int(kwargs.get("maximum_workflow_invocation_duration", 2678400))
# Per-user Job concurrency limitations
self.cache_user_job_count = string_as_bool(kwargs.get('cache_user_job_count', False))
self.user_job_limit = int(kwargs.get('user_job_limit', 0))
self.registered_user_job_limit = int(kwargs.get('registered_user_job_limit', self.user_job_limit))
self.anonymous_user_job_limit = int(kwargs.get('anonymous_user_job_limit', self.user_job_limit))
self.default_cluster_job_runner = kwargs.get('default_cluster_job_runner', 'local:///')
self.pbs_application_server = kwargs.get('pbs_application_server', "")
self.pbs_dataset_server = kwargs.get('pbs_dataset_server', "")
self.pbs_dataset_path = kwargs.get('pbs_dataset_path', "")
@@ -592,12 +579,8 @@ class Configuration(object):
self.galaxy_infrastructure_url_set = galaxy_infrastructure_url_set
# Store advanced job management config
self.job_manager = kwargs.get('job_manager', self.server_name).strip()
self.job_handlers = [x.strip() for x in kwargs.get('job_handlers', self.server_name).split(',')]
self.default_job_handlers = [x.strip() for x in kwargs.get('default_job_handlers', ','.join(self.job_handlers)).split(',')]
# Store per-tool runner configs
self.tool_handlers = self.__read_tool_job_config(global_conf_parser, 'galaxy:tool_handlers', 'name')
self.tool_runners = self.__read_tool_job_config(global_conf_parser, 'galaxy:tool_runners', 'url')
# Galaxy messaging (AMQP) configuration options
self.amqp = {}
try:
@@ -734,41 +717,6 @@ class Configuration(object):
self.datatypes_config = self.datatypes_config_file
self.tool_configs = self.tool_config_file
def __read_tool_job_config(self, global_conf_parser, section, key):
try:
tool_runners_config = global_conf_parser.items(section)
# Process config to group multiple configs for the same tool.
rval = {}
for entry in tool_runners_config:
tool_config, val = entry
tool = None
runner_dict = {}
if tool_config.find("[") != -1:
# Found tool with additional params; put params in dict.
tool, params = tool_config[:-1].split("[")
param_dict = {}
for param in params.split(","):
name, value = param.split("@")
param_dict[name] = value
runner_dict['params'] = param_dict
else:
tool = tool_config
# Add runner URL.
runner_dict[key] = val
# Create tool entry if necessary.
if tool not in rval:
rval[tool] = []
# Add entry to runners.
rval[tool].append(runner_dict)
return rval
except configparser.NoSectionError:
return {}
def get(self, key, default):
return self.config_dict.get(key, default)
+33 -79
View File
@@ -46,6 +46,7 @@ TOOL_PROVIDED_JOB_METADATA_KEYS = ['name', 'info', 'dbkey']
# Override with config.default_job_shell.
DEFAULT_JOB_SHELL = '/bin/bash'
DEFAULT_LOCAL_WORKERS = 4
DEFAULT_CLEANUP_JOB = "always"
@@ -139,7 +140,15 @@ class JobConfiguration(object, ConfiguresHandlers):
self.resource_groups = {}
self.default_resource_group = None
self.resource_parameters = {}
self.limits = Bunch()
self.limits = Bunch(registered_user_concurrent_jobs=None,
anonymous_user_concurrent_jobs=None,
walltime=None,
walltime_delta=None,
total_walltime={},
output_size=None,
destination_user_concurrent_jobs={},
destination_total_concurrent_jobs={})
self._is_default = False
default_resubmits = []
default_resubmit_condition = self.app.config.default_job_resubmission_condition
@@ -159,10 +168,11 @@ class JobConfiguration(object, ConfiguresHandlers):
tree = load(job_config_file)
self.__parse_job_conf_xml(tree)
except IOError:
log.warning('Job configuration "%s" does not exist, using legacy'
' job configuration from Galaxy config file "%s" instead'
% (self.app.config.job_config_file, self.app.config.config_file))
self.__parse_job_conf_legacy()
log.warning('Job configuration "%s" does not exist, using default'
' job configuration (this server will run jobs)',
self.app.config.job_config_file)
self._is_default = True
self.__set_default_job_conf()
except Exception as e:
raise config_exception(e, job_config_file)
@@ -284,15 +294,6 @@ class JobConfiguration(object, ConfiguresHandlers):
total_walltime=str,
output_size=util.size_to_bytes)
self.limits = Bunch(registered_user_concurrent_jobs=None,
anonymous_user_concurrent_jobs=None,
walltime=None,
walltime_delta=None,
total_walltime={},
output_size=None,
destination_user_concurrent_jobs={},
destination_total_concurrent_jobs={})
# Parse job limits
limits = root.find('limits')
if limits is not None:
@@ -334,76 +335,29 @@ class JobConfiguration(object, ConfiguresHandlers):
self.handler_runner_plugins[handler_id] = []
self.handler_runner_plugins[handler_id].append(plugin.get('id'))
def __parse_job_conf_legacy(self):
"""Loads the old-style job configuration from options in the galaxy config file (by default, config/galaxy.ini).
"""
log.debug('Loading job configuration from %s' % self.app.config.config_file)
# Always load local
self.runner_plugins = [dict(id='local', load='local', workers=self.app.config.local_job_queue_workers)]
def __set_default_job_conf(self):
# Run jobs locally
self.runner_plugins = [dict(id='local', load='local', workers=DEFAULT_LOCAL_WORKERS)]
# Load tasks if configured
if self.app.config.use_tasked_jobs:
self.runner_plugins.append(dict(id='tasks', load='tasks', workers=self.app.config.local_task_queue_workers))
for runner in self.app.config.start_job_runners:
self.runner_plugins.append(dict(id=runner, load=runner, workers=self.app.config.cluster_job_queue_workers))
self.runner_plugins.append(dict(id='tasks', load='tasks', workers=DEFAULT_LOCAL_WORKERS))
# Set the handlers
if self.app.application_stack.has_pool(self.app.application_stack.pools.JOB_HANDLERS):
self.default_handler_id = None
else:
for id in self.app.config.job_handlers:
self.handlers[id] = (id,)
self.handlers['default_job_handlers'] = self.app.config.default_job_handlers
self.default_handler_id = 'default_job_handlers'
# Set tool handler configs
for id, tool_handlers in self.app.config.tool_handlers.items():
self.tools[id] = list()
for handler_config in tool_handlers:
# rename the 'name' key to 'handler'
handler_config['handler'] = handler_config.pop('name')
self.tools[id].append(JobToolConfiguration(**handler_config))
# Set tool runner configs
for id, tool_runners in self.app.config.tool_runners.items():
# Might have been created in the handler parsing above
if id not in self.tools:
self.tools[id] = list()
for runner_config in tool_runners:
url = runner_config['url']
if url not in self.destinations:
# Create a new "legacy" JobDestination - it will have its URL converted to a destination params once the appropriate plugin has loaded
self.destinations[url] = (JobDestination(id=url, runner=url.split(':', 1)[0], url=url, legacy=True, converted=False),)
for tool_conf in self.tools[id]:
if tool_conf.params == runner_config.get('params', {}):
tool_conf['destination'] = url
break
else:
# There was not an existing config (from the handlers section) with the same params
# rename the 'url' key to 'destination'
runner_config['destination'] = runner_config.pop('url')
self.tools[id].append(JobToolConfiguration(**runner_config))
self.destinations[self.app.config.default_cluster_job_runner] = (JobDestination(id=self.app.config.default_cluster_job_runner,
runner=self.app.config.default_cluster_job_runner.split(':', 1)[0],
url=self.app.config.default_cluster_job_runner,
legacy=True,
converted=False),)
self.default_destination_id = self.app.config.default_cluster_job_runner
# Set the job limits
self.limits = Bunch(registered_user_concurrent_jobs=self.app.config.registered_user_job_limit,
anonymous_user_concurrent_jobs=self.app.config.anonymous_user_job_limit,
walltime=self.app.config.job_walltime,
walltime_delta=self.app.config.job_walltime_delta,
total_walltime={},
output_size=self.app.config.output_size_limit,
destination_user_concurrent_jobs={},
destination_total_concurrent_jobs={})
# FIXME: is this conditional necessary?
if not self.app.application_stack.has_pool(self.app.application_stack.pools.JOB_HANDLERS):
self.app.application_stack.register_postfork_function(self.make_self_default_handler)
# Set the destination
self.default_destination_id = 'local'
self.destinations['local'] = [JobDestination(id='local', runner='local')]
log.debug('Done loading job configuration')
@property
def is_default(self):
return self._is_default
def make_self_default_handler(self):
self.default_handler_id = self.app.config.server_name
self.handlers[self.app.config.server_name] = [self.app.config.server_name]
def get_tool_resource_xml(self, tool_id, tool_type):
""" Given a tool id, return XML elements describing parameters to
insert into job resources.
-7
View File
@@ -105,7 +105,6 @@ class Configuration(object):
self.smtp_username = kwargs.get('smtp_username', None)
self.smtp_password = kwargs.get('smtp_password', None)
self.smtp_ssl = kwargs.get('smtp_ssl', None)
self.start_job_runners = kwargs.get('start_job_runners', None)
self.email_from = kwargs.get('email_from', None)
self.nginx_upload_path = kwargs.get('nginx_upload_path', False)
self.log_actions = string_as_bool(kwargs.get('log_actions', 'False'))
@@ -123,12 +122,6 @@ class Configuration(object):
self.log_events = False
self.cloud_controller_instance = False
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 = []
# Error logging with sentry
self.sentry_dsn = kwargs.get('sentry_dsn', None)
# Where the tool shed hgweb.config file is stored - the default is the Galaxy installation directory.
-1
View File
@@ -195,7 +195,6 @@ def setup_galaxy_config(
galaxy_data_manager_data_path=galaxy_data_manager_data_path,
id_secret='changethisinproductiontoo',
job_config_file=job_config_file,
job_queue_workers=5,
job_working_directory=job_working_directory,
library_import_dir=library_import_dir,
log_destination="stdout",
-1
View File
@@ -84,7 +84,6 @@ class ToolShedTestDriver(driver_util.TestDriver):
datatype_converters_config_file='datatype_converters_conf.xml.sample',
file_path=shed_file_path,
hgweb_config_dir=hgweb_config_dir,
job_queue_workers=5,
id_secret='changethisinproductiontoo',
log_destination="stdout",
new_file_path=new_repos_path,