From 731ccfbe725d219059ca98eac2d5e6b6a95826c6 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 16 Jan 2020 16:19:55 -0500 Subject: [PATCH 01/20] Reduce duplication in test_kubernetes_staging --- test/integration/test_kubernetes_staging.py | 17 +++++------------ 1 file changed, 5 insertions(+), 12 deletions(-) diff --git a/test/integration/test_kubernetes_staging.py b/test/integration/test_kubernetes_staging.py index 864dd5fb8b9..30d658684d5 100644 --- a/test/integration/test_kubernetes_staging.py +++ b/test/integration/test_kubernetes_staging.py @@ -105,7 +105,7 @@ def job_config(template_str, jobs_directory): @integration_util.skip_unless_kubernetes() @integration_util.skip_unless_fixed_port() -class KubernetesStagingContainerIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, MulledJobTestCases): +class BaseKubernetesStagingTest(BaseJobEnvironmentIntegrationTestCase, MulledJobTestCases): def setUp(self): super(KubernetesStagingContainerIntegrationTestCase, self).setUp() @@ -117,6 +117,9 @@ class KubernetesStagingContainerIntegrationTestCase(BaseJobEnvironmentIntegratio cls.jobs_directory = os.path.realpath(tempfile.mkdtemp()) super(KubernetesStagingContainerIntegrationTestCase, cls).setUpClass() + +class KubernetesStagingContainerIntegrationTestCase(BaseKubernetesStagingTest): + @classmethod def handle_galaxy_config_kwds(cls, config): config["jobs_directory"] = cls.jobs_directory @@ -134,17 +137,7 @@ class KubernetesStagingContainerIntegrationTestCase(BaseJobEnvironmentIntegratio @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): From 069f794fc7787236cf21ecb49a9d6f0b2b243744 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 1 Apr 2020 16:04:12 -0400 Subject: [PATCH 02/20] Rework conditions for running Kubernetes staging tests. --- lib/galaxy_test/driver/integration_util.py | 8 ++++++ test/integration/test_kubernetes_staging.py | 28 +++++++++++++-------- 2 files changed, 26 insertions(+), 10 deletions(-) diff --git a/lib/galaxy_test/driver/integration_util.py b/lib/galaxy_test/driver/integration_util.py index 59e3ef17b86..e2da2defbd6 100644 --- a/lib/galaxy_test/driver/integration_util.py +++ b/lib/galaxy_test/driver/integration_util.py @@ -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 diff --git a/test/integration/test_kubernetes_staging.py b/test/integration/test_kubernetes_staging.py index 30d658684d5..74a7b16abbf 100644 --- a/test/integration/test_kubernetes_staging.py +++ b/test/integration/test_kubernetes_staging.py @@ -15,14 +15,19 @@ import os 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_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 CONTAINERIZED_TEMPLATE = """ @@ -32,7 +37,6 @@ runners: workers: 1 pulsar_k8s: load: galaxy.jobs.runners.pulsar:PulsarKubernetesJobRunner - galaxy_url: ${infrastructure_url} amqp_url: ${amqp_url} execution: @@ -63,7 +67,6 @@ runners: workers: 1 pulsar_k8s: load: galaxy.jobs.runners.pulsar:PulsarKubernetesJobRunner - galaxy_url: ${infrastructure_url} amqp_url: ${amqp_url} execution: @@ -87,16 +90,12 @@ 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) job_conf_str = job_conf_template.substitute(jobs_directory=jobs_directory, tool_directory=TOOL_DIR, k8s_config_path=integration_util.k8s_config_path(), 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,18 +103,20 @@ def job_config(template_str, jobs_directory): @integration_util.skip_unless_kubernetes() -@integration_util.skip_unless_fixed_port() +@integration_util.skip_unless_amqp() 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.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(BaseKubernetesStagingTest): @@ -128,6 +129,7 @@ class KubernetesStagingContainerIntegrationTestCase(BaseKubernetesStagingTest): config["default_job_shell"] = '/bin/sh' # Disable local tool dependency resolution. config["tool_dependency_dir"] = "none" + set_infrastucture_url(config) @skip_without_tool("job_environment_default") def test_job_environment(self): @@ -149,6 +151,7 @@ class KubernetesDependencyResolutionIntegrationTestCase(BaseKubernetesStagingTes # 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) @@ -157,6 +160,11 @@ class KubernetesDependencyResolutionIntegrationTestCase(BaseKubernetesStagingTes 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 From f814780f5e7c1cd5475e2abab9c54663e950e7da Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 30 Mar 2020 09:54:25 -0400 Subject: [PATCH 03/20] Refactor galaxy instance ID parsing for k8s for reuse. --- lib/galaxy/jobs/runners/kubernetes.py | 22 +++-------------- lib/galaxy/jobs/runners/util/pykube_util.py | 26 +++++++++++++++++++++ 2 files changed, 29 insertions(+), 19 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index ad68c60db77..77e43c471e3 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -19,6 +19,7 @@ from galaxy.jobs.runners import ( from galaxy.jobs.runners.util.pykube_util import ( DEFAULT_JOB_API_VERSION, ensure_pykube, + galaxy_instance_id, Job, job_object_dict, Pod, @@ -201,25 +202,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}(?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. diff --git a/lib/galaxy/jobs/runners/util/pykube_util.py b/lib/galaxy/jobs/runners/util/pykube_util.py index 0203b9f5046..b0e7fba5795 100644 --- a/lib/galaxy/jobs/runners/util/pykube_util.py +++ b/lib/galaxy/jobs/runners/util/pykube_util.py @@ -1,5 +1,7 @@ """Interface layer for pykube library shared between Galaxy and Pulsar.""" +import logging import os +import re import uuid try: @@ -17,6 +19,8 @@ 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" @@ -77,9 +81,31 @@ 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: + if re.match(r"(?!-)[a-z\d-]{1,20}(?20 characters) or it includes non DNS acceptable characters, ignoring it.') + return None + + __all__ = ( "DEFAULT_JOB_API_VERSION", "ensure_pykube", + "galaxy_instance_id", "Job", "job_object_dict", "Pod", From 5d8fb30448a5e258ae95c8816584b77b896c7e31 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 30 Mar 2020 09:55:51 -0400 Subject: [PATCH 04/20] Refactor local job cancel test case for reuse. --- .../test_local_job_cancellation.py | 31 ++++++++++--------- 1 file changed, 17 insertions(+), 14 deletions(-) diff --git a/test/integration/test_local_job_cancellation.py b/test/integration/test_local_job_cancellation.py index a0c58090794..c5ac96c9d4c 100644 --- a/test/integration/test_local_job_cancellation.py +++ b/test/integration/test_local_job_cancellation.py @@ -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,12 +25,21 @@ class LocalJobCancellationTestCase(integration_util.IntegrationTestCase): assert_ok=False, ).json() job_dict = running_response["jobs"][0] - return job_dict + return job_dict["id"] + + +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', + job_id = self._setup_cat_data_and_sleep(history_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) app = self._app @@ -48,7 +51,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 +59,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 +83,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 From 1c512af82e08d4b947213a61f09f4dbac905c87a Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 30 Mar 2020 11:11:01 -0400 Subject: [PATCH 05/20] Set random instance ID in K8S staging tests. Before we were using UUID but Pulsar switching to using stable IDs - so we can use this method to avoid duplicate job names instead. --- test/integration/test_kubernetes_staging.py | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/test/integration/test_kubernetes_staging.py b/test/integration/test_kubernetes_staging.py index 74a7b16abbf..ccc4bf30c8e 100644 --- a/test/integration/test_kubernetes_staging.py +++ b/test/integration/test_kubernetes_staging.py @@ -12,6 +12,7 @@ rabbitmq installed via Homebrew, and if a fixed port is set for the test. """ import os +import random import string import tempfile @@ -44,6 +45,7 @@ execution: environments: pulsar_k8s_environment: k8s_config_path: ${k8s_config_path} + k8s_galaxy_instance_id: ${instance_id} runner: pulsar_k8s docker_enabled: true docker_default_container_id: busybox:ubuntu-14.04 @@ -74,6 +76,7 @@ execution: environments: pulsar_k8s_environment: k8s_config_path: ${k8s_config_path} + k8s_galaxy_instance_id: ${instance_id} runner: pulsar_k8s pulsar_app_config: message_queue_url: '${container_amqp_url}' @@ -91,8 +94,10 @@ tools: def job_config(template_str, jobs_directory): job_conf_template = string.Template(template_str) 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(), amqp_url=AMQP_URL, container_amqp_url=container_amqp_url, From 63515fe15eff4f7abd01a195f4939847c6c5ed35 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 1 Apr 2020 09:11:14 -0400 Subject: [PATCH 06/20] Test case for stopping Pulsar Kubernetes job. --- lib/galaxy/jobs/runners/kubernetes.py | 25 +++-------- lib/galaxy/jobs/runners/util/pykube_util.py | 43 +++++++++++++++++-- test/integration/test_kubernetes_staging.py | 41 +++++++++++++++++- .../test_local_job_cancellation.py | 9 ++-- 4 files changed, 91 insertions(+), 27 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 77e43c471e3..aaac45ebf11 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -19,6 +19,7 @@ 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, @@ -26,6 +27,7 @@ from galaxy.jobs.runners.util.pykube_util import ( produce_unique_k8s_job_name, pull_policy, pykube_client_from_dict, + stop_job, ) from galaxy.util.bytesize import ByteSize @@ -493,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 [] @@ -534,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" diff --git a/lib/galaxy/jobs/runners/util/pykube_util.py b/lib/galaxy/jobs/runners/util/pykube_util.py index b0e7fba5795..a4cbe9273ad 100644 --- a/lib/galaxy/jobs/runners/util/pykube_util.py +++ b/lib/galaxy/jobs/runners/util/pykube_util.py @@ -23,6 +23,9 @@ 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(): @@ -65,6 +68,36 @@ def pull_policy(params): return None +def find_job_object_by_name(pykube_api, job_name, namespace=None): + filter_kwd = dict(selector="app=%s" % job_name) + if namespace is not None: + filter_kwd["namespace"] = namespace + + jobs = Job.objects(pykube_api).filter(**filter_kwd) + job = None + if len(jobs.response['items']) > 0: + job = Job(pykube_api, jobs.response['items'][0]) + return job + + +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), @@ -94,17 +127,18 @@ def galaxy_instance_id(params): setup of a Job that is being recovered or restarted after a downtime/reboot. """ if "k8s_galaxy_instance_id" in params: - if re.match(r"(?!-)[a-z\d-]{1,20}(?20 characters) or it includes non DNS acceptable characters, ignoring it.') + log.error(INSTANCE_ID_INVALID_MESSAGE % raw_value) return None __all__ = ( "DEFAULT_JOB_API_VERSION", "ensure_pykube", + "find_job_object_by_name", "galaxy_instance_id", "Job", "job_object_dict", @@ -112,4 +146,5 @@ __all__ = ( "produce_unique_k8s_job_name", "pull_policy", "pykube_client_from_dict", + "stop_job", ) diff --git a/test/integration/test_kubernetes_staging.py b/test/integration/test_kubernetes_staging.py index ccc4bf30c8e..02488e784e8 100644 --- a/test/integration/test_kubernetes_staging.py +++ b/test/integration/test_kubernetes_staging.py @@ -22,6 +22,7 @@ from galaxy.jobs.runners.util.pykube_util import ( ) from galaxy_test.base.populators import skip_without_tool from galaxy_test.driver import integration_util +from .test_local_job_cancellation import CancelsJob from .test_containerized_jobs import EXTENDED_TIMEOUT, MulledJobTestCases from .test_job_environments import BaseJobEnvironmentIntegrationTestCase from .test_local_job_cancellation import CancelsJob @@ -124,23 +125,59 @@ class BaseKubernetesStagingTest(BaseJobEnvironmentIntegrationTestCase, MulledJob super(BaseKubernetesStagingTest, cls).setUpClass() -class KubernetesStagingContainerIntegrationTestCase(BaseKubernetesStagingTest): +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() diff --git a/test/integration/test_local_job_cancellation.py b/test/integration/test_local_job_cancellation.py index c5ac96c9d4c..aa0ecc94019 100644 --- a/test/integration/test_local_job_cancellation.py +++ b/test/integration/test_local_job_cancellation.py @@ -27,6 +27,11 @@ class CancelsJob(object): job_dict = running_response["jobs"][0] 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): @@ -39,9 +44,7 @@ class LocalJobCancellationTestCase(CancelsJob, integration_util.IntegrationTestC def test_cancel_job_with_admin_message(self): with self.dataset_populator.test_history() as history_id: job_id = self._setup_cat_data_and_sleep(history_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) + self._wait_for_job_running(job_id) app = self._app sa_session = app.model.context.current Job = app.model.Job From 8fe2b918c8de6a49aa0853a748ede676a2396b49 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 1 Apr 2020 13:55:33 -0400 Subject: [PATCH 07/20] A couple more docs for Kuberentes Pulsar execution. --- test/unit/jobs/job_conf.sample_advanced.yml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/test/unit/jobs/job_conf.sample_advanced.yml b/test/unit/jobs/job_conf.sample_advanced.yml index bdffa03ed33..79635c9ec34 100644 --- a/test/unit/jobs/job_conf.sample_advanced.yml +++ b/test/unit/jobs/job_conf.sample_advanced.yml @@ -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: From 52ee3b67d6b0f982d5a30a9b69d56688df792a01 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 2 Apr 2020 10:37:07 -0400 Subject: [PATCH 08/20] Rev pulsar for fixes. --- .../dependencies/pipfiles/default/pinned-requirements.txt | 2 +- lib/galaxy/jobs/runners/pulsar.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt b/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt index 25237754d4b..8fa66c01da6 100644 --- a/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt @@ -123,7 +123,7 @@ pbr==5.4.4 prettytable==0.7.2 prov==1.5.1 psutil==5.6.7 -pulsar-galaxy-lib==0.14.0.dev1 +pulsar-galaxy-lib==0.14.0.dev2 pyasn1-modules==0.2.7 pyasn1==0.4.8 pycparser==2.19 diff --git a/lib/galaxy/jobs/runners/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py index ebd72250e58..b003b8724b4 100644 --- a/lib/galaxy/jobs/runners/pulsar.py +++ b/lib/galaxy/jobs/runners/pulsar.py @@ -868,7 +868,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, From 910807d1c9c023f2939d5d750f5ceb901b9093d7 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 2 Apr 2020 16:08:33 +0200 Subject: [PATCH 09/20] Log node state if history doesn't finish --- test/integration/test_kubernetes_runner.py | 16 +++++++++++++++- 1 file changed, 15 insertions(+), 1 deletion(-) diff --git a/test/integration/test_kubernetes_runner.py b/test/integration/test_kubernetes_runner.py index 000b1029737..09b325b92b0 100644 --- a/test/integration/test_kubernetes_runner.py +++ b/test/integration/test_kubernetes_runner.py @@ -12,7 +12,10 @@ import time import pytest from galaxy.util import unicodify -from galaxy_test.base.populators import skip_without_tool +from galaxy_test.base.populators import ( + DatasetPopulator, + skip_without_tool, +) from galaxy_test.driver import integration_util from .test_containerized_jobs import MulledJobTestCases from .test_job_environments import BaseJobEnvironmentIntegrationTestCase @@ -151,11 +154,22 @@ def job_config(jobs_directory): return Config(job_conf.name) +class KubernetesDatasetPopulator(DatasetPopulator): + + def wait_for_history(self, *args, **kwargs): + try: + super(KubernetesDatasetPopulator, self).wait_for_history(*args, **kwargs) + except AssertionError: + print("Kubernetes status:\n %s" % unicodify(subprocess.check_output(['kubectl', 'describe', 'nodes']))) + raise + + @integration_util.skip_unless_kubernetes() class BaseKubernetesIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, MulledJobTestCases): def setUp(self): super(BaseKubernetesIntegrationTestCase, self).setUp() + self.dataset_populator = KubernetesDatasetPopulator(self.galaxy_interactor) self.history_id = self.dataset_populator.new_history() @classmethod From 25c30405d913317285f0ffc673189d1d65df06aa Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 2 Apr 2020 19:53:59 +0200 Subject: [PATCH 10/20] Delete 15G dotnet folder, use minikube 1.9.0 --- .github/workflows/integration.yaml | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/.github/workflows/integration.yaml b/.github/workflows/integration.yaml index 25edfff1b8c..fb3ff8d1648 100644 --- a/.github/workflows/integration.yaml +++ b/.github/workflows/integration.yaml @@ -27,10 +27,15 @@ jobs: steps: - name: Prune unused docker image, volumes and containers run: docker system prune -a -f + - name: Clean dotnet folder for space + if: matrix.subset == 'kubernetes' + run: rm -Rf /usr/share/dotnet - name: Setup Minikube if: matrix.subset == 'kubernetes' id: minikube - uses: CodingNagger/minikube-setup-action@v1.0.2 + uses: CodingNagger/minikube-setup-action@v1.0.3 + with: + minikube-version: "1.9.0-0_amd64" - name: Launch Minikube if: matrix.subset == 'kubernetes' run: eval ${{ steps.minikube.outputs.launcher }} From a0d108b844ef20254c36d4fd92a75838fbb509a0 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 26 Mar 2020 09:35:23 -0400 Subject: [PATCH 11/20] Outline a test of kubernetes staging + interactive tool. --- test/integration/test_interactivetools_api.py | 50 +++++++++++++++++++ 1 file changed, 50 insertions(+) diff --git a/test/integration/test_interactivetools_api.py b/test/integration/test_interactivetools_api.py index 762ecd458a0..6411eca9cca 100644 --- a/test/integration/test_interactivetools_api.py +++ b/test/integration/test_interactivetools_api.py @@ -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,46 @@ 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() +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) From 35ec208ea61fa858e089d4bb59f95955a970afab Mon Sep 17 00:00:00 2001 From: John Chilton Date: Sun, 12 Apr 2020 12:35:09 -0400 Subject: [PATCH 12/20] Fix for interactivetool_proxy_host added in 9504. --- lib/galaxy/config/__init__.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/config/__init__.py b/lib/galaxy/config/__init__.py index 6f529186bff..3e2b4ae5dd6 100644 --- a/lib/galaxy/config/__init__.py +++ b/lib/galaxy/config/__init__.py @@ -636,7 +636,7 @@ class GalaxyAppConfiguration(BaseAppConfiguration): # InteractiveTools propagator mapping file self.interactivetools_map = self.resolve_path(kwargs.get("interactivetools_map", os.path.join(self.data_dir, "interactivetools_map.sqlite"))) self.interactivetool_prefix = kwargs.get("interactivetools_prefix", "interactivetool") - self.interactivetool_proxy_host = kwargs.get("interactivetool_proxy_host", None) + self.interactivetool_proxy_host = kwargs.get("interactivetools_proxy_host", None) self.containers_conf = parse_containers_config(self.containers_config_file) From 22c4a2989262b64e958ab427377e6caa1e2e2cbb Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 14 Apr 2020 09:40:26 -0400 Subject: [PATCH 13/20] Have Pulsar in Kubernetes mode poll for ITs if needed. --- lib/galaxy/jobs/runners/pulsar.py | 28 ++++++++++++++++++--- lib/galaxy/jobs/runners/util/pykube_util.py | 21 +++++++++++----- 2 files changed, 40 insertions(+), 9 deletions(-) diff --git a/lib/galaxy/jobs/runners/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py index b003b8724b4..68befa97b86 100644 --- a/lib/galaxy/jobs/runners/pulsar.py +++ b/lib/galaxy/jobs/runners/pulsar.py @@ -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", @@ -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) diff --git a/lib/galaxy/jobs/runners/util/pykube_util.py b/lib/galaxy/jobs/runners/util/pykube_util.py index a4cbe9273ad..f2e06f78c59 100644 --- a/lib/galaxy/jobs/runners/util/pykube_util.py +++ b/lib/galaxy/jobs/runners/util/pykube_util.py @@ -69,15 +69,23 @@ def pull_policy(params): def find_job_object_by_name(pykube_api, job_name, namespace=None): - filter_kwd = dict(selector="app=%s" % job_name) + 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 - jobs = Job.objects(pykube_api).filter(**filter_kwd) - job = None - if len(jobs.response['items']) > 0: - job = Job(pykube_api, jobs.response['items'][0]) - return job + 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"): @@ -139,6 +147,7 @@ __all__ = ( "DEFAULT_JOB_API_VERSION", "ensure_pykube", "find_job_object_by_name", + "find_pod_object_by_name", "galaxy_instance_id", "Job", "job_object_dict", From 0fc027710e3b4bd06cb424e3e48d7ae90faa3114 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 13 Apr 2020 11:35:56 -0400 Subject: [PATCH 14/20] Use namespaces in Kubernetes tests. --- test/integration/test_kubernetes_staging.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/test/integration/test_kubernetes_staging.py b/test/integration/test_kubernetes_staging.py index 02488e784e8..5e60c99a1f7 100644 --- a/test/integration/test_kubernetes_staging.py +++ b/test/integration/test_kubernetes_staging.py @@ -30,6 +30,7 @@ from .test_local_job_cancellation import CancelsJob TOOL_DIR = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir, 'tools')) 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 = """ @@ -47,6 +48,7 @@ execution: 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 @@ -78,6 +80,7 @@ execution: 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}' @@ -100,6 +103,7 @@ def job_config(template_str, 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, ) From 461a7ce272a0cec77fe795a8ce3fd5d9fad4a144 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 14 Apr 2020 11:40:53 -0400 Subject: [PATCH 15/20] Rev Pulsar for ITs client support. --- .../dependencies/pipfiles/default/pinned-requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt b/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt index 8fa66c01da6..4a223d125e4 100644 --- a/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt @@ -123,7 +123,7 @@ pbr==5.4.4 prettytable==0.7.2 prov==1.5.1 psutil==5.6.7 -pulsar-galaxy-lib==0.14.0.dev2 +pulsar-galaxy-lib==0.14.0.dev3 pyasn1-modules==0.2.7 pyasn1==0.4.8 pycparser==2.19 From e6163ba9465117347cfcb4948fc085527103d642 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 14 Apr 2020 14:34:03 -0400 Subject: [PATCH 16/20] Better job_wrapper abstraction for guest_ports. --- lib/galaxy/jobs/__init__.py | 8 ++++++++ lib/galaxy/jobs/runners/__init__.py | 2 +- test/unit/jobs/test_runner_local.py | 1 + 3 files changed, 10 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 71dd60e7852..7c2567eb0ff 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -1122,6 +1122,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: diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index 2ae1550ffea..e14b07cfffb 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -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, diff --git a/test/unit/jobs/test_runner_local.py b/test/unit/jobs/test_runner_local.py index 173e710cbd2..3e19045fb5d 100644 --- a/test/unit/jobs/test_runner_local.py +++ b/test/unit/jobs/test_runner_local.py @@ -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( From 97db452eb38fa9e51bb2fd28c99ee44ee1e89fe1 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 15 Apr 2020 13:30:32 -0400 Subject: [PATCH 17/20] Fixup config option names. --- lib/galaxy/config/__init__.py | 4 ++-- lib/galaxy/managers/interactivetool.py | 6 +++--- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/lib/galaxy/config/__init__.py b/lib/galaxy/config/__init__.py index 3e2b4ae5dd6..d37f24bd3d2 100644 --- a/lib/galaxy/config/__init__.py +++ b/lib/galaxy/config/__init__.py @@ -635,8 +635,8 @@ class GalaxyAppConfiguration(BaseAppConfiguration): # InteractiveTools propagator mapping file self.interactivetools_map = self.resolve_path(kwargs.get("interactivetools_map", os.path.join(self.data_dir, "interactivetools_map.sqlite"))) - self.interactivetool_prefix = kwargs.get("interactivetools_prefix", "interactivetool") - self.interactivetool_proxy_host = kwargs.get("interactivetools_proxy_host", None) + self.interactivetools_prefix = kwargs.get("interactivetools_prefix", "interactivetool") + self.interactivetools_proxy_host = kwargs.get("interactivetools_proxy_host", None) self.containers_conf = parse_containers_config(self.containers_config_file) diff --git a/lib/galaxy/managers/interactivetool.py b/lib/galaxy/managers/interactivetool.py index 3c7c2d1aa5a..08e44a9ce5b 100644 --- a/lib/galaxy/managers/interactivetool.py +++ b/lib/galaxy/managers/interactivetool.py @@ -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 From a29d82c4d4d0b1a705adfb85a11e04fd6619a81a Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 21 Apr 2020 09:29:17 -0400 Subject: [PATCH 18/20] More debug for k8s CI. --- test/integration/test_kubernetes_staging.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/test/integration/test_kubernetes_staging.py b/test/integration/test_kubernetes_staging.py index 5e60c99a1f7..25994b8a0d0 100644 --- a/test/integration/test_kubernetes_staging.py +++ b/test/integration/test_kubernetes_staging.py @@ -25,6 +25,7 @@ from galaxy_test.driver import integration_util from .test_local_job_cancellation import CancelsJob 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')) @@ -120,6 +121,7 @@ class BaseKubernetesStagingTest(BaseJobEnvironmentIntegrationTestCase, MulledJob def setUp(self): super(BaseKubernetesStagingTest, self).setUp() + self.dataset_populator = KubernetesDatasetPopulator(self.galaxy_interactor) self.history_id = self.dataset_populator.new_history() @classmethod From 52bd3ad9660250f3ffc8f5f3b3bd22b44c5b9614 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 21 Apr 2020 13:09:38 -0400 Subject: [PATCH 19/20] Return to skipping Kubernetes staging tests on Github actions... ... still working out some kind of issue with these. --- lib/galaxy_test/driver/integration_util.py | 7 +++++++ test/integration/test_interactivetools_api.py | 1 + test/integration/test_kubernetes_staging.py | 3 +-- 3 files changed, 9 insertions(+), 2 deletions(-) diff --git a/lib/galaxy_test/driver/integration_util.py b/lib/galaxy_test/driver/integration_util.py index e2da2defbd6..caf05cf8a74 100644 --- a/lib/galaxy_test/driver/integration_util.py +++ b/lib/galaxy_test/driver/integration_util.py @@ -61,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.""" diff --git a/test/integration/test_interactivetools_api.py b/test/integration/test_interactivetools_api.py index 6411eca9cca..7216d6ee78c 100644 --- a/test/integration/test_interactivetools_api.py +++ b/test/integration/test_interactivetools_api.py @@ -153,6 +153,7 @@ class InteractiveToolsRemoteProxyIntegrationTestCase(BaseInteractiveToolsIntegra @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 diff --git a/test/integration/test_kubernetes_staging.py b/test/integration/test_kubernetes_staging.py index 25994b8a0d0..187e6b27f51 100644 --- a/test/integration/test_kubernetes_staging.py +++ b/test/integration/test_kubernetes_staging.py @@ -115,6 +115,7 @@ def job_config(template_str, jobs_directory): @integration_util.skip_unless_kubernetes() @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 @@ -185,8 +186,6 @@ class KubernetesStagingContainerIntegrationTestCase(CancelsJob, BaseKubernetesSt return active -@integration_util.skip_unless_kubernetes() -@integration_util.skip_unless_fixed_port() class KubernetesDependencyResolutionIntegrationTestCase(BaseKubernetesStagingTest): @classmethod From e3ddd908be01c1d2ee2526ada08c00bb53946b7b Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 23 Apr 2020 14:44:44 -0400 Subject: [PATCH 20/20] How did that revert? --- test/integration/test_kubernetes_staging.py | 1 - 1 file changed, 1 deletion(-) diff --git a/test/integration/test_kubernetes_staging.py b/test/integration/test_kubernetes_staging.py index 187e6b27f51..14523a30e07 100644 --- a/test/integration/test_kubernetes_staging.py +++ b/test/integration/test_kubernetes_staging.py @@ -22,7 +22,6 @@ from galaxy.jobs.runners.util.pykube_util import ( ) from galaxy_test.base.populators import skip_without_tool from galaxy_test.driver import integration_util -from .test_local_job_cancellation import CancelsJob from .test_containerized_jobs import EXTENDED_TIMEOUT, MulledJobTestCases from .test_job_environments import BaseJobEnvironmentIntegrationTestCase from .test_kubernetes_runner import KubernetesDatasetPopulator