Merge branch 'release_20.01' into dev

This commit is contained in:
Nicola Soranzo
2020-03-26 16:59:14 +00:00
17 changed files with 213 additions and 70 deletions
+5 -1
View File
@@ -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
+16 -1
View File
@@ -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(
+1 -1
View File
@@ -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
+39 -33
View File
@@ -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):
+1 -1
View File
@@ -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?
+2 -2
View File
@@ -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,
+8 -8
View File
@@ -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):
+1 -1
View File
@@ -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
+2 -1
View File
@@ -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):
@@ -1,6 +1,9 @@
<tool id="__IMPORT_HISTORY__" name="Import History" version="0.1" tool_type="import_history">
<tool id="__IMPORT_HISTORY__" name="Import History" version="0.1" tool_type="import_history" profile="16.04">
<type class="ImportHistoryTool" module="galaxy.tools"/>
<action module="galaxy.tools.actions.history_imp_exp" class="ImportHistoryToolAction"/>
<requirements>
<requirement type="package" version="2.23.0">requests</requirement>
</requirements>
<command>#from base64 import b64encode#
python '$__tool_directory__/unpack_tar_gz_archive.py'
'${ b64encode(str($__ARCHIVE_SOURCE__).encode('utf-8')).decode('utf-8')}'
@@ -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)
+5 -4
View File
@@ -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
@@ -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):
+4 -1
View File
@@ -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)
+1
View File
@@ -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):
+33
View File
@@ -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)
@@ -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"])