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.
This commit is contained in:
John Chilton
2019-06-17 11:47:27 -04:00
parent 90f3c078ec
commit ac1b25ed4c
7 changed files with 320 additions and 54 deletions
+22 -42
View File
@@ -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."""
+62 -8
View File
@@ -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."""
@@ -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",
)
+16 -2
View File
@@ -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."""
+2 -2
View File
@@ -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)
+117
View File
@@ -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
@@ -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