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