From 0488972359ffb66d1ac11be5e93649fd79e5be2e Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 6 Mar 2020 10:48:28 +0100 Subject: [PATCH 01/16] Set status code even if creating a html message body That makes it way easier to understand that for instance an exported history has not been shared and so the archive can't be downloaded. --- lib/galaxy/web/framework/webapp.py | 1 + 1 file changed, 1 insertion(+) diff --git a/lib/galaxy/web/framework/webapp.py b/lib/galaxy/web/framework/webapp.py index 273eb40cc69..d77b230a275 100644 --- a/lib/galaxy/web/framework/webapp.py +++ b/lib/galaxy/web/framework/webapp.py @@ -106,6 +106,7 @@ class WebApplication(base.WebApplication): if isinstance(e, MessageException): # In the case of a controller exception, sanitize to make sure # unsafe html input isn't reflected back to the user + trans.response.status = e.status_code return trans.show_message(sanitize_html(e.err_msg), e.type) def make_body_iterable(self, trans, body): From 2d8ab3bc6b6da44535bb9faaf63efd2d1cdfe2d3 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 6 Mar 2020 10:51:15 +0100 Subject: [PATCH 02/16] Use exit to determine failure, let exceptions bubble up --- .../imp_exp/imp_history_from_archive.xml | 5 ++- .../tools/imp_exp/unpack_tar_gz_archive.py | 35 ++++++++----------- 2 files changed, 18 insertions(+), 22 deletions(-) diff --git a/lib/galaxy/tools/imp_exp/imp_history_from_archive.xml b/lib/galaxy/tools/imp_exp/imp_history_from_archive.xml index 94e165c7d5a..c749b6115f6 100644 --- a/lib/galaxy/tools/imp_exp/imp_history_from_archive.xml +++ b/lib/galaxy/tools/imp_exp/imp_history_from_archive.xml @@ -1,6 +1,9 @@ - + + + requests + #from base64 import b64encode# python '$__tool_directory__/unpack_tar_gz_archive.py' '${ b64encode(str($__ARCHIVE_SOURCE__).encode('utf-8')).decode('utf-8')}' diff --git a/lib/galaxy/tools/imp_exp/unpack_tar_gz_archive.py b/lib/galaxy/tools/imp_exp/unpack_tar_gz_archive.py index f1a8d1daf4b..5611fe6d2be 100644 --- a/lib/galaxy/tools/imp_exp/unpack_tar_gz_archive.py +++ b/lib/galaxy/tools/imp_exp/unpack_tar_gz_archive.py @@ -10,7 +10,6 @@ from __future__ import print_function import math import optparse import os -import sys import tarfile import tempfile from base64 import b64decode @@ -26,22 +25,19 @@ def url_to_file(url, dest_file): """ Transfer a file from a remote URL to a temporary file. """ - try: - url_reader = requests.get(url, stream=True) - CHUNK = 10 * 1024 # 10k - total = 0 - fp = open(dest_file, 'wb') - for chunk in url_reader.iter_content(chunk_size=CHUNK): - if chunk: - fp.write(chunk) - total += CHUNK - if total > MAX_SIZE: - break - fp.close() - return dest_file - except Exception as e: - print("Exception getting file from URL: %s" % e, file=sys.stderr) - return None + url_reader = requests.get(url, stream=True) + assert url_reader.ok, "History import failed, server returned '%s'" % url_reader.reason + CHUNK = 10 * 1024 # 10k + total = 0 + fp = open(dest_file, 'wb') + for chunk in url_reader.iter_content(chunk_size=CHUNK): + if chunk: + fp.write(chunk) + total += CHUNK + if total > MAX_SIZE: + break + fp.close() + return dest_file def check_archive(archive_file, dest_dir): @@ -96,7 +92,4 @@ if __name__ == "__main__": parser.add_option('-F', '--file', dest='is_file', action="store_true", help='Source is a file.') parser.add_option('-e', '--encoded', dest='is_b64encoded', action="store_true", default=False, help='Source and destination dir values are base64 encoded.') (options, args) = parser.parse_args() - try: - main(options, args) - except Exception as e: - print("Error unpacking tar/gz archive: %s" % e, file=sys.stderr) + main(options, args) From 3d5ec396db64cc895185a483093f3925c3acfe88 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 6 Mar 2020 11:01:34 +0100 Subject: [PATCH 03/16] Fix history export for outputs_to_working_directory --- lib/galaxy/jobs/__init__.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index f485b52e88b..71dd60e7852 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -1939,7 +1939,7 @@ class JobWrapper(HasResourceParameters): self.output_paths = [t[2] for t in results] self.output_hdas_and_paths = dict([(t[0], t[1:]) for t in results]) if special: - false_path = dataset_path_rewriter.rewrite_dataset_path(special.dataset, 'output') + false_path = dataset_path_rewriter.rewrite_dataset_path(special, 'output') dsp = DatasetPath(special.dataset.id, special.dataset.file_name, false_path) self.output_paths.append(dsp) return self.output_paths From d4e6ebe9f722709e1d47b949205878e8b969ae67 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 6 Mar 2020 16:46:47 +0100 Subject: [PATCH 04/16] Use with statement --- lib/galaxy/tools/imp_exp/unpack_tar_gz_archive.py | 15 +++++++-------- 1 file changed, 7 insertions(+), 8 deletions(-) diff --git a/lib/galaxy/tools/imp_exp/unpack_tar_gz_archive.py b/lib/galaxy/tools/imp_exp/unpack_tar_gz_archive.py index 5611fe6d2be..f3fad868169 100644 --- a/lib/galaxy/tools/imp_exp/unpack_tar_gz_archive.py +++ b/lib/galaxy/tools/imp_exp/unpack_tar_gz_archive.py @@ -29,14 +29,13 @@ def url_to_file(url, dest_file): assert url_reader.ok, "History import failed, server returned '%s'" % url_reader.reason CHUNK = 10 * 1024 # 10k total = 0 - fp = open(dest_file, 'wb') - for chunk in url_reader.iter_content(chunk_size=CHUNK): - if chunk: - fp.write(chunk) - total += CHUNK - if total > MAX_SIZE: - break - fp.close() + with open(dest_file, 'wb') as fp: + for chunk in url_reader.iter_content(chunk_size=CHUNK): + if chunk: + fp.write(chunk) + total += CHUNK + if total > MAX_SIZE: + break return dest_file From 8ef36e974b6bab9684806aa4d15ea3ea0bb1e53b Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 22 Mar 2020 16:41:38 +0100 Subject: [PATCH 05/16] Fix setting annotation that contains only integers This fixes https://sentry.galaxyproject.org/sentry/main/issues/572478/: ``` TypeError: argument cannot be of 'int' type, must be of text type File "galaxy/web/framework/decorators.py", line 157, in decorator rval = func(self, trans, *args, **kwargs) File "galaxy/webapps/galaxy/controllers/dataset.py", line 353, in set_edit annotation = sanitize_html(payload.get('annotation')) File "galaxy/util/sanitize_html.py", line 48, in sanitize_html return bleach.clean(htmlSource, **kwd) File "bleach/__init__.py", line 84, in clean return cleaner.clean(text) File "bleach/sanitizer.py", line 162, in clean raise TypeError(message) ``` --- lib/galaxy/util/sanitize_html.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/util/sanitize_html.py b/lib/galaxy/util/sanitize_html.py index b2c66b2cdb3..1ff4b8c2ec4 100644 --- a/lib/galaxy/util/sanitize_html.py +++ b/lib/galaxy/util/sanitize_html.py @@ -3,6 +3,9 @@ HTML Sanitizer (lists of acceptable_* ripped from feedparser) """ import bleach +from galaxy.util import unicodify + + _acceptable_elements = ['a', 'abbr', 'acronym', 'address', 'area', 'article', 'aside', 'audio', 'b', 'big', 'blockquote', 'br', 'button', 'canvas', 'caption', 'center', 'cite', 'code', 'col', 'colgroup', 'command', @@ -45,4 +48,4 @@ def sanitize_html(htmlSource, allow_data_urls=False): kwd = dict(tags=_acceptable_elements, attributes=_acceptable_attributes, strip=True) if allow_data_urls: kwd["protocols"] = bleach.ALLOWED_PROTOCOLS + ["data"] - return bleach.clean(htmlSource, **kwd) + return bleach.clean(unicodify(htmlSource), **kwd) From 6fffe50261082c617f4c2df59ff0ff93b0050d45 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 12 Mar 2020 16:08:20 -0400 Subject: [PATCH 06/16] Setup rabbitmq in Integration tests. --- .github/workflows/integration.yaml | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/.github/workflows/integration.yaml b/.github/workflows/integration.yaml index 3d11cf634ca..4af6a9cd727 100644 --- a/.github/workflows/integration.yaml +++ b/.github/workflows/integration.yaml @@ -2,8 +2,8 @@ name: Integration on: [push, pull_request] env: GALAXY_TEST_DBURI: 'postgres://postgres:postgres@localhost:5432/galaxy?client_encoding=utf8' + GALAXY_TEST_AMQP_URL: 'amqp://rabbitmq:5672//' jobs: - test: name: Test runs-on: ubuntu-18.04 @@ -20,6 +20,10 @@ jobs: POSTGRES_DB: postgres ports: - 5432:5432 + rabbitmq: + image: rabbitmq + ports: + - 5672:5672 steps: - name: Prune unused docker image, volumes and containers run: docker system prune -a -f From 08f1a705b17fb788b9bce53f891a4ea6427fcb27 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 12 Mar 2020 13:43:07 -0400 Subject: [PATCH 07/16] Allow embedded Pulsar to use a MQ. --- lib/galaxy/jobs/runners/pulsar.py | 72 +++++++++++++++++-------------- 1 file changed, 39 insertions(+), 33 deletions(-) diff --git a/lib/galaxy/jobs/runners/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py index a13935f403b..ebd72250e58 100644 --- a/lib/galaxy/jobs/runners/pulsar.py +++ b/lib/galaxy/jobs/runners/pulsar.py @@ -54,6 +54,7 @@ __all__ = ( 'PulsarRESTJobRunner', 'PulsarMQJobRunner', 'PulsarEmbeddedJobRunner', + 'PulsarEmbeddedMQJobRunner', ) MINIMUM_PULSAR_VERSIONS = { @@ -187,6 +188,7 @@ class PulsarJobRunner(AsynchronousJobRunner): runner_name = "PulsarJobRunner" default_build_pulsar_app = False + use_mq = False def __init__(self, app, nworkers, **kwds): """Start the job runner.""" @@ -202,8 +204,15 @@ class PulsarJobRunner(AsynchronousJobRunner): self._monitor() def _monitor(self): - # Extension point allow MQ variant to setup callback instead - self._init_monitor_thread() + if self.use_mq: + # This is a message queue driven runner, don't monitor + # just setup required callback. + self._init_noop_monitor() + + self.client_manager.ensure_has_status_update_callback(self.__async_update) + self.client_manager.ensure_has_ack_consumers() + else: + self._init_monitor_thread() def __init_client_manager(self): pulsar_conf = self.runner_params.get('pulsar_app_config', None) @@ -813,36 +822,6 @@ class PulsarJobRunner(AsynchronousJobRunner): metadata_kwds['datatypes_config'] = datatypes_config return metadata_kwds - -class PulsarLegacyJobRunner(PulsarJobRunner): - """Flavor of Pulsar job runner mimicking behavior of old LWR runner.""" - - destination_defaults = dict( - rewrite_parameters="false", - dependency_resolution="local", - ) - - -class PulsarMQJobRunner(PulsarJobRunner): - """Flavor of Pulsar job runner with sensible defaults for message queue communication.""" - - destination_defaults = dict( - default_file_action="remote_transfer", - rewrite_parameters="true", - dependency_resolution="remote", - jobs_directory=PARAMETER_SPECIFICATION_REQUIRED, - url=PARAMETER_SPECIFICATION_IGNORED, - private_token=PARAMETER_SPECIFICATION_IGNORED - ) - - def _monitor(self): - # This is a message queue driven runner, don't monitor - # just setup required callback. - self._init_noop_monitor() - - self.client_manager.ensure_has_status_update_callback(self.__async_update) - self.client_manager.ensure_has_ack_consumers() - def __async_update(self, full_status): galaxy_job_id = None try: @@ -862,6 +841,29 @@ class PulsarMQJobRunner(PulsarJobRunner): # Nothing else to do? - Attempt to fail the job? +class PulsarLegacyJobRunner(PulsarJobRunner): + """Flavor of Pulsar job runner mimicking behavior of old LWR runner.""" + + destination_defaults = dict( + rewrite_parameters="false", + dependency_resolution="local", + ) + + +class PulsarMQJobRunner(PulsarJobRunner): + """Flavor of Pulsar job runner with sensible defaults for message queue communication.""" + use_mq = True + + destination_defaults = dict( + default_file_action="remote_transfer", + rewrite_parameters="true", + dependency_resolution="remote", + jobs_directory=PARAMETER_SPECIFICATION_REQUIRED, + url=PARAMETER_SPECIFICATION_IGNORED, + private_token=PARAMETER_SPECIFICATION_IGNORED + ) + + KUBERNETES_DESTINATION_DEFAULTS = { "default_file_action": "remote_transfer", "rewrite_parameters": "true", @@ -901,7 +903,7 @@ class PulsarRESTJobRunner(PulsarJobRunner): class PulsarEmbeddedJobRunner(PulsarJobRunner): - """Flavor of Puslar job runnner that runs Pulsar's server code directly within Galaxy. + """Flavor of Puslar job runner that runs Pulsar's server code directly within Galaxy. This is an appropriate job runner for when the desire is to use Pulsar staging but their is not need to run a remote service. @@ -915,6 +917,10 @@ class PulsarEmbeddedJobRunner(PulsarJobRunner): default_build_pulsar_app = True +class PulsarEmbeddedMQJobRunner(PulsarMQJobRunner): + default_build_pulsar_app = True + + class PulsarComputeEnvironment(ComputeEnvironment): def __init__(self, pulsar_client, job_wrapper, remote_job_config): From 80138afedc991c5c466c2390cf0823861701a627 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 12 Mar 2020 13:52:11 -0400 Subject: [PATCH 08/16] Integration test for embedded Pulsar in MQ mode. --- .github/workflows/integration.yaml | 2 +- test/integration/test_pulsar_embedded_mq.py | 80 +++++++++++++++++++++ 2 files changed, 81 insertions(+), 1 deletion(-) create mode 100644 test/integration/test_pulsar_embedded_mq.py diff --git a/.github/workflows/integration.yaml b/.github/workflows/integration.yaml index 4af6a9cd727..25edfff1b8c 100644 --- a/.github/workflows/integration.yaml +++ b/.github/workflows/integration.yaml @@ -2,7 +2,7 @@ name: Integration on: [push, pull_request] env: GALAXY_TEST_DBURI: 'postgres://postgres:postgres@localhost:5432/galaxy?client_encoding=utf8' - GALAXY_TEST_AMQP_URL: 'amqp://rabbitmq:5672//' + GALAXY_TEST_AMQP_URL: 'amqp://localhost:5672//' jobs: test: name: Test diff --git a/test/integration/test_pulsar_embedded_mq.py b/test/integration/test_pulsar_embedded_mq.py new file mode 100644 index 00000000000..f6f043bd0b7 --- /dev/null +++ b/test/integration/test_pulsar_embedded_mq.py @@ -0,0 +1,80 @@ +"""Integration tests for the Pulsar embedded runner with outputs to working directory.""" + +import os +import string +import tempfile + +import pytest + +from galaxy.util import safe_makedirs +from galaxy_test.driver import integration_util + +SCRIPT_DIRECTORY = os.path.abspath(os.path.dirname(__file__)) +EMBEDDED_PULSAR_JOB_CONFIG_FILE = os.path.join(SCRIPT_DIRECTORY, "embedded_pulsar_mq_job_conf.yml") +AMQP_URL = os.environ.get("GALAXY_TEST_AMQP_URL", "amqp://guest:guest@localhost:5672//") + +JOB_CONF_TEMPLATE = """ +runners: + local: + load: galaxy.jobs.runners.local:LocalJobRunner + workers: 1 + pulsar: + load: galaxy.jobs.runners.pulsar:PulsarEmbeddedMQJobRunner + pulsar_app_config: + tool_dependency_dir: none + conda_auto_init: false + conda_auto_install: false + message_queue_url: ${amqp_url} + staging_directory: ${jobs_directory} + amqp_url: ${amqp_url} +execution: + default: pulsar_mq_environment + environments: + pulsar_mq_environment: + runner: pulsar + rewrite_parameters: true + dependency_resolution: none + default_file_action: remote_transfer + jobs_directory: ${jobs_directory} + local_environment: + runner: local +tools: + - id: upload1 + environment: local_environment +""" + + +class EmbeddedMessageQueuePulsarIntegrationInstance(integration_util.IntegrationInstance): + """Describe a Galaxy test instance with embedded pulsar configured. + + $ Setup RabbitMQ (e.g. https://www.rabbitmq.com/install-homebrew.html) + $ GALAXY_TEST_AMQP_URL='amqp://guest:guest@localhost:5672//' pytest -s test/integration/test_pulsar_embedded_mq.py + """ + + framework_tool_and_types = True + # Test leverages $UWSGI_PORT in job code, need to set this up. + require_uwsgi = True + + @classmethod + def handle_galaxy_config_kwds(cls, config): + amqp_url = os.environ.get("GALAXY_TEST_AMQP_URL", None) + if amqp_url is None: + pytest.skip("External AMQP URL not configured for test") + + jobs_directory = os.path.join(cls._test_driver.mkdtemp(), "pulsar_staging") + safe_makedirs(jobs_directory) + job_conf_template = string.Template(JOB_CONF_TEMPLATE) + job_conf_str = job_conf_template.substitute( + amqp_url=AMQP_URL, + jobs_directory=jobs_directory, + ) + with tempfile.NamedTemporaryFile(suffix="_mq_job_conf.yml", mode="w", delete=False) as job_conf: + job_conf.write(job_conf_str) + config["job_config_file"] = job_conf.name + infrastructure_url = "http://localhost:$UWSGI_PORT" + config["galaxy_infrastructure_url"] = infrastructure_url + + +instance = integration_util.integration_module_instance(EmbeddedMessageQueuePulsarIntegrationInstance) + +test_tools = integration_util.integration_tool_runner(["simple_constructs"]) From 04027a89966f348c04f5a318283d1c3958f08066 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 25 Mar 2020 13:44:49 +0100 Subject: [PATCH 09/16] Hide LDDA when creating collection in library interface Fixes https://github.com/galaxyproject/galaxy/issues/9435 --- lib/galaxy/managers/collections.py | 2 +- lib/galaxy/model/__init__.py | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/managers/collections.py b/lib/galaxy/managers/collections.py index 08cef87405b..91b137c6ea7 100644 --- a/lib/galaxy/managers/collections.py +++ b/lib/galaxy/managers/collections.py @@ -434,7 +434,7 @@ class DatasetCollectionManager(object): self.tag_handler.apply_item_tags(user=trans.user, item=element, tags_str=tag_str) elif src_type == 'ldda': element = self.ldda_manager.get(trans, encoded_id, check_accessible=True) - element = element.to_history_dataset_association(trans.history, add_to_history=True) + element = element.to_history_dataset_association(trans.history, add_to_history=True, visible=not hide_source_items) self.tag_handler.apply_item_tags(user=trans.user, item=element, tags_str=tag_str) elif src_type == 'hdca': # TODO: Option to copy? Force copy? Copy or allow if not owned? diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 621124e9ccd..d51b38cc251 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -3528,7 +3528,7 @@ class LibraryDatasetDatasetAssociation(DatasetInstance, HasName, RepresentById): self.library_dataset = library_dataset self.user = user - def to_history_dataset_association(self, target_history, parent_id=None, add_to_history=False): + def to_history_dataset_association(self, target_history, parent_id=None, add_to_history=False, visible=None): sa_session = object_session(self) hda = HistoryDatasetAssociation(name=self.name, info=self.info, @@ -3538,7 +3538,7 @@ class LibraryDatasetDatasetAssociation(DatasetInstance, HasName, RepresentById): extension=self.extension, dbkey=self.dbkey, dataset=self.dataset, - visible=self.visible, + visible=visible if visible is not None else self.visible, deleted=self.deleted, parent_id=parent_id, copied_from_library_dataset_dataset_association=self, From 3b5a1a03d2f4b38c1ef9d0f8bc5955441b5fba51 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 25 Mar 2020 15:57:50 +0100 Subject: [PATCH 10/16] Add API tests library datasets as HDCA import --- lib/galaxy_test/api/test_libraries.py | 33 +++++++++++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/lib/galaxy_test/api/test_libraries.py b/lib/galaxy_test/api/test_libraries.py index 2bfe863c01d..8f0edbe60b7 100644 --- a/lib/galaxy_test/api/test_libraries.py +++ b/lib/galaxy_test/api/test_libraries.py @@ -319,6 +319,39 @@ class LibrariesApiTestCase(ApiTestCase, TestsDatasets): self._assert_has_keys(create_response.json(), "file_ext") assert create_response.json()["file_ext"] == "txt" + def test_ldda_collection_import_to_history(self): + self._import_to_history(visible=True) + + def test_ldda_collection_import_to_history_hide_source(self): + self._import_to_history(visible=False) + + def _import_to_history(self, visible=True): + ld = self._create_dataset_in_folder_in_library("ForHistoryImport").json() + history_id = self.dataset_populator.new_history() + url = "histories/%s/contents" % history_id + collection_name = 'new_collection_name' + element_identifer = 'new_element_identifier' + payload = { + "collection_type": "list", + "history_content_type": "dataset_collection", + "model_class": "HistoryDatasetCollectionAssociation", + "history_id": history_id, + "name": collection_name, + "hide_source_items": not visible, + "element_identifiers": json.dumps([{ + "id": ld['id'], + "name": element_identifer, + "src": "ldda"}]), + "type": "dataset_collection", + "elements": [] + } + new_collection = self._post(url, payload).json() + assert new_collection['name'] == collection_name + assert new_collection['element_count'] == 1 + element = new_collection['elements'][0] + assert element['element_identifier'] == element_identifer + assert element['object']['visible'] == visible + def test_create_datasets_in_library_from_collection(self): library = self.library_populator.new_private_library("ForCreateDatasetsFromCollection") folder_response = self._create_folder(library) From 0ca974c0d5aaf922f9920caa16fa473ddb864a0e Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 4 Mar 2020 12:04:18 +0100 Subject: [PATCH 11/16] Save integrated tool panel once on toolbox reload This should make toolbox reloads work on main and other large instances. --- lib/galaxy/config_watchers.py | 13 ++++++++++++- lib/galaxy/queue_worker.py | 6 +++--- lib/galaxy/tools/toolbox/base.py | 5 +++-- 3 files changed, 18 insertions(+), 6 deletions(-) diff --git a/lib/galaxy/config_watchers.py b/lib/galaxy/config_watchers.py index 23db380d6a2..ddd1036644f 100644 --- a/lib/galaxy/config_watchers.py +++ b/lib/galaxy/config_watchers.py @@ -1,3 +1,4 @@ +import logging from os.path import dirname from galaxy.queue_worker import job_rule_modules @@ -7,6 +8,8 @@ from galaxy.tools.toolbox.watcher import ( ) from galaxy.util.watcher import get_watcher +log = logging.getLogger(__name__) + class ConfigWatchers(object): """Contains ToolConfWatcher, ToolWatcher and ToolDataWatcher objects.""" @@ -21,8 +24,16 @@ class ConfigWatchers(object): # If there are multiple ToolConfWatcher objects for the same handler or web process a race condition occurs between the two cache_cleanup functions. # 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. + + def reload_toolbox(): + try: + self.app.queue_worker.send_local_control_task('reload_toolbox', get_response=True), + except Exception: + log.exception() + self.app.queue_worker.send_control_task('reload_toolbox', kwargs={'save_integrated_tool_panel': False}), + self.tool_config_watcher = get_tool_conf_watcher( - reload_callback=lambda: self.app.queue_worker.send_control_task('reload_toolbox'), + reload_callback=reload_toolbox, tool_cache=self.app.tool_cache, ) self.data_manager_config_watcher = get_tool_conf_watcher( diff --git a/lib/galaxy/queue_worker.py b/lib/galaxy/queue_worker.py index 59335726060..1c101bc2eff 100644 --- a/lib/galaxy/queue_worker.py +++ b/lib/galaxy/queue_worker.py @@ -42,7 +42,7 @@ def send_local_control_task(app, task, get_response=False, kwargs=None): """ if kwargs is None: kwargs = {} - log.info("Queuing async task %s for %s." % (task, app.config.server_name)) + log.info("Queuing %s task %s for %s." % ("sync" if get_response else "async", task, app.config.server_name)) payload = {'task': task, 'kwargs': kwargs} routing_key = 'control.%s@%s' % (app.config.server_name, socket.gethostname()) @@ -352,8 +352,8 @@ class GalaxyQueueWorker(ConsumerProducerMixin, threading.Thread): def send_control_task(self, task, noop_self=False, get_response=False, routing_key='control.*', kwargs=None): return send_control_task(app=self.app, task=task, noop_self=noop_self, get_response=get_response, routing_key=routing_key, kwargs=kwargs) - def send_local_control_task(self, task, kwargs=None): - return send_local_control_task(app=self.app, task=task, kwargs=kwargs) + def send_local_control_task(self, task, get_response=False, kwargs=None): + return send_local_control_task(app=self.app, get_response=get_response, task=task, kwargs=kwargs) @property def declare_queues(self): diff --git a/lib/galaxy/tools/toolbox/base.py b/lib/galaxy/tools/toolbox/base.py index 60ed3c0e925..f0af4bcbd97 100644 --- a/lib/galaxy/tools/toolbox/base.py +++ b/lib/galaxy/tools/toolbox/base.py @@ -75,7 +75,7 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin): workflows optionally in labelled sections. """ - def __init__(self, config_filenames, tool_root_dir, app): + def __init__(self, config_filenames, tool_root_dir, app, save_integrated_tool_panel=True): """ Create a toolbox from the config files named by `config_filenames`, using `tool_root_dir` as the base directory for finding individual tool config files. @@ -119,7 +119,8 @@ class AbstractToolBox(Dictifiable, ManagesIntegratedToolPanelMixin): if self.app.name == 'galaxy' and self._integrated_tool_panel_config_has_contents: # Load self._tool_panel based on the order in self._integrated_tool_panel. self._load_tool_panel() - self._save_integrated_tool_panel() + if save_integrated_tool_panel: + self._save_integrated_tool_panel() def create_tool(self, config_file, tool_shed_repository=None, guid=None, **kwds): raise NotImplementedError() From 4a2ae58966c6bafc6f03ad9aa91153fce9a31c5c Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 4 Mar 2020 12:45:27 +0100 Subject: [PATCH 12/16] Don't run toolbox reload twice on watcher --- lib/galaxy/config_watchers.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/config_watchers.py b/lib/galaxy/config_watchers.py index ddd1036644f..2e16735ee7f 100644 --- a/lib/galaxy/config_watchers.py +++ b/lib/galaxy/config_watchers.py @@ -26,11 +26,15 @@ class ConfigWatchers(object): # and will be blind to further changes in these tools. def reload_toolbox(): + save_integrated_tool_panel = False try: + # Run and wait for toolbox reload on the process that watches the config files. + # The toolbpox reload will update the integrated_tool_panel_file self.app.queue_worker.send_local_control_task('reload_toolbox', get_response=True), except Exception: - log.exception() - self.app.queue_worker.send_control_task('reload_toolbox', kwargs={'save_integrated_tool_panel': False}), + save_integrated_tool_panel = True + log.exception("Exception occured while reloading toolbox") + self.app.queue_worker.send_control_task('reload_toolbox', noop_self=True, kwargs={'save_integrated_tool_panel': save_integrated_tool_panel}), self.tool_config_watcher = get_tool_conf_watcher( reload_callback=reload_toolbox, From a1ffcb6b30134f8a580bf9e4b4484465f429fda4 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 4 Mar 2020 12:47:58 +0100 Subject: [PATCH 13/16] Log when writing tool panel config file --- lib/galaxy/tools/toolbox/integrated_panel.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/lib/galaxy/tools/toolbox/integrated_panel.py b/lib/galaxy/tools/toolbox/integrated_panel.py index 9e3020282bd..fa1b7949b5e 100644 --- a/lib/galaxy/tools/toolbox/integrated_panel.py +++ b/lib/galaxy/tools/toolbox/integrated_panel.py @@ -1,3 +1,4 @@ +import logging import os import shutil import string @@ -11,6 +12,8 @@ from .panel import ( ToolPanelElements ) +log = logging.getLogger(__name__) + INTEGRATED_TOOL_PANEL_DESCRIPTION = """ This is Galaxy's integrated tool panel and should be modified directly only for reordering tools inside a section. Each time Galaxy starts up, this file is @@ -52,6 +55,7 @@ class ManagesIntegratedToolPanelMixin(object): use this file to manage the tool panel, we'll not use xml_to_string() since it doesn't write XML quite right. """ destination = os.path.abspath(self._integrated_tool_panel_config) + log.debug("Writing integrated tool panel config file to '%s'", destination) tracking_directory = self._integrated_tool_panel_tracking_directory if tracking_directory: if not os.path.exists(tracking_directory): From 173d19c23a7e027700488b8cc0deee25e3fcc006 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 4 Mar 2020 13:24:19 +0100 Subject: [PATCH 14/16] Pass save_integrated_tool_panel down the ToolBox --- lib/galaxy/queue_worker.py | 8 ++++---- lib/galaxy/tools/__init__.py | 3 ++- lib/galaxy/tools/toolbox/base.py | 4 ++-- 3 files changed, 8 insertions(+), 7 deletions(-) diff --git a/lib/galaxy/queue_worker.py b/lib/galaxy/queue_worker.py index 1c101bc2eff..da2538508e8 100644 --- a/lib/galaxy/queue_worker.py +++ b/lib/galaxy/queue_worker.py @@ -163,19 +163,19 @@ def reload_tool(app, **kwargs): log.error("Reload tool invoked without tool id.") -def reload_toolbox(app, **kwargs): +def reload_toolbox(app, save_integrated_tool_panel=True, **kwargs): reload_timer = util.ExecutionTimer() log.debug("Executing toolbox reload on '%s'", app.config.server_name) reload_count = app.toolbox._reload_count if hasattr(app, 'tool_cache'): app.tool_cache.cleanup() - _get_new_toolbox(app) + _get_new_toolbox(app, save_integrated_tool_panel) app.toolbox._reload_count = reload_count + 1 send_local_control_task(app, 'rebuild_toolbox_search_index') log.debug("Toolbox reload %s", reload_timer) -def _get_new_toolbox(app): +def _get_new_toolbox(app, save_integrated_tool_panel=True): """ Generate a new toolbox, by constructing a toolbox from the config files, and then adding pre-existing data managers from the old toolbox to the new toolbox. @@ -186,7 +186,7 @@ def _get_new_toolbox(app): app.tool_shed_repository_cache.rebuild() tool_configs = app.config.tool_configs - new_toolbox = tools.ToolBox(tool_configs, app.config.tool_path, app) + new_toolbox = tools.ToolBox(tool_configs, app.config.tool_path, app, save_integrated_tool_panel=save_integrated_tool_panel) new_toolbox.data_manager_tools = app.toolbox.data_manager_tools app.datatypes_registry.load_datatype_converters(new_toolbox, use_cached=True) app.datatypes_registry.load_external_metadata_tool(new_toolbox) diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 711b52c977b..208fe5f5df4 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -241,7 +241,7 @@ class ToolBox(BaseGalaxyToolBox): how to construct them, action types, dependency management, etc.... """ - def __init__(self, config_filenames, tool_root_dir, app): + def __init__(self, config_filenames, tool_root_dir, app, save_integrated_tool_panel=True): self._reload_count = 0 self.tool_location_fetcher = ToolLocationFetcher() # This is here to deal with the old default value, which doesn't make @@ -253,6 +253,7 @@ class ToolBox(BaseGalaxyToolBox): config_filenames=config_filenames, tool_root_dir=tool_root_dir, app=app, + save_integrated_tool_panel=save_integrated_tool_panel, ) def can_load_config_file(self, config_filename): diff --git a/lib/galaxy/tools/toolbox/base.py b/lib/galaxy/tools/toolbox/base.py index f0af4bcbd97..a19a5f91b42 100644 --- a/lib/galaxy/tools/toolbox/base.py +++ b/lib/galaxy/tools/toolbox/base.py @@ -1190,8 +1190,8 @@ class BaseGalaxyToolBox(AbstractToolBox): shouldn't really depend on. """ - def __init__(self, config_filenames, tool_root_dir, app): - super(BaseGalaxyToolBox, self).__init__(config_filenames, tool_root_dir, app) + def __init__(self, config_filenames, tool_root_dir, app, save_integrated_tool_panel=True): + super(BaseGalaxyToolBox, self).__init__(config_filenames, tool_root_dir, app, save_integrated_tool_panel) old_toolbox = getattr(app, 'toolbox', None) if old_toolbox: self.dependency_manager = old_toolbox.dependency_manager From 1604f2ab28c6757d92a347e2c8b364b8f586c383 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 4 Mar 2020 16:16:05 +0100 Subject: [PATCH 15/16] Don't block forever if connection can't be acquired --- lib/galaxy/queue_worker.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/queue_worker.py b/lib/galaxy/queue_worker.py index da2538508e8..87eeada1ecb 100644 --- a/lib/galaxy/queue_worker.py +++ b/lib/galaxy/queue_worker.py @@ -114,7 +114,7 @@ class ControlTask(object): callback_queue = [self.callback_queue] self.correlation_id = uuid() try: - with producers[self.connection].acquire(block=True) as producer: + with producers[self.connection].acquire(block=True, timeout=10) as producer: producer.publish( payload, exchange=None if local else self.exchange, From 7d110f1d71e4a00fb7a7e6cff594b45f61f52e0d Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Tue, 17 Mar 2020 17:56:21 +0100 Subject: [PATCH 16/16] Exchange routing key needs to be wildcard ... But why that wasn't a problem before is a mystery to me. --- lib/galaxy/queues.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/queues.py b/lib/galaxy/queues.py index abb9b993027..5b74c0e496d 100644 --- a/lib/galaxy/queues.py +++ b/lib/galaxy/queues.py @@ -32,7 +32,7 @@ def control_queues_from_config(config): """ hostname = socket.gethostname() process_name = "{server_name}@{hostname}".format(server_name=config.server_name, hostname=hostname) - exchange_queue = Queue("control.%s" % process_name, galaxy_exchange, routing_key='control.%s' % process_name) + exchange_queue = Queue("control.%s" % process_name, galaxy_exchange, routing_key='control.*') non_exchange_queue = Queue("control.%s" % process_name, routing_key='control.%s' % process_name) return exchange_queue, non_exchange_queue