Always use send_control_task via queue_worker

This commit is contained in:
mvdbeek
2019-08-18 16:01:32 +02:00
parent 30027a094b
commit f3276bf1db
8 changed files with 37 additions and 39 deletions
+5 -8
View File
@@ -1,9 +1,6 @@
from os.path import dirname
from galaxy.queue_worker import (
job_rule_modules,
send_control_task,
)
from galaxy.queue_worker import job_rule_modules
from galaxy.tools.toolbox.watcher import (
get_tool_conf_watcher,
get_tool_watcher,
@@ -27,11 +24,11 @@ class ConfigWatchers(object):
# If the reload_data_managers callback wins, the cache will miss the tools that had been removed from the cache
# and will be blind to further changes in these tools.
self.tool_config_watcher = get_tool_conf_watcher(
reload_callback=lambda: send_control_task(self.app, 'reload_toolbox'),
reload_callback=lambda: self.app.queue_worker.send_control_task('reload_toolbox'),
tool_cache=self.app.tool_cache,
)
self.data_manager_config_watcher = get_tool_conf_watcher(
reload_callback=lambda: send_control_task(self.app, 'reload_data_managers'),
reload_callback=lambda: self.app.queue_worker.send_control_task('reload_data_managers'),
)
self.tool_data_watcher = get_watcher(self.app.config, 'watch_tool_data_dir', monitor_what_str='data tables')
self.tool_watcher = get_tool_watcher(self, app.config)
@@ -63,14 +60,14 @@ class ConfigWatchers(object):
for tool_data_path in self.tool_data_paths:
self.tool_data_watcher.watch_directory(
tool_data_path,
callback=lambda path: send_control_task(self.app, 'reload_tool_data_tables', kwargs={'path': path}),
callback=lambda path: self.app.queue_worker.send_control_task('reload_tool_data_tables', kwargs={'path': path}),
require_extensions=('.loc',),
recursive=True,
)
for job_rules_directory in self.job_rules_paths:
self.job_rule_watcher.watch_directory(
job_rules_directory,
callback=lambda: send_control_task(self.app, 'reload_job_rules'),
callback=lambda: self.app.queue_worker.send_control_task('reload_job_rules'),
recursive=True,
ignore_extensions=('.pyc', '.pyo', '.pyd'))
self.active = True
+10 -10
View File
@@ -6,9 +6,6 @@ import os
from six import string_types
from galaxy import util
from galaxy.queue_worker import (
send_control_task
)
from galaxy.tools.data import TabularToolDataTable
from galaxy.util.odict import odict
from galaxy.util.template import fill_template
@@ -334,10 +331,11 @@ class DataManager(object):
self.process_move(data_table_name, name, output_ref_values[name].extra_files_path, **data_table_value)
data_table_value[name] = self.process_value_translation(data_table_name, name, **data_table_value)
data_table.add_entry(data_table_value, persist=True, entry_source=self)
send_control_task(self.data_managers.app,
'reload_tool_data_tables',
noop_self=True,
kwargs={'table_name': data_table_name})
self.data_managers.app.queue_worker.send_control_task(
'reload_tool_data_tables',
noop_self=True,
kwargs={'table_name': data_table_name}
)
if self.undeclared_tables and data_tables_dict:
# We handle the data move, by just moving all the data out of the extra files path
# moving a directory and the target already exists, we move the contents instead
@@ -356,9 +354,11 @@ class DataManager(object):
if name in path_column_names:
data_table_value[name] = os.path.abspath(os.path.join(self.data_managers.app.config.galaxy_data_manager_data_path, value))
data_table.add_entry(data_table_value, persist=True, entry_source=self)
send_control_task(self.data_managers.app, 'reload_tool_data_tables',
noop_self=True,
kwargs={'table_name': data_table_name})
self.data_managers.app.queue_worker.send_control_task(
'reload_tool_data_tables',
noop_self=True,
kwargs={'table_name': data_table_name}
)
else:
for data_table_name, data_table_values in data_tables_dict.items():
# tool returned extra data table entries, but data table was not declared in data manager
@@ -7,7 +7,6 @@ import logging
import os
from galaxy.managers import configuration, users
from galaxy.queue_worker import send_control_task
from galaxy.web import (
expose_api,
expose_api_anonymous_and_sessionless,
@@ -132,7 +131,7 @@ class ConfigurationController(BaseAPIController):
PUT /api/configuration/toolbox
Reload the Galaxy toolbox (but not individual tools).
"""
send_control_task(self.app.toolbox.app, 'reload_toolbox')
self.app.queue_worker.send_control_task('reload_toolbox')
def _tool_conf_to_dict(conf):
@@ -3,7 +3,6 @@ API operations on annotations.
"""
import logging
from galaxy import queue_worker
from galaxy.web import legacy_expose_api, require_admin
from galaxy.webapps.base.controller import BaseAPIController
@@ -45,10 +44,11 @@ class DisplayApplicationsController(BaseAPIController):
:type ids: list
"""
ids = payload.get('ids')
queue_worker.send_control_task(trans.app,
trans.app.queue_worker.send_control_task(
'reload_display_application',
noop_self=True,
kwargs={'display_application_ids': ids})
kwargs={'display_application_ids': ids}
)
reloaded, failed = trans.app.datatypes_registry.reload_display_applications(ids)
if not reloaded and failed:
message = 'Unable to reload any of the %i requested display applications ("%s").' % (len(failed), '", "'.join(failed))
+10 -7
View File
@@ -1,6 +1,5 @@
import os
import galaxy.queue_worker
from galaxy import (
exceptions,
web
@@ -42,9 +41,11 @@ class ToolData(BaseAPIController):
decoded_tool_data_id = id
data_table = trans.app.tool_data_tables.data_tables.get(decoded_tool_data_id)
data_table.reload_from_files()
galaxy.queue_worker.send_control_task(trans.app, 'reload_tool_data_tables',
noop_self=True,
kwargs={'table_name': decoded_tool_data_id})
trans.app.queue_worker.send_control_task(
'reload_tool_data_tables',
noop_self=True,
kwargs={'table_name': decoded_tool_data_id}
)
return self._data_table(decoded_tool_data_id).to_dict(view='element')
@web.require_admin
@@ -87,9 +88,11 @@ class ToolData(BaseAPIController):
return "Invalid data table item ( %s ) specified. Wrong number of columns (%s given, %s required)." % (str(values), str(len(split_values)), str(len(data_table.get_column_name_list())))
data_table.remove_entry(split_values)
galaxy.queue_worker.send_control_task(trans.app, 'reload_tool_data_tables',
noop_self=True,
kwargs={'table_name': decoded_tool_data_id})
trans.app.queue_worker.send_control_task(
'reload_tool_data_tables',
noop_self=True,
kwargs={'table_name': decoded_tool_data_id}
)
return self._data_table(decoded_tool_data_id).to_dict(view='element')
@web.require_admin
+1 -2
View File
@@ -2,7 +2,6 @@ import logging
import os
from json import dumps, loads
import galaxy.queue_worker
from galaxy import exceptions, managers, util, web
from galaxy.managers.collections_util import dictify_dataset_collection_instance
from galaxy.tools import global_tool_errors
@@ -220,7 +219,7 @@ class ToolsController(BaseAPIController, UsesVisualizationMixin):
GET /api/tools/{tool_id}/reload
Reload specified tool.
"""
galaxy.queue_worker.send_control_task(trans.app, 'reload_tool', noop_self=True, kwargs={'tool_id': id})
trans.app.queue_worker.send_control_task('reload_tool', noop_self=True, kwargs={'tool_id': id})
message, status = trans.app.toolbox.reload_tool_by_id(id)
if status == 'error':
raise exceptions.MessageException(message)
@@ -8,7 +8,6 @@ from string import punctuation as PUNCTUATION
import six
from sqlalchemy import and_, false, or_
import galaxy.queue_worker
from galaxy import (
model,
util,
@@ -1609,7 +1608,7 @@ class AdminGalaxy(controller.JSAppLauncher, AdminActions, UsesQuotaMixin, QuotaP
def jobs_control(self, trans, job_lock=None, **kwd):
if job_lock is not None:
job_lock = True if job_lock == 'true' else False
galaxy.queue_worker.send_control_task(trans.app, 'admin_job_lock', kwargs={'job_lock': job_lock}, get_response=True)
trans.app.queue_worker.send_control_task('admin_job_lock', kwargs={'job_lock': job_lock}, get_response=True)
job_lock = trans.app.job_manager.job_lock
return {'job_lock': job_lock}
@@ -1757,7 +1756,7 @@ class AdminGalaxy(controller.JSAppLauncher, AdminActions, UsesQuotaMixin, QuotaP
new_whitelist = sorted([tid for tid in tools_to_whitelist if tid in trans.app.toolbox.tools_by_id])
f.write("\n".join(new_whitelist))
trans.app.config.sanitize_whitelist = new_whitelist
galaxy.queue_worker.send_control_task(trans.app, 'reload_sanitize_whitelist', noop_self=True)
trans.app.queue_worker.send_control_task('reload_sanitize_whitelist', noop_self=True)
# dispatch a message to reload list for other processes
return trans.fill_template('/webapps/galaxy/admin/sanitize_whitelist.mako',
sanitize_all=trans.app.config.sanitize_all_html,
@@ -4,7 +4,6 @@ from json import loads
import paste.httpexceptions
from six import string_types
import galaxy.queue_worker
from galaxy import web
from galaxy.util import nice_size, unicodify
from galaxy.webapps.base.controller import BaseUIController
@@ -174,9 +173,11 @@ class DataManager(BaseUIController):
table_name = table_name.split(",")
# Reload the tool data tables
table_names = self.app.tool_data_tables.reload_tables(table_names=table_name)
galaxy.queue_worker.send_control_task(trans.app, 'reload_tool_data_tables',
noop_self=True,
kwargs={'table_name': table_name})
trans.app.queue_worker.send_control_task(
'reload_tool_data_tables',
noop_self=True,
kwargs={'table_name': table_name}
)
data = None
if table_names:
message = "Reloaded data table%s '%s'." % ('s'[len(table_names) == 1:],