diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index 8bf847bac50..1fb54543705 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -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) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 34244dc15a6..43728b8af06 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -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. diff --git a/lib/galaxy/webapps/tool_shed/config.py b/lib/galaxy/webapps/tool_shed/config.py index 0ff3e2f3061..6a46410542f 100644 --- a/lib/galaxy/webapps/tool_shed/config.py +++ b/lib/galaxy/webapps/tool_shed/config.py @@ -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. diff --git a/test/base/driver_util.py b/test/base/driver_util.py index e99844eedf2..b9c83a50d28 100644 --- a/test/base/driver_util.py +++ b/test/base/driver_util.py @@ -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", diff --git a/test/shed_functional/functional_tests.py b/test/shed_functional/functional_tests.py index 1805acca543..3d655732272 100644 --- a/test/shed_functional/functional_tests.py +++ b/test/shed_functional/functional_tests.py @@ -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,