diff --git a/.github/workflows/integration.yaml b/.github/workflows/integration.yaml index eae99cebd44..7e72077fbb2 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://localhost: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 diff --git a/lib/galaxy/config_watchers.py b/lib/galaxy/config_watchers.py index 23db380d6a2..2e16735ee7f 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,20 @@ 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(): + 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: + 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=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/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index d21bc5a2fa6..2451197b3aa 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -1958,7 +1958,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 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): 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 f8cafc59af6..a065fd5af96 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -3547,7 +3547,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, @@ -3557,7 +3557,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, diff --git a/lib/galaxy/queue_worker.py b/lib/galaxy/queue_worker.py index 78c7acc65b3..df671acdfb9 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()) @@ -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, @@ -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) @@ -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/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 diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 7665ac7fc7e..f19b350272f 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/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..f3fad868169 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,18 @@ 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') + 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 + 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 - fp.close() - return dest_file - except Exception as e: - print("Exception getting file from URL: %s" % e, file=sys.stderr) - return None + return dest_file def check_archive(archive_file, dest_dir): @@ -96,7 +91,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) diff --git a/lib/galaxy/tools/toolbox/base.py b/lib/galaxy/tools/toolbox/base.py index de683796ec7..4794e20d779 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() @@ -1197,8 +1198,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 diff --git a/lib/galaxy/tools/toolbox/integrated_panel.py b/lib/galaxy/tools/toolbox/integrated_panel.py index 29a3fb86f15..ee10251f7d7 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 @@ -12,6 +13,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 @@ -53,6 +56,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): 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) diff --git a/lib/galaxy/web/framework/webapp.py b/lib/galaxy/web/framework/webapp.py index cf9b52e7e0e..85500107426 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): 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) 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"])