From ac1b25ed4c069abf1c8d810ba26d5be43717766d Mon Sep 17 00:00:00 2001 From: John Chilton Date: Sun, 28 Apr 2019 11:49:55 -0400 Subject: [PATCH] Kubernetes 2-container pod staging. The tests included here are really very manual and brittle tests, but I'm not sure how to do much better. See documentation in the file docstring for more information. --- lib/galaxy/jobs/runners/kubernetes.py | 64 ++++------- lib/galaxy/jobs/runners/pulsar.py | 70 ++++++++++-- lib/galaxy/jobs/runners/util/pykube_util.py | 89 +++++++++++++++ test/base/integration_util.py | 18 ++- test/integration/test_kubernetes_runner.py | 4 +- test/integration/test_kubernetes_staging.py | 117 ++++++++++++++++++++ test/unit/jobs/job_conf.sample_advanced.yml | 12 ++ 7 files changed, 320 insertions(+), 54 deletions(-) create mode 100644 lib/galaxy/jobs/runners/util/pykube_util.py create mode 100644 test/integration/test_kubernetes_staging.py diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 7c94edb654b..1dba8c47281 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -14,22 +14,18 @@ from galaxy.jobs.runners import ( AsynchronousJobState, JobState ) +from galaxy.jobs.runners.util.pykube_util import ( + DEFAULT_JOB_API_VERSION, + ensure_pykube, + Job, + job_object_dict, + Pod, + produce_unique_k8s_job_name, + pull_policy, + pykube_client_from_dict, +) from galaxy.util.bytesize import ByteSize -# pykube imports: -try: - from pykube.config import KubeConfig - from pykube.http import HTTPClient - from pykube.objects import ( - Job, - Pod - ) -except ImportError as exc: - KubeConfig = None - K8S_IMPORT_MESSAGE = ('The Python pykube package is required to use ' - 'this feature, please install it or correct the ' - 'following error:\nImportError %s' % str(exc)) - log = logging.getLogger(__name__) __all__ = ('KubernetesJobRunner', ) @@ -43,15 +39,15 @@ class KubernetesJobRunner(AsynchronousJobRunner): def __init__(self, app, nworkers, **kwargs): # Check if pykube was importable, fail if not - assert KubeConfig is not None, K8S_IMPORT_MESSAGE + ensure_pykube() runner_param_specs = dict( - k8s_config_path=dict(map=str, default=os.environ.get('KUBECONFIG', None)), + k8s_config_path=dict(map=str, default=None), k8s_use_service_account=dict(map=bool, default=False), k8s_persistent_volume_claims=dict(map=str), k8s_namespace=dict(map=str, default="default"), k8s_galaxy_instance_id=dict(map=str), k8s_timeout_seconds_job_deletion=dict(map=int, valid=lambda x: int > 0, default=30), - k8s_job_api_version=dict(map=str, default="batch/v1"), + k8s_job_api_version=dict(map=str, default=DEFAULT_JOB_API_VERSION), k8s_supplemental_group_id=dict(map=str), k8s_pull_policy=dict(map=str, default="Default"), k8s_run_as_user_id=dict(map=str, valid=lambda s: s == "$uid" or s.isdigit()), @@ -71,11 +67,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): """Start the job runner parent object """ super(KubernetesJobRunner, self).__init__(app, nworkers, **kwargs) - if "k8s_use_service_account" in self.runner_params and self.runner_params["k8s_use_service_account"]: - self._pykube_api = HTTPClient(KubeConfig.from_service_account()) - else: - self._pykube_api = HTTPClient(KubeConfig.from_file(self.runner_params["k8s_config_path"])) - + self._pykube_api = pykube_client_from_dict(self.runner_params) self._galaxy_instance_id = self.__get_galaxy_instance_id() self._run_as_user_id = self.__get_run_as_user_id() @@ -122,18 +114,11 @@ class KubernetesJobRunner(AsynchronousJobRunner): # Construction of the Kubernetes Job object follows: http://kubernetes.io/docs/user-guide/persistent-volumes/ k8s_job_name = self.__produce_unique_k8s_job_name(job_wrapper.get_id_tag()) - k8s_job_obj = { - "apiVersion": self.runner_params['k8s_job_api_version'], - "kind": "Job", - "metadata": { - # metadata.name is the name of the pod resource created, and must be unique - # http://kubernetes.io/docs/user-guide/configuring-containers/ - "name": k8s_job_name, - "namespace": self.runner_params['k8s_namespace'], - "labels": {"app": k8s_job_name} - }, - "spec": self.__get_k8s_job_spec(ajs) - } + k8s_job_obj = job_object_dict( + self.runner_params, + k8s_job_name, + self.__get_k8s_job_spec(ajs) + ) # Checks if job exists and is trusted, or if it needs re-creation. job = Job(self._pykube_api, k8s_job_obj) @@ -168,10 +153,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): self.monitor_queue.put(ajs) def __get_pull_policy(self): - if "k8s_pull_policy" in self.runner_params: - if self.runner_params['k8s_pull_policy'] in ["Always", "IfNotPresent", "Never"]: - return self.runner_params['k8s_pull_policy'] - return None + return pull_policy(self.runner_params) def __get_run_as_user_id(self): if "k8s_run_as_user_id" in self.runner_params: @@ -234,10 +216,8 @@ class KubernetesJobRunner(AsynchronousJobRunner): def __produce_unique_k8s_job_name(self, galaxy_internal_job_id): # wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers. - instance_id = "" - if self._galaxy_instance_id and len(self._galaxy_instance_id) > 0: - instance_id = self._galaxy_instance_id + "-" - return "galaxy-" + instance_id + galaxy_internal_job_id + instance_id = self._galaxy_instance_id or '' + return produce_unique_k8s_job_name(app_prefix='galaxy', instance_id=instance_id, job_id=galaxy_internal_job_id) def __get_k8s_job_spec(self, ajs): """Creates the k8s Job spec. For a Job spec, the only requirement is to have a .spec.template.""" diff --git a/lib/galaxy/jobs/runners/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py index eda05554b6f..13189d5179b 100644 --- a/lib/galaxy/jobs/runners/pulsar.py +++ b/lib/galaxy/jobs/runners/pulsar.py @@ -56,6 +56,7 @@ __all__ = ( MINIMUM_PULSAR_VERSIONS = { '_default_': packaging.version.parse("0.7.0.dev3"), 'remote_metadata': packaging.version.parse("0.8.0"), + 'remote_container_handling': packaging.version.parse("0.9.1.dev0") # probably 0.10 ultimately? } NO_REMOTE_GALAXY_FOR_METADATA_MESSAGE = "Pulsar misconfiguration - Pulsar client configured to set metadata remotely, but remote Pulsar isn't properly configured with a galaxy_home directory." @@ -84,6 +85,10 @@ PULSAR_PARAM_SPECS = dict( map=specs.to_bool_or_none, default=None, ), + remote_container_handling=dict( + map=specs.to_bool, + default=False, + ), amqp_url=dict( map=specs.to_str_or_none, default=None, @@ -273,7 +278,7 @@ class PulsarJobRunner(AsynchronousJobRunner): job_destination = job_wrapper.job_destination self._populate_parameter_defaults(job_destination) - command_line, client, remote_job_config, compute_environment = self.__prepare_job(job_wrapper, job_destination) + command_line, client, remote_job_config, compute_environment, remote_container = self.__prepare_job(job_wrapper, job_destination) if not command_line: return @@ -293,6 +298,7 @@ class PulsarJobRunner(AsynchronousJobRunner): else: metadata_directory = os.path.join(job_wrapper.working_directory, "metadata") + remote_pulsar_app_config = job_destination.params.get("pulsar_app_config", {}) client_job_description = ClientJobDescription( command_line=command_line, input_files=self.get_input_files(job_wrapper), @@ -306,6 +312,8 @@ class PulsarJobRunner(AsynchronousJobRunner): rewrite_paths=rewrite_paths, arbitrary_files=unstructured_path_rewrites, touch_outputs=output_names, + remote_pulsar_app_config=remote_pulsar_app_config, + container=None if not remote_container else remote_container.container_id, ) job_id = pulsar_submit_job(client, client_job_description, remote_job_config) log.info("Pulsar job submitted with job_id %s" % job_id) @@ -327,6 +335,7 @@ class PulsarJobRunner(AsynchronousJobRunner): def __needed_features(self, client): return { 'remote_metadata': PulsarJobRunner.__remote_metadata(client), + 'remote_container_handling': PulsarJobRunner.__remote_container_handling(client), } def __prepare_job(self, job_wrapper, job_destination): @@ -335,12 +344,20 @@ class PulsarJobRunner(AsynchronousJobRunner): client = None remote_job_config = None compute_environment = None + remote_container = None fail_or_resubmit = False try: client = self.get_client_from_wrapper(job_wrapper) tool = job_wrapper.tool remote_job_config = client.setup(tool.id, tool.version, tool.requires_galaxy_python_environment) + remote_container_handling = PulsarJobRunner.__remote_container_handling(client) + if remote_container_handling: + # Handle this remotely and don't pass it to build_command + remote_container = self._find_container( + job_wrapper, + ) + needed_features = self.__needed_features(client) PulsarJobRunner.check_job_config(remote_job_config, check_features=needed_features) rewrite_parameters = PulsarJobRunner.__rewrite_parameters(client) @@ -363,13 +380,15 @@ class PulsarJobRunner(AsynchronousJobRunner): # calculated at some other level. remote_job_directory = os.path.abspath(os.path.join(remote_working_directory, os.path.pardir)) remote_tool_directory = os.path.abspath(os.path.join(remote_job_directory, "tool_files")) - container = self._find_container( - job_wrapper, - compute_working_directory=remote_working_directory, - compute_tool_directory=remote_tool_directory, - compute_job_directory=remote_job_directory, - ) job_wrapper.disable_commands_in_new_shell() + container = None + if remote_container is None: + container = self._find_container( + job_wrapper, + compute_working_directory=remote_working_directory, + compute_tool_directory=remote_tool_directory, + compute_job_directory=remote_job_directory, + ) # Pulsar handles ``create_tool_working_directory`` and # ``include_work_dir_outputs`` details. @@ -399,7 +418,7 @@ class PulsarJobRunner(AsynchronousJobRunner): job_state = self._job_state(job_wrapper.get_job(), job_wrapper) self.work_queue.put((self.fail_job, job_state)) - return command_line, client, remote_job_config, compute_environment + return command_line, client, remote_job_config, compute_environment, remote_container def __prepare_input_files_locally(self, job_wrapper): """Run task splitting commands locally.""" @@ -675,6 +694,11 @@ class PulsarJobRunner(AsynchronousJobRunner): remote_metadata = string_as_bool_or_none(pulsar_client.destination_params.get("remote_metadata", False)) return remote_metadata + @staticmethod + def __remote_container_handling(pulsar_client): + remote_container_handling = string_as_bool_or_none(pulsar_client.destination_params.get("remote_container_handling", False)) + return remote_container_handling + @staticmethod def __use_remote_datatypes_conf(pulsar_client): """Use remote metadata datatypes instead of Galaxy's. @@ -786,6 +810,36 @@ class PulsarMQJobRunner(PulsarJobRunner): # Nothing else to do? - Attempt to fail the job? +class PulsarKubernetesCoexecutionJobRunner(PulsarMQJobRunner): + """Flavor of Pulsar job runner with sensible defaults for Kubernetes Pod co-execution.""" + + destination_defaults = dict( + default_file_action="remote_transfer", + rewrite_parameters="true", + dependency_resolution="none", + jobs_directory="/pulsar_staging", + pulsar_container_image="galaxy/pulsar-pod-staging:0.10.0", + remote_container_handling=True, + k8s_enabled=True, + url=PARAMETER_SPECIFICATION_IGNORED, + private_token=PARAMETER_SPECIFICATION_IGNORED, + ) + + def _populate_parameter_defaults(self, job_destination): + super(PulsarKubernetesCoexecutionJobRunner, self)._populate_parameter_defaults(job_destination) + params = job_destination.params + # Set some sensible defaults for Pulsar application that runs in staging container. + if "pulsar_app_config" not in params: + params["pulsar_app_config"] = {} + pulsar_app_config = params["pulsar_app_config"] + if "staging_directory" not in pulsar_app_config: + # coexecution always uses a fixed path for staging directory + pulsar_app_config["staging_directory"] = params.get("jobs_directory") + if "manager" not in pulsar_app_config and "managers" not in pulsar_app_config: + # coexecution always uses coexecution manager + pulsar_app_config["manager"] = {"type": "coexecution"} + + class PulsarRESTJobRunner(PulsarJobRunner): """Flavor of Pulsar job runner with sensible defaults for RESTful usage.""" diff --git a/lib/galaxy/jobs/runners/util/pykube_util.py b/lib/galaxy/jobs/runners/util/pykube_util.py new file mode 100644 index 00000000000..a86601c324f --- /dev/null +++ b/lib/galaxy/jobs/runners/util/pykube_util.py @@ -0,0 +1,89 @@ +"""Interface layer for pykube library shared between Galaxy and Pulsar.""" +import os +import uuid + +try: + from pykube.config import KubeConfig + from pykube.http import HTTPClient + from pykube.objects import ( + Job, + Pod + ) +except ImportError as exc: + KubeConfig = None + Job = None + Pod = None + K8S_IMPORT_MESSAGE = ('The Python pykube package is required to use ' + 'this feature, please install it or correct the ' + 'following error:\nImportError %s' % str(exc)) + +DEFAULT_JOB_API_VERSION = "batch/v1" +DEFAULT_NAMESPACE = "default" + + +def ensure_pykube(): + if KubeConfig is None: + raise Exception(K8S_IMPORT_MESSAGE) + + +def pykube_client_from_dict(params): + if "k8s_use_service_account" in params and params["k8s_use_service_account"]: + pykube_client = HTTPClient(KubeConfig.from_service_account()) + else: + config_path = params.get("k8s_config_path") + if config_path is None: + config_path = os.environ.get('KUBECONFIG', None) + if config_path is None: + config_path = '~/.kube/config' + pykube_client = HTTPClient(KubeConfig.from_file(config_path)) + return pykube_client + + +def produce_unique_k8s_job_name(app_prefix=None, instance_id=None, job_id=None): + if job_id is None: + job_id = str(uuid.uuid4()) + + job_name = "" + if app_prefix: + job_name += "%s-" % app_prefix + + if instance_id and instance_id > 0: + job_name += "%s-" % instance_id + + return job_name + job_id + + +def pull_policy(params): + # If this doesn't validate it returns None, that seems odd? + if "k8s_pull_policy" in params: + if params['k8s_pull_policy'] in ["Always", "IfNotPresent", "Never"]: + return params['k8s_pull_policy'] + return None + + +def job_object_dict(params, job_name, spec): + k8s_job_obj = { + "apiVersion": params.get('k8s_job_api_version', DEFAULT_JOB_API_VERSION), + "kind": "Job", + "metadata": { + # metadata.name is the name of the pod resource created, and must be unique + # http://kubernetes.io/docs/user-guide/configuring-containers/ + "name": job_name, + "namespace": params.get('k8s_namespace', DEFAULT_NAMESPACE), + "labels": {"app": job_name} + }, + "spec": spec, + } + return k8s_job_obj + + +__all__ = ( + "DEFAULT_JOB_API_VERSION", + "ensure_pykube", + "Job", + "job_object_dict", + "Pod", + "produce_unique_k8s_job_name", + "pull_policy", + "pykube_client_from_dict", +) diff --git a/test/base/integration_util.py b/test/base/integration_util.py index d137d73b618..24724825829 100644 --- a/test/base/integration_util.py +++ b/test/base/integration_util.py @@ -17,8 +17,11 @@ from .driver_util import GalaxyTestDriver NO_APP_MESSAGE = "test_case._app called though no Galaxy has been configured." -def skip_if_jenkins(cls): +def _identity(func): + return func + +def skip_if_jenkins(cls): if os.environ.get("BUILD_NUMBER", ""): return skip @@ -27,7 +30,7 @@ def skip_if_jenkins(cls): def skip_unless_executable(executable): if which(executable): - return lambda func: func + return _identity return skip("PATH doesn't contain executable %s" % executable) @@ -39,6 +42,17 @@ def skip_unless_kubernetes(): return skip_unless_executable("kubectl") +def k8s_config_path(): + return os.environ.get('GALAXY_TEST_KUBE_CONFIG_PATH', '~/.kube/config') + + +def skip_unless_fixed_port(): + if os.environ.get("GALAXY_TEST_PORT"): + return _identity + + return skip("GALAXY_TEST_PORT must be set for this test.") + + class IntegrationInstance(UsesApiTestCaseMixin): """Unit test case with utilities for spinning up Galaxy.""" diff --git a/test/integration/test_kubernetes_runner.py b/test/integration/test_kubernetes_runner.py index 34e469ee716..6546ac3a20a 100644 --- a/test/integration/test_kubernetes_runner.py +++ b/test/integration/test_kubernetes_runner.py @@ -1,4 +1,4 @@ -"""Integration tests for the CLI shell plugins and runners.""" +"""Integration tests for the Kubernetes runner.""" # Tested on docker for mac 18.06.1-ce-mac73 using the default kubernetes setup import collections import json @@ -109,7 +109,7 @@ def job_config(jobs_directory): """) job_conf_str = job_conf_template.substitute(jobs_directory=jobs_directory, tool_directory=TOOL_DIR, - k8s_config_path=os.environ.get('GALAXY_TEST_KUBE_CONFIG_PATH', '~/.kube/config'), + k8s_config_path=integration_util.k8s_config_path(), ) with tempfile.NamedTemporaryFile(suffix="_kubernetes_integration_job_conf.xml", mode="w", delete=False) as job_conf: job_conf.write(job_conf_str) diff --git a/test/integration/test_kubernetes_staging.py b/test/integration/test_kubernetes_staging.py new file mode 100644 index 00000000000..cdbff077d3d --- /dev/null +++ b/test/integration/test_kubernetes_staging.py @@ -0,0 +1,117 @@ +"""Integration tests for Kubernetes pod staging. + +This case is way more brittle than the other integration tests distributed with Galaxy. +This is because it requires Kubernetes, RabbitMQ, and the test itself needs to know what +the test host IP address will be relative to a container inside Kubernetes in order to +communicate job status updates back. + +For this reason, this test will only work out of the box currently with Docker for Mac, +rabbitmq installed via Homebrew, and if a fixed port is set for the test. + + GALAXY_TEST_PORT=9234 pytest test/integration/test_kubernetes_staging.py + +""" +import os +import string +import tempfile + +from base import integration_util # noqa: I100,I202 +from base.populators import skip_without_tool +from .test_containerized_jobs import MulledJobTestCases # noqa: I201 +from .test_job_environments import BaseJobEnvironmentIntegrationTestCase # noqa: I201 + +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", "docker.for.mac.localhost") + + +def job_config(jobs_directory): + job_conf_template = string.Template(""" +runners: + local: + load: galaxy.jobs.runners.local:LocalJobRunner + workers: 1 + pulsar_k8s: + load: galaxy.jobs.runners.pulsar:PulsarKubernetesCoexecutionJobRunner + galaxy_url: ${infrastructure_url} + amqp_url: ${amqp_url} + +execution: + default: pulsar_k8s_environment + environments: + pulsar_k8s_environment: + k8s_config_path: ${k8s_config_path} + runner: pulsar_k8s + docker_enabled: true + docker_default_container_id: busybox:ubuntu-14.04 + pulsar_app_config: + message_queue_url: '${container_amqp_url}' + env: + - name: SOME_ENV_VAR + value: '42' + local_environment: + runner: local +tools: + - id: upload1 + environment: local_environment +""") + 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) + return job_conf.name + + +@integration_util.skip_unless_kubernetes() +@integration_util.skip_unless_fixed_port() +class BaseKubernetesStagingIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase, MulledJobTestCases): + + def setUp(self): + super(BaseKubernetesStagingIntegrationTestCase, 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(BaseKubernetesStagingIntegrationTestCase, cls).setUpClass() + + @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(cls.jobs_directory) + + config["default_job_shell"] = '/bin/sh' + # Disable tool dependency resolution. + config["tool_dependency_dir"] = "none" + config["enable_beta_mulled_containers"] = "true" + + @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' + + +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 + # similar code found in Pulsar integration_tests.py. + infrastructure_uri = uri + if GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST: + if "0.0.0.0" in infrastructure_uri: + infrastructure_uri = infrastructure_uri.replace("0.0.0.0", GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST) + elif "localhost" in infrastructure_uri: + infrastructure_uri = infrastructure_uri.replace("localhost", GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST) + elif "127.0.0.1" in infrastructure_uri: + infrastructure_uri = infrastructure_uri.replace("127.0.0.1", GALAXY_TEST_KUBERNETES_INFRASTRUCTURE_HOST) + return infrastructure_uri diff --git a/test/unit/jobs/job_conf.sample_advanced.yml b/test/unit/jobs/job_conf.sample_advanced.yml index 7426838b7f3..9362d95d2c8 100644 --- a/test/unit/jobs/job_conf.sample_advanced.yml +++ b/test/unit/jobs/job_conf.sample_advanced.yml @@ -125,6 +125,13 @@ runners: # higher value (in seconds) (or `None` to use blocking connections). #amqp_consumer_timeout: None + pulsar_k8s: + # Use a two-container Kubernetes to run jobs - one for staging and one + # for tool execution. + load: galaxy.jobs.runners.pulsar:PulsarKubernetesCoexecutionJobRunner + galaxy_url: http://docker.for.mac.localhost:8080 + amqp_url: amqp://guest:guest@localhost:5672// + pulsar_legacy: # Pulsar job runner with default parameters matching those # of old LWR job runner. If your Pulsar server is running on a @@ -437,6 +444,11 @@ execution: # all the same destination parameters as the RESTful client documented # above (though url and private_token are ignored when using a MQ). + pulsar_k8s_environment: + runner: pulsar_k8s + docker_enabled: true # probably shouldn't be needed but is still + docker_default_container_id: 'conda/miniconda2' + # Example CLI runners. ssh_torque: runner: cli