Merge branch 'release_20.01' into dev

This commit is contained in:
mvdbeek
2020-04-27 15:40:13 +02:00
16 changed files with 311 additions and 95 deletions
+4 -4
View File
@@ -12,8 +12,8 @@
<filter>${$site_id.startswith( 'local_' ) or $dataset.dbkey in $site_dbkeys}</filter>
<!-- We define url and params as normal, but values defined in dynamic_param are available by specified name -->
<url>${redirect_url}</url>
<param type="data" name="bgzip_file" url="galaxy_${DATASET_HASH}.vcf.gz" format="vcf_bgzip" />
<param type="data" name="tabix_file" dataset="bgzip_file" url="galaxy_${DATASET_HASH}.vcf.gz.tbi" format="tabix" />
<param type="data" name="bgzip_file" url="galaxy_${DATASET_HASH}.vcf.gz" format="vcf_bgzip" allow_cors="true" mimetype="application/octet-stream"/>
<param type="data" name="tabix_file" metadata="tabix_index" url="galaxy_${DATASET_HASH}.vcf.gz.tbi" allow_cors="true" mimetype="application/octet-stream"/>
<param type="template" name="site_organism" strip="True" >
#if ($dataset.dbkey in $site_dbkeys)
$site_organisms[ $site_dbkeys.index( $bgzip_file.dbkey ) ]
@@ -94,8 +94,8 @@
<filter>${ $dataset.dbkey == $value }</filter>
<!-- We define url and params as normal, but values defined in dynamic_param are available by specified name -->
<url>http://www.broadinstitute.org/igv/projects/current/igv.php?sessionURL=${bgzip_file.qp}&amp;genome=${qp($bgzip_file.dbkey)}&amp;merge=true&amp;name=${qp( ( $bgzip_file.name or $DATASET_HASH ).replace( ',', ';' ) )}</url>
<param type="data" name="bgzip_file" url="galaxy_${DATASET_HASH}.vcf.gz" format="vcf_bgzip" />
<param type="data" name="tabix_file" dataset="bgzip_file" url="galaxy_${DATASET_HASH}.vcf.gz.tbi" format="tabix" />
<param type="data" name="bgzip_file" url="galaxy_${DATASET_HASH}.vcf.gz" format="vcf_bgzip" mimetype="application/octet-stream"/>
<param type="data" name="tabix_file" metadata="tabix_index" url="galaxy_${DATASET_HASH}.vcf.gz.tbi" format="vcf_bgzip" mimetype="application/octet-stream"/>
</dynamic_links>
</display>
<!-- Dan Blankenberg -->
+2 -2
View File
@@ -671,8 +671,8 @@ class GalaxyAppConfiguration(BaseAppConfiguration, CommonConfigurationMixin):
# InteractiveTools propagator mapping file
self.interactivetools_map = self._in_root_dir(kwargs.get("interactivetools_map", self._in_data_dir("interactivetools_map.sqlite")))
self.interactivetool_prefix = kwargs.get("interactivetools_prefix", "interactivetool")
self.interactivetool_proxy_host = kwargs.get("interactivetool_proxy_host", None)
self.interactivetools_prefix = kwargs.get("interactivetools_prefix", "interactivetool")
self.interactivetools_proxy_host = kwargs.get("interactivetool_proxy_host", None)
self.containers_conf = parse_containers_config(self.containers_config_file)
@@ -122,7 +122,7 @@ pbr==5.4.5
prettytable==0.7.2
prov==1.5.1
psutil==5.7.0
pulsar-galaxy-lib==0.14.0.dev1
pulsar-galaxy-lib==0.14.0.dev3
pyasn1-modules==0.2.8
pyasn1==0.4.8
pycparser==2.20
+8
View File
@@ -1141,6 +1141,14 @@ class JobWrapper(HasResourceParameters):
raise Exception('(%s) Unable to create job working directory',
job.id)
@property
def guest_ports(self):
if hasattr(self, "interactivetools"):
guest_ports = [ep.get('port') for ep in self.interactivetools]
return guest_ports
else:
return []
@property
def working_directory(self):
if self.__working_directory is None:
+1 -1
View File
@@ -425,7 +425,7 @@ class BaseJobRunner(object):
compute_tmp_directory = job_wrapper.tmp_directory()
tool = job_wrapper.tool
guest_ports = [ep.get('port') for ep in getattr(job_wrapper, 'interactivetools', [])]
guest_ports = job_wrapper.guest_ports
tool_info = ToolInfo(
tool.containers,
tool.requirements,
+10 -37
View File
@@ -19,12 +19,15 @@ from galaxy.jobs.runners import (
from galaxy.jobs.runners.util.pykube_util import (
DEFAULT_JOB_API_VERSION,
ensure_pykube,
find_job_object_by_name,
galaxy_instance_id,
Job,
job_object_dict,
Pod,
produce_unique_k8s_job_name,
pull_policy,
pykube_client_from_dict,
stop_job,
)
from galaxy.util.bytesize import ByteSize
@@ -201,25 +204,8 @@ class KubernetesJobRunner(AsynchronousJobRunner):
return None
def __get_galaxy_instance_id(self):
"""
Gets the id of the Galaxy instance. This will be added to Jobs and Pods names, so it needs to be DNS friendly,
this means: `The Internet standards (Requests for Comments) for protocols mandate that component hostname labels
may contain only the ASCII letters 'a' through 'z' (in a case-insensitive manner), the digits '0' through '9',
and the minus sign ('-').`
It looks for the value set on self.runner_params['k8s_galaxy_instance_id'], which might or not be set. The
idea behind this is to allow the Galaxy instance to trust (or not) existing k8s Jobs and Pods that match the
setup of a Job that is being recovered or restarted after a downtime/reboot.
:return:
:rtype:
"""
if "k8s_galaxy_instance_id" in self.runner_params:
if re.match(r"(?!-)[a-z\d-]{1,20}(?<!-)$", self.runner_params['k8s_galaxy_instance_id']):
return self.runner_params['k8s_galaxy_instance_id']
else:
log.error("Galaxy instance '" + self.runner_params['k8s_galaxy_instance_id'] + "' is either too long "
+ '(>20 characters) or it includes non DNS acceptable characters, ignoring it.')
return None
"""Parse the ID of the Galaxy instance from runner params."""
return galaxy_instance_id(self.runner_params)
def __produce_unique_k8s_job_name(self, galaxy_internal_job_id):
# wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers.
@@ -509,19 +495,7 @@ class KubernetesJobRunner(AsynchronousJobRunner):
def __cleanup_k8s_job(self, job):
k8s_cleanup_job = self.runner_params['k8s_cleanup_job']
job_failed = (job.obj['status']['failed'] > 0
if 'failed' in job.obj['status'] else False)
# Scale down the job just in case even if cleanup is never
job.scale(replicas=0)
if (k8s_cleanup_job == "always" or
(k8s_cleanup_job == "onsuccess" and not job_failed)):
delete_options = {
"apiVersion": "v1",
"kind": "DeleteOptions",
"propagationPolicy": "Background"
}
r = job.api.delete(json=delete_options, **job.api_kwargs())
job.api.raise_for_status(r)
stop_job(job, k8s_cleanup_job)
def __job_failed_due_to_walltime_limit(self, job):
conditions = job.obj['status'].get('conditions') or []
@@ -550,11 +524,10 @@ class KubernetesJobRunner(AsynchronousJobRunner):
"""Attempts to delete a dispatched job to the k8s cluster"""
job = job_wrapper.get_job()
try:
jobs = Job.objects(self._pykube_api).filter(
selector="app=" + self.__produce_unique_k8s_job_name(job.get_id_tag()),
namespace=self.runner_params['k8s_namespace'])
if len(jobs.response['items']) > 0:
job_to_delete = Job(self._pykube_api, jobs.response['items'][0])
name = self.__produce_unique_k8s_job_name(job.get_id_tag())
namespace = self.runner_params['k8s_namespace']
job_to_delete = find_job_object_by_name(self._pykube_api, name, namespace)
if job_to_delete:
self.__cleanup_k8s_job(job_to_delete)
# TODO assert whether job parallelism == 0
# assert not job_to_delete.exists(), "Could not delete job,"+job.job_runner_external_id+" it still exists"
+26 -4
View File
@@ -189,6 +189,7 @@ class PulsarJobRunner(AsynchronousJobRunner):
runner_name = "PulsarJobRunner"
default_build_pulsar_app = False
use_mq = False
poll = True
def __init__(self, app, nworkers, **kwds):
"""Start the job runner."""
@@ -207,12 +208,13 @@ class PulsarJobRunner(AsynchronousJobRunner):
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:
if self.poll:
self._init_monitor_thread()
else:
self._init_noop_monitor()
def __init_client_manager(self):
pulsar_conf = self.runner_params.get('pulsar_app_config', None)
@@ -262,6 +264,23 @@ class PulsarJobRunner(AsynchronousJobRunner):
return JobDestination(runner="pulsar", params=url_to_destination_params(url))
def check_watched_item(self, job_state):
if self.use_mq:
# Might still need to check pod IPs.
job_wrapper = job_state.job_wrapper
guest_ports = job_wrapper.guest_ports
if len(guest_ports) > 0:
client = self.get_client_from_state(job_state)
job_ip = client.job_ip()
if job_ip:
ports_dict = {}
for guest_port in guest_ports:
ports_dict[str(guest_port)] = dict(host=job_ip, port=guest_port, protocol="http")
self.app.interactivetool_manager.configure_entry_points(job_wrapper.get_job(), ports_dict)
return job_state
else:
return self.check_watched_item_state(job_state)
def check_watched_item_state(self, job_state):
try:
client = self.get_client_from_state(job_state)
status = client.get_status()
@@ -366,6 +385,7 @@ class PulsarJobRunner(AsynchronousJobRunner):
remote_pulsar_app_config=remote_pulsar_app_config,
job_directory_files=job_directory_files,
container=None if not remote_container else remote_container.container_id,
guest_ports=job_wrapper.guest_ports,
)
job_id = pulsar_submit_job(client, client_job_description, remote_job_config)
log.info("Pulsar job submitted with job_id %s" % job_id)
@@ -853,6 +873,7 @@ class PulsarLegacyJobRunner(PulsarJobRunner):
class PulsarMQJobRunner(PulsarJobRunner):
"""Flavor of Pulsar job runner with sensible defaults for message queue communication."""
use_mq = True
poll = False
destination_defaults = dict(
default_file_action="remote_transfer",
@@ -868,7 +889,7 @@ KUBERNETES_DESTINATION_DEFAULTS = {
"default_file_action": "remote_transfer",
"rewrite_parameters": "true",
"jobs_directory": "/pulsar_staging",
"pulsar_container_image": "galaxy/pulsar-pod-staging:0.13.0",
"pulsar_container_image": "galaxy/pulsar-pod-staging:0.14.0",
"remote_container_handling": True,
"k8s_enabled": True,
"url": PARAMETER_SPECIFICATION_IGNORED,
@@ -878,6 +899,7 @@ KUBERNETES_DESTINATION_DEFAULTS = {
class PulsarKubernetesJobRunner(PulsarMQJobRunner):
destination_defaults = KUBERNETES_DESTINATION_DEFAULTS
poll = True # Poll so we can check API for pod IP for ITs.
def _populate_parameter_defaults(self, job_destination):
super(PulsarKubernetesJobRunner, self)._populate_parameter_defaults(job_destination)
@@ -1,5 +1,7 @@
"""Interface layer for pykube library shared between Galaxy and Pulsar."""
import logging
import os
import re
import uuid
try:
@@ -17,8 +19,13 @@ except ImportError as exc:
'this feature, please install it or correct the '
'following error:\nImportError %s' % str(exc))
log = logging.getLogger(__name__)
DEFAULT_JOB_API_VERSION = "batch/v1"
DEFAULT_NAMESPACE = "default"
INSTANCE_ID_INVALID_MESSAGE = ("Galaxy instance [%s] is either too long "
"(>20 characters) or it includes non DNS "
"acceptable characters, ignoring it.")
def ensure_pykube():
@@ -61,6 +68,44 @@ def pull_policy(params):
return None
def find_job_object_by_name(pykube_api, job_name, namespace=None):
return _find_object_by_name(Job, pykube_api, job_name, namespace=namespace)
def find_pod_object_by_name(pykube_api, pod_name, namespace=None):
return _find_object_by_name(Pod, pykube_api, pod_name, namespace=namespace)
def _find_object_by_name(clazz, pykube_api, object_name, namespace=None):
filter_kwd = dict(selector="app=%s" % object_name)
if namespace is not None:
filter_kwd["namespace"] = namespace
objs = clazz.objects(pykube_api).filter(**filter_kwd)
obj = None
if len(objs.response['items']) > 0:
obj = clazz(pykube_api, objs.response['items'][0])
return obj
def stop_job(job, cleanup="always"):
job_failed = (job.obj['status']['failed'] > 0
if 'failed' in job.obj['status'] else False)
# Scale down the job just in case even if cleanup is never
job.scale(replicas=0)
api_delete = cleanup == "always"
if not api_delete and cleanup == "onsuccess" and not job_failed:
api_delete = True
if api_delete:
delete_options = {
"apiVersion": "v1",
"kind": "DeleteOptions",
"propagationPolicy": "Background"
}
r = job.api.delete(json=delete_options, **job.api_kwargs())
job.api.raise_for_status(r)
def job_object_dict(params, job_name, spec):
k8s_job_obj = {
"apiVersion": params.get('k8s_job_api_version', DEFAULT_JOB_API_VERSION),
@@ -77,13 +122,38 @@ def job_object_dict(params, job_name, spec):
return k8s_job_obj
def galaxy_instance_id(params):
"""Parse and validate the id of the Galaxy instance from supplied dict.
This will be added to Jobs and Pods names, so it needs to be DNS friendly,
this means: `The Internet standards (Requests for Comments) for protocols mandate that component hostname labels
may contain only the ASCII letters 'a' through 'z' (in a case-insensitive manner), the digits '0' through '9',
and the minus sign ('-').`
It looks for the value set on params['k8s_galaxy_instance_id'], which might or not be set. The
idea behind this is to allow the Galaxy instance to trust (or not) existing k8s Jobs and Pods that match the
setup of a Job that is being recovered or restarted after a downtime/reboot.
"""
if "k8s_galaxy_instance_id" in params:
raw_value = params['k8s_galaxy_instance_id']
if re.match(r"(?!-)[a-z\d-]{1,20}(?<!-)$", raw_value):
return raw_value
else:
log.error(INSTANCE_ID_INVALID_MESSAGE % raw_value)
return None
__all__ = (
"DEFAULT_JOB_API_VERSION",
"ensure_pykube",
"find_job_object_by_name",
"find_pod_object_by_name",
"galaxy_instance_id",
"Job",
"job_object_dict",
"Pod",
"produce_unique_k8s_job_name",
"pull_policy",
"pykube_client_from_dict",
"stop_job",
)
+3 -3
View File
@@ -263,10 +263,10 @@ class InteractiveToolManager(object):
protocol = trans.request.host_url.split('//', 1)[0]
entry_point_encoded_id = trans.security.encode_id(entry_point.id)
entry_point_class = entry_point.__class__.__name__.lower()
entry_point_prefix = self.app.config.interactivetool_prefix
interactivetool_proxy_host = self.app.config.interactivetool_proxy_host or request_host
entry_point_prefix = self.app.config.interactivetools_prefix
interactivetools_proxy_host = self.app.config.interactivetools_proxy_host or request_host
rval = '%s//%s-%s.%s.%s.%s/' % (protocol, entry_point_encoded_id,
entry_point.token, entry_point_class, entry_point_prefix, interactivetool_proxy_host)
entry_point.token, entry_point_class, entry_point_prefix, interactivetools_proxy_host)
if entry_point.entry_url:
rval = '%s/%s' % (rval.rstrip('/'), entry_point.entry_url.lstrip('/'))
return rval
+21 -3
View File
@@ -31,6 +31,7 @@ from galaxy.tools.parameters import (
from galaxy.tools.parameters.basic import (
BaseDataToolParameter,
BooleanToolParameter,
ColorToolParameter,
ConnectedValue,
DataCollectionToolParameter,
DataToolParameter,
@@ -815,8 +816,8 @@ class InputParameterModule(WorkflowModule):
parameter_type_cond.test_param = input_parameter_type
cases = []
for param_type in ["text", "integer", "float"]:
default_source = dict(name="default", label="Default Value", type=parameter_type)
for param_type in ["text", "integer", "float", "boolean", "color"]:
default_source = dict(name="default", label="Default Value", type=param_type)
if param_type == "text":
if parameter_type == "text":
default = parameter_def.get("default") or ""
@@ -838,7 +839,22 @@ class InputParameterModule(WorkflowModule):
default = 0.0
default_source["value"] = default
input_default_value = FloatToolParameter(None, default_source)
# color parameter defaults?
elif param_type == "boolean":
if parameter_type == "boolean":
default = parameter_def.get("default") or False
else:
default = False
default_source["value"] = default
default_source["checked"] = default
input_default_value = BooleanToolParameter(None, default_source)
elif param_type == "color":
if parameter_type == 'color':
default = parameter_def.get('default') or '#000000'
else:
default = '#000000'
default_source["value"] = default
input_default_value = ColorToolParameter(None, default_source)
optional_value = optional_param(optional)
optional_cond = Conditional()
optional_cond.name = "optional"
@@ -1023,6 +1039,8 @@ class InputParameterModule(WorkflowModule):
if optional:
default_value = parameter_def.get("default", self.default_default_value)
parameter_kwds["value"] = default_value
if parameter_type == 'boolean':
parameter_kwds['checked'] = default_value
if "value" not in parameter_kwds and parameter_type in ["integer", "float"]:
parameter_kwds["value"] = str(0)
@@ -15,6 +15,8 @@ from .api import UsesApiTestCaseMixin
from .driver_util import GalaxyTestDriver
NO_APP_MESSAGE = "test_case._app called though no Galaxy has been configured."
# Following should be for Homebrew Rabbitmq and Docker on Mac "amqp://guest:guest@localhost:5672//"
AMQP_URL = os.environ.get("GALAXY_TEST_AMQP_URL", None)
def _identity(func):
@@ -28,6 +30,12 @@ def skip_if_jenkins(cls):
return cls
def skip_unless_amqp():
if AMQP_URL is not None:
return _identity
return pytest.mark.skip("AMQP_URL is not set, required for this test.")
def skip_unless_executable(executable):
if which(executable):
return _identity
@@ -53,6 +61,13 @@ def skip_unless_fixed_port():
return pytest.mark.skip("GALAXY_TEST_PORT must be set for this test.")
def skip_if_github_workflow():
if os.environ.get("GITHUB_ACTIONS", None) is None:
return _identity
return pytest.mark.skip("This test is skipped for Github actions.")
class IntegrationInstance(UsesApiTestCaseMixin):
"""Unit test case with utilities for spinning up Galaxy."""
@@ -1,6 +1,7 @@
"""Integration tests for realtime tools."""
import os
import tempfile
import pytest
import requests
@@ -10,11 +11,17 @@ from galaxy_test.base.populators import (
DatasetPopulator,
wait_on,
)
from galaxy_test.driver import integration_util
from .test_containerized_jobs import (
ContainerizedIntegrationTestCase,
disable_dependency_resolution,
DOCKERIZED_JOB_CONFIG_FILE,
)
from .test_kubernetes_staging import (
CONTAINERIZED_TEMPLATE,
job_config,
set_infrastucture_url,
)
SCRIPT_DIRECTORY = os.path.abspath(os.path.dirname(__file__))
EMBEDDED_PULSAR_JOB_CONFIG_FILE_DOCKER = os.path.join(SCRIPT_DIRECTORY, "embedded_pulsar_docker_job_conf.yml")
@@ -142,3 +149,47 @@ class InteractiveToolsRemoteProxyIntegrationTestCase(BaseInteractiveToolsIntegra
config["interactivetools_proxy_host"] = interactivetools_proxy_host
config["interactivetools_map"] = interactivetools_map
disable_dependency_resolution(config)
@integration_util.skip_unless_kubernetes()
@integration_util.skip_unless_amqp()
@integration_util.skip_if_github_workflow()
class KubeInteractiveToolsRemoteProxyIntegrationTestCase(BaseInteractiveToolsIntegrationTestCase, RunsInterativeToolTests):
"""
$ git clone https://github.com/galaxyproject/gx-it-proxy.git $HOME/gx-it-proxy
$ cd $HOME/gx-it-proxy/docker/k8s
$ # Setup proxy inside K8 cluster with kubectl - including forwarding port 8910
$ bash run.sh
$ cd ../.. # back session.
$ # Need new DB for every test.
$ rm -rf $HOME/gxitk8proxy.sqlite
$ ./lib/createdb.js --sessions $HOME/gxitk8proxy.sqlite
$ ./lib/main.js --port 9002 --ip 0.0.0.0 --verbose --sessions $HOME/gxitk8proxy.sqlite --forwardIP localhost --forwardPort 8910 &
$ cd back/to/galaxy
$ GALAXY_TEST_K8S_EXTERNAL_PROXY_HOST="localhost:9002" GALAXY_TEST_K8S_EXTERNAL_PROXY_MAP="$HOME/gxitk8proxy.sqlite" pytest -s test/integration/test_interactivetools_api.py::KubeInteractiveToolsRemoteProxyIntegrationTestCase
"""
require_uwsgi = True
@classmethod
def setUpClass(cls):
# realpath for docker deployed in a VM on Mac, also done in driver_util.
cls.jobs_directory = os.path.realpath(tempfile.mkdtemp())
super(KubeInteractiveToolsRemoteProxyIntegrationTestCase, cls).setUpClass()
@classmethod
def handle_galaxy_config_kwds(cls, config):
interactivetools_map = os.environ.get("GALAXY_TEST_K8S_EXTERNAL_PROXY_MAP")
interactivetools_proxy_host = os.environ.get("GALAXY_TEST_K8S_EXTERNAL_PROXY_HOST")
if not interactivetools_map or not interactivetools_proxy_host:
pytest.skip("External proxy not configured for test [map=%s,host=%s]" % (interactivetools_map, interactivetools_proxy_host))
config["interactivetools_proxy_host"] = interactivetools_proxy_host
config["interactivetools_map"] = interactivetools_map
config["jobs_directory"] = cls.jobs_directory
config["file_path"] = cls.jobs_directory
config["job_config_file"] = job_config(CONTAINERIZED_TEMPLATE, cls.jobs_directory)
config["default_job_shell"] = '/bin/sh'
set_infrastucture_url(config)
disable_dependency_resolution(config)
+71 -24
View File
@@ -12,17 +12,25 @@ rabbitmq installed via Homebrew, and if a fixed port is set for the test.
"""
import os
import random
import string
import tempfile
from galaxy.jobs.runners.util.pykube_util import (
Job,
pykube_client_from_dict,
)
from galaxy_test.base.populators import skip_without_tool
from galaxy_test.driver import integration_util
from .test_containerized_jobs import EXTENDED_TIMEOUT, MulledJobTestCases
from .test_job_environments import BaseJobEnvironmentIntegrationTestCase
from .test_kubernetes_runner import KubernetesDatasetPopulator
from .test_local_job_cancellation import CancelsJob
TOOL_DIR = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir, 'tools'))
AMQP_URL = os.environ.get("GALAXY_TEST_AMQP_URL", "amqp://guest:guest@localhost:5672//")
GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST = os.environ.get("GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST", "host.docker.internal")
AMQP_URL = integration_util.AMQP_URL
GALAXY_TEST_KUBERNETES_NAMESPACE = os.environ.get("GALAXY_TEST_K8S_NAMESPACE", "default")
CONTAINERIZED_TEMPLATE = """
@@ -32,7 +40,6 @@ runners:
workers: 1
pulsar_k8s:
load: galaxy.jobs.runners.pulsar:PulsarKubernetesJobRunner
galaxy_url: ${infrastructure_url}
amqp_url: ${amqp_url}
execution:
@@ -40,6 +47,8 @@ execution:
environments:
pulsar_k8s_environment:
k8s_config_path: ${k8s_config_path}
k8s_galaxy_instance_id: ${instance_id}
k8s_namespace: ${k8s_namespace}
runner: pulsar_k8s
docker_enabled: true
docker_default_container_id: busybox:ubuntu-14.04
@@ -63,7 +72,6 @@ runners:
workers: 1
pulsar_k8s:
load: galaxy.jobs.runners.pulsar:PulsarKubernetesJobRunner
galaxy_url: ${infrastructure_url}
amqp_url: ${amqp_url}
execution:
@@ -71,6 +79,8 @@ execution:
environments:
pulsar_k8s_environment:
k8s_config_path: ${k8s_config_path}
k8s_galaxy_instance_id: ${instance_id}
k8s_namespace: ${k8s_namespace}
runner: pulsar_k8s
pulsar_app_config:
message_queue_url: '${container_amqp_url}'
@@ -87,16 +97,15 @@ tools:
def job_config(template_str, jobs_directory):
job_conf_template = string.Template(template_str)
port = os.environ.get('GALAXY_TEST_PORT')
assert port
infrastructure_url = "http://%s:%s" % (GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST, port)
container_amqp_url = to_infrastructure_uri(AMQP_URL)
instance_id = ''.join(random.choice(string.ascii_lowercase) for i in range(8))
job_conf_str = job_conf_template.substitute(jobs_directory=jobs_directory,
tool_directory=TOOL_DIR,
instance_id=instance_id,
k8s_config_path=integration_util.k8s_config_path(),
k8s_namespace=GALAXY_TEST_KUBERNETES_NAMESPACE,
amqp_url=AMQP_URL,
container_amqp_url=container_amqp_url,
infrastructure_url=infrastructure_url,
)
with tempfile.NamedTemporaryFile(suffix="_kubernetes_integration_job_conf.yml", mode="w", delete=False) as job_conf:
job_conf.write(job_conf_str)
@@ -104,47 +113,79 @@ def job_config(template_str, jobs_directory):
@integration_util.skip_unless_kubernetes()
@integration_util.skip_unless_fixed_port()
class KubernetesStagingContainerIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, MulledJobTestCases):
@integration_util.skip_unless_amqp()
@integration_util.skip_if_github_workflow()
class BaseKubernetesStagingTest(BaseJobEnvironmentIntegrationTestCase, MulledJobTestCases):
# Test leverages $UWSGI_PORT in job code, need to set this up.
require_uwsgi = True
def setUp(self):
super(KubernetesStagingContainerIntegrationTestCase, self).setUp()
super(BaseKubernetesStagingTest, self).setUp()
self.dataset_populator = KubernetesDatasetPopulator(self.galaxy_interactor)
self.history_id = self.dataset_populator.new_history()
@classmethod
def setUpClass(cls):
# realpath for docker deployed in a VM on Mac, also done in driver_util.
cls.jobs_directory = os.path.realpath(tempfile.mkdtemp())
super(KubernetesStagingContainerIntegrationTestCase, cls).setUpClass()
super(BaseKubernetesStagingTest, cls).setUpClass()
class KubernetesStagingContainerIntegrationTestCase(CancelsJob, BaseKubernetesStagingTest):
@classmethod
def handle_galaxy_config_kwds(cls, config):
config["jobs_directory"] = cls.jobs_directory
config["file_path"] = cls.jobs_directory
config["job_config_file"] = job_config(CONTAINERIZED_TEMPLATE, cls.jobs_directory)
cls.job_config_file = job_config(CONTAINERIZED_TEMPLATE, cls.jobs_directory)
config["job_config_file"] = cls.job_config_file
config["default_job_shell"] = '/bin/sh'
# Disable local tool dependency resolution.
config["tool_dependency_dir"] = "none"
set_infrastucture_url(config)
@property
def instance_id(self):
import yaml
config = yaml.load(open(self.job_config_file, "r"))
return config["execution"]["environments"]["pulsar_k8s_environment"]["k8s_galaxy_instance_id"]
@skip_without_tool("cat_data_and_sleep")
def test_job_cancel(self):
with self.dataset_populator.test_history() as history_id:
job_id = self._setup_cat_data_and_sleep(history_id)
self._wait_for_job_running(job_id)
assert self._active_kubernetes_jobs == 1
delete_response = self.dataset_populator.cancel_job(job_id)
assert delete_response.json() is True
import time
time.sleep(5)
assert self._active_kubernetes_jobs == 0
@skip_without_tool("job_environment_default")
def test_job_environment(self):
job_env = self._run_and_get_environment_properties()
assert job_env.some_env == '42'
@property
def _active_kubernetes_jobs(self):
pykube_api = pykube_client_from_dict({})
# TODO: namespace.
jobs = Job.objects(pykube_api).filter()
active = 0
for job in jobs:
if self.instance_id not in job.obj["metadata"]["name"]:
continue
status = job.obj["status"]
active += status.get("active", 0)
return active
@integration_util.skip_unless_kubernetes()
@integration_util.skip_unless_fixed_port()
class KubernetesDependencyResolutionIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, MulledJobTestCases):
def setUp(self):
super(KubernetesDependencyResolutionIntegrationTestCase, self).setUp()
self.history_id = self.dataset_populator.new_history()
@classmethod
def setUpClass(cls):
# realpath for docker deployed in a VM on Mac, also done in driver_util.
cls.jobs_directory = os.path.realpath(tempfile.mkdtemp())
super(KubernetesDependencyResolutionIntegrationTestCase, cls).setUpClass()
class KubernetesDependencyResolutionIntegrationTestCase(BaseKubernetesStagingTest):
@classmethod
def handle_galaxy_config_kwds(cls, config):
@@ -156,6 +197,7 @@ class KubernetesDependencyResolutionIntegrationTestCase(BaseJobEnvironmentIntegr
# Disable tool dependency resolution.
config["tool_dependency_dir"] = "none"
config["enable_beta_mulled_containers"] = "true"
set_infrastucture_url(config)
def test_mulled_simple(self):
self.dataset_populator.run_tool("mulled_example_simple", {}, self.history_id)
@@ -164,6 +206,11 @@ class KubernetesDependencyResolutionIntegrationTestCase(BaseJobEnvironmentIntegr
assert "0.7.15-r1140" in output
def set_infrastucture_url(config):
infrastructure_url = "http://%s:$UWSGI_PORT" % GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST
config["galaxy_infrastructure_url"] = infrastructure_url
def to_infrastructure_uri(uri):
# remap MQ or file server URI hostnames for in-container versions, this is sloppy
# should actually parse the URI and rebuild with correct host
+22 -16
View File
@@ -10,15 +10,9 @@ from galaxy_test.base.populators import (
from galaxy_test.driver import integration_util
class LocalJobCancellationTestCase(integration_util.IntegrationTestCase):
class CancelsJob(object):
framework_tool_and_types = True
def setUp(self):
super(LocalJobCancellationTestCase, self).setUp()
self.dataset_populator = DatasetPopulator(self.galaxy_interactor)
def setup_cat_data_and_sleep(self, history_id):
def _setup_cat_data_and_sleep(self, history_id):
hda1 = self.dataset_populator.new_dataset(history_id, content="1 2 3")
running_inputs = {
"input1": {"src": "hda", "id": hda1["id"]},
@@ -31,14 +25,26 @@ class LocalJobCancellationTestCase(integration_util.IntegrationTestCase):
assert_ok=False,
).json()
job_dict = running_response["jobs"][0]
return job_dict
return job_dict["id"]
def _wait_for_job_running(self, job_id):
self.galaxy_interactor.wait_for(lambda: self._get("jobs/%s" % job_id).json()['state'] != 'running',
what="Wait for job to start running",
maxseconds=60)
class LocalJobCancellationTestCase(CancelsJob, integration_util.IntegrationTestCase):
framework_tool_and_types = True
def setUp(self):
super(LocalJobCancellationTestCase, self).setUp()
self.dataset_populator = DatasetPopulator(self.galaxy_interactor)
def test_cancel_job_with_admin_message(self):
with self.dataset_populator.test_history() as history_id:
job_dict = self.setup_cat_data_and_sleep(history_id)
self.galaxy_interactor.wait_for(lambda: self._get("jobs/%s" % job_dict['id']).json()['state'] != 'running',
what="Wait for job to start running",
maxseconds=60)
job_id = self._setup_cat_data_and_sleep(history_id)
self._wait_for_job_running(job_id)
app = self._app
sa_session = app.model.context.current
Job = app.model.Job
@@ -48,7 +54,7 @@ class LocalJobCancellationTestCase(integration_util.IntegrationTestCase):
job.set_state(app.model.Job.states.DELETED_NEW)
sa_session.add(job)
sa_session.flush()
self.galaxy_interactor.wait_for(lambda: self._get("jobs/%s" % job_dict['id']).json()['state'] != 'error',
self.galaxy_interactor.wait_for(lambda: self._get("jobs/%s" % job_id).json()['state'] != 'error',
what="Wait for job to end in error",
maxseconds=60)
@@ -56,7 +62,7 @@ class LocalJobCancellationTestCase(integration_util.IntegrationTestCase):
"""
"""
with self.dataset_populator.test_history() as history_id:
job_dict = self.setup_cat_data_and_sleep(history_id)
job_id = self._setup_cat_data_and_sleep(history_id)
app = self._app
sa_session = app.model.context.current
@@ -80,7 +86,7 @@ class LocalJobCancellationTestCase(integration_util.IntegrationTestCase):
pid_exists = psutil.pid_exists(external_id)
assert pid_exists
delete_response = self.dataset_populator.cancel_job(job_dict["id"])
delete_response = self.dataset_populator.cancel_job(job_id)
assert delete_response.json() is True
state = None
@@ -792,6 +792,11 @@ execution:
docker_default_container_id: 'conda/miniconda2'
# Specify a non-default Pulsar staging container.
#pulsar_container_image: 'galaxy/pulsar-pod-staging:0.12.0'
# Generate job names with a string unique to this Galaxy (see
# Kubernetes runner description).
#k8s_galaxy_instance_id: mycoolgalaxy
# Path to Kubernetes configuration fil (see Kubernetes runner description.)
#k8s_config_path: /path/to/kubeconfig
# Example CLI runners.
ssh_torque:
+1
View File
@@ -145,6 +145,7 @@ class MockJobWrapper(object):
self.cleanup_job = "never"
self.tmp_dir_creation_statement = ""
self.use_metadata_binary = False
self.guest_ports = []
# Cruft for setting metadata externally, axe at some point.
self.external_output_metadata = bunch.Bunch(