From ca82e6d746a685209b9585b3f69128a55198f801 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Sat, 25 Jun 2016 15:05:04 +0100 Subject: [PATCH 1/8] Deletes existing jobs with same ID when re-submitting. --- lib/galaxy/jobs/runners/kubernetes.py | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 10c70244cc4..338e909fe14 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -89,7 +89,14 @@ class KubernetesJobRunner(AsynchronousJobRunner): "spec": self.__get_k8s_job_spec(job_wrapper) } + # Checks if job exists + job = Job(self._pykube_api, k8s_job_obj) + if job.exists(): + job.delete() # Creates the Kubernetes Job + # TODO if a job with that ID exists, what should we do? + # TODO do we trust that this is the same job and use that? + # TODO or create a new job as we cannot make sure Job(self._pykube_api, k8s_job_obj).create() # define job attributes in the AsyncronousJobState for follow-up From 01fd25b9f4b79b4f2868efc51ad76e19fc132b5d Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Sat, 25 Jun 2016 15:16:36 +0100 Subject: [PATCH 2/8] Changes use of produce_unique_k8s_job_name to be able to use both job and job_wrappers. --- lib/galaxy/jobs/runners/kubernetes.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 338e909fe14..230ba639e01 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -73,7 +73,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): job_destination = job_wrapper.job_destination # 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) + k8s_job_name = self.__produce_unique_k8s_job_name(job_wrapper.get_id_tag()) k8s_job_obj = { "apiVersion": "extensions/v1beta1", "kind": "Job", @@ -108,9 +108,9 @@ class KubernetesJobRunner(AsynchronousJobRunner): external_runjob_script = None return external_runjob_script - def __produce_unique_k8s_job_name(self, job_wrapper): + def __produce_unique_k8s_job_name(self, galaxy_internal_job_id): # wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers. - return "galaxy-" + job_wrapper.get_id_tag() + return "galaxy-" + galaxy_internal_job_id def __get_k8s_job_spec(self, job_wrapper): """Creates the k8s Job spec. For a Job spec, the only requirement is to have a .spec.template.""" @@ -123,7 +123,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): (see pod selector) and an appropriate restart policy.""" k8s_spec_template = { "metadata": { - "labels": {"app": self.__produce_unique_k8s_job_name(job_wrapper)} + "labels": {"app": self.__produce_unique_k8s_job_name(job_wrapper.get_id_tag())} }, "spec": { "volumes": self.__get_k8s_mountable_volumes(job_wrapper), From 409eb373cd9233f7925c64e9c5716859a2363ba1 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Sat, 25 Jun 2016 15:17:58 +0100 Subject: [PATCH 3/8] Fixes stop job method so that it correctly scales down the Kubernetes job. --- lib/galaxy/jobs/runners/kubernetes.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 230ba639e01..bc69aa86edb 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -284,8 +284,9 @@ class KubernetesJobRunner(AsynchronousJobRunner): def stop_job(self, job): """Attempts to delete a dispatched job to the k8s cluster""" try: - jobs = Job.objects(self._pykube_api).filter(selector="app=" + job.job_runner_external_id) - if jobs.response['items'].len() >= 0: + jobs = Job.objects(self._pykube_api).filter(selector="app=" + + self.__produce_unique_k8s_job_name(job.get_id_tag())) + if len(jobs.response['items']) >= 0: job_to_delete = Job(self._pykube_api, jobs.response['items'][0]) job_to_delete.scale(replicas=0) # TODO assert whether job parallelism == 0 From 65a937965a22c46922f9b731475561e38809e588 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Sun, 26 Jun 2016 08:21:24 +0100 Subject: [PATCH 4/8] Improves handling of Kubernetes logs --- lib/galaxy/jobs/runners/kubernetes.py | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index bc69aa86edb..aea9c8f45f8 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -4,6 +4,8 @@ Offload jobs to a Kubernetes cluster. import logging +from pykube.exceptions import HTTPError + from galaxy import model from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner from os import environ as os_environ @@ -269,10 +271,14 @@ class KubernetesJobRunner(AsynchronousJobRunner): pod_r = Pod.objects(self._pykube_api).filter(selector="app=" + job_state.job_id) logs = "" for pod_obj in pod_r.response['items']: - pod = Pod(self._pykube_api, pod_obj) - logs += "\n\n==== Pod " + pod.name + " log start ====\n\n" - logs += pod.logs(timestamps=True) - logs += "\n\n==== Pod " + pod.name + " log end ====" + try: + pod = Pod(self._pykube_api, pod_obj) + logs += "\n\n==== Pod " + pod.name + " log start ====\n\n" + logs += pod.logs(timestamps=True) + logs += "\n\n==== Pod " + pod.name + " log end ====" + except Exception as detail: + log.info("Could not write pods "+job_state.job_id+" log file due to HTTPError "+str(detail)) + logs_file_path = job_state.output_file logs_file = open(logs_file_path, mode="w") if isinstance(logs, text_type): From e8bf2df21fe7fb7405ea283a41eae22e6d58509f Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Sun, 26 Jun 2016 08:26:15 +0100 Subject: [PATCH 5/8] Corrects usage of pods retrials parameters set in runner and container destinations. --- lib/galaxy/jobs/runners/kubernetes.py | 14 +++++++++++--- 1 file changed, 11 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index aea9c8f45f8..2bb058e8dc6 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -211,9 +211,17 @@ class KubernetesJobRunner(AsynchronousJobRunner): jobs = Job.objects(self._pykube_api).filter(selector="app=" + job_state.job_id) if len(jobs.response['items']) == 1: job = Job(self._pykube_api, jobs.response['items'][0]) + job_destination = job_state.job_wrapper.job_destination succeeded = 0 active = 0 failed = 0 + + max_pod_retrials = 1 + if 'k8s_pod_retrials' in self.runner_params: + max_pod_retrials = int(self.runner_params['k8s_pod_retrials']) + if 'max_pod_retrials' in job_destination.params: + max_pod_retrials = int(job_destination.params['max_pod_retrials']) + if 'succeeded' in job.obj['status']: succeeded = job.obj['status']['succeeded'] if 'active' in job.obj['status']: @@ -230,16 +238,16 @@ class KubernetesJobRunner(AsynchronousJobRunner): job_state.running = False self.mark_as_finished(job_state) return None - - elif active > 0 or succeeded + active + failed == 0: + elif active > 0 and failed <= max_pod_retrials: job_state.running = True return job_state - elif failed > job_state.job_destination.params['max_pod_retrials']: + elif failed > max_pod_retrials: self.__produce_log_file(job_state) error_file = open(job_state.error_file, 'w') error_file.write("Exceeded max number of Kubernetes pod retrials allowed for job\n") error_file.close() job_state.running = False + job_state.fail_message = "More pods failed than allowed." self.mark_as_failed(job_state) job.scale(replicas=0) return None From 89b5f3f47f14fa82065b83326332313ec8cca629 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Sun, 26 Jun 2016 09:57:07 +0100 Subject: [PATCH 6/8] Uses exception instead of pykube.exceptions.HTTPError to manage missing pod's log exception. --- lib/galaxy/jobs/runners/kubernetes.py | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 2bb058e8dc6..3e221b150b1 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -4,8 +4,6 @@ Offload jobs to a Kubernetes cluster. import logging -from pykube.exceptions import HTTPError - from galaxy import model from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner from os import environ as os_environ @@ -285,7 +283,8 @@ class KubernetesJobRunner(AsynchronousJobRunner): logs += pod.logs(timestamps=True) logs += "\n\n==== Pod " + pod.name + " log end ====" except Exception as detail: - log.info("Could not write pods "+job_state.job_id+" log file due to HTTPError "+str(detail)) + log.info("Could not write pod\'s " + pod_obj['metadata']['name'] + + " log file due to HTTPError "+str(detail)) logs_file_path = job_state.output_file logs_file = open(logs_file_path, mode="w") From 9db3453e6559cb883af22793b741a84f3abd071c Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Sun, 26 Jun 2016 09:59:18 +0100 Subject: [PATCH 7/8] Overrides parent fail_job to save stdout/stderr of pods when Kubernetes job fails (otherwise lost). --- lib/galaxy/jobs/runners/kubernetes.py | 29 ++++++++++++++++++++++++++- 1 file changed, 28 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 3e221b150b1..deb365c9efc 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -245,7 +245,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): error_file.write("Exceeded max number of Kubernetes pod retrials allowed for job\n") error_file.close() job_state.running = False - job_state.fail_message = "More pods failed than allowed." + job_state.fail_message = "More pods failed than allowed. See stdout for pods details." self.mark_as_failed(job_state) job.scale(replicas=0) return None @@ -273,6 +273,33 @@ class KubernetesJobRunner(AsynchronousJobRunner): self.mark_as_failed(job_state) return job_state + def fail_job(self, job_state): + """ + Kubernetes runner overrides fail_job (called by mark_as_failed) to rescue the pod's log files which are left as + stdout (pods logs are the natural stdout and stderr of the running processes inside the pods) and are + deleted in the parent implementation as part of the failing the job process. + + :param job_state: + :return: + """ + + # First we rescue the pods logs + with open(job_state.output_file, 'r') as outfile: + stdout_content = outfile.read() + + if getattr(job_state, 'stop_job', True): + self.stop_job(self.sa_session.query(self.app.model.Job).get(job_state.job_wrapper.job_id)) + self._handle_runner_state('failure', job_state) + # Not convinced this is the best way to indicate this state, but + # something necessary + if not job_state.runner_state_handled: + job_state.job_wrapper.fail( + message=getattr(job_state, 'fail_message', 'Job failed'), + stdout=stdout_content, stderr='See stdout for pod\'s stderr.' + ) + if job_state.job_wrapper.cleanup_job == "always": + job_state.cleanup() + def __produce_log_file(self, job_state): pod_r = Pod.objects(self._pykube_api).filter(selector="app=" + job_state.job_id) logs = "" From 5eb60b3448f2cdf8706aa2ab43e4d51dcfc0a0cb Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 28 Jun 2016 09:59:13 +0100 Subject: [PATCH 8/8] Adds missing whitespace to line in Kubernetes runner. --- lib/galaxy/jobs/runners/kubernetes.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index deb365c9efc..146ebe24300 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -311,7 +311,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): logs += "\n\n==== Pod " + pod.name + " log end ====" except Exception as detail: log.info("Could not write pod\'s " + pod_obj['metadata']['name'] + - " log file due to HTTPError "+str(detail)) + " log file due to HTTPError " + str(detail)) logs_file_path = job_state.output_file logs_file = open(logs_file_path, mode="w")