From 8cf3f4fc32f799273af19bc20fccb06aa3871896 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 23 Apr 2020 16:02:32 +0200 Subject: [PATCH 01/24] Fix workfow parameter connections and default values in workflow editor and run form. The main problem was that `default_source = dict(name="default", label="Default Value", type=parameter_type)`, where `parameter_type` is the top-level parameter that is always `text`. This has changed with the Conditional introduced in 20.01. In addition BooleanToolParameters use `check` instead of `value`. I also added ColorToolParameter. One thing for dev is that I think we should pull the default values out of the optional conditional so we can set defaults for required parameters. Fixes https://github.com/galaxyproject/galaxy/issues/9646 --- lib/galaxy/workflow/modules.py | 24 +++++++++++++++++++++--- 1 file changed, 21 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/workflow/modules.py b/lib/galaxy/workflow/modules.py index 31a2348c926..c1609f029be 100644 --- a/lib/galaxy/workflow/modules.py +++ b/lib/galaxy/workflow/modules.py @@ -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) From 731ccfbe725d219059ca98eac2d5e6b6a95826c6 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 16 Jan 2020 16:19:55 -0500 Subject: [PATCH 02/24] 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 03/24] 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 04/24] 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 05/24] 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 06/24] 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 07/24] 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 08/24] 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 09/24] 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 10/24] 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 11/24] 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 12/24] 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 13/24] 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 14/24] 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 15/24] 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 16/24] 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 17/24] 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 18/24] 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 19/24] 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 20/24] 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 21/24] 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 From 33d017337b52315eee9b81cecf55b8a10838c057 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 2 Apr 2020 19:53:59 +0200 Subject: [PATCH 22/24] 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 bf4fb7735970568d56bd427e208b084a92ad87cb Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 2 Apr 2020 19:53:59 +0200 Subject: [PATCH 23/24] 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 58c82ad7f37eb577ef66a0546415c76873064bc2 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 24 Apr 2020 13:19:56 +0200 Subject: [PATCH 24/24] Fix vcf_bgzip mimetype, allow cors, tbi from metadata element The mimetype would be text/plain otherwise, which is not compatible with byte-range requests. `allow_cors` doesn't quite seem to work in my tests, but we'll need this for jbrowse. Also skips an additional conversion by using the tabix index, which again is helpful when loading datasets from a shared history. --- display_applications/igv/vcf.xml | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/display_applications/igv/vcf.xml b/display_applications/igv/vcf.xml index f2c1e07966a..3c66ffd677a 100644 --- a/display_applications/igv/vcf.xml +++ b/display_applications/igv/vcf.xml @@ -12,8 +12,8 @@ ${$site_id.startswith( 'local_' ) or $dataset.dbkey in $site_dbkeys} ${redirect_url} - - + + #if ($dataset.dbkey in $site_dbkeys) $site_organisms[ $site_dbkeys.index( $bgzip_file.dbkey ) ] @@ -94,8 +94,8 @@ ${ $dataset.dbkey == $value } http://www.broadinstitute.org/igv/projects/current/igv.php?sessionURL=${bgzip_file.qp}&genome=${qp($bgzip_file.dbkey)}&merge=true&name=${qp( ( $bgzip_file.name or $DATASET_HASH ).replace( ',', ';' ) )} - - + +