diff --git a/lib/galaxy/config/__init__.py b/lib/galaxy/config/__init__.py index 6f529186bff..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("interactivetool_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/dependencies/pipfiles/default/pinned-requirements.txt b/lib/galaxy/dependencies/pipfiles/default/pinned-requirements.txt index 25237754d4b..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.dev1 +pulsar-galaxy-lib==0.14.0.dev3 pyasn1-modules==0.2.7 pyasn1==0.4.8 pycparser==2.19 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/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index ad68c60db77..aaac45ebf11 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -19,12 +19,15 @@ 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, Pod, produce_unique_k8s_job_name, pull_policy, pykube_client_from_dict, + stop_job, ) from galaxy.util.bytesize import ByteSize @@ -201,25 +204,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. @@ -509,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 [] @@ -550,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/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py index ebd72250e58..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", @@ -868,7 +889,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, @@ -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 0203b9f5046..f2e06f78c59 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,8 +19,13 @@ 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" +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(): @@ -61,6 +68,44 @@ def pull_policy(params): return None +def find_job_object_by_name(pykube_api, job_name, namespace=None): + 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 + + 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"): + 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), @@ -77,13 +122,38 @@ 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: + raw_value = params['k8s_galaxy_instance_id'] + if re.match(r"(?!-)[a-z\d-]{1,20}(?