diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index c3204b34d95..df3c4d3165d 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -75,7 +75,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", @@ -91,7 +91,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 @@ -103,9 +110,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.""" @@ -118,7 +125,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), @@ -204,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']: @@ -223,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. See stdout for pods details." self.mark_as_failed(job_state) job.scale(replicas=0) return None @@ -260,14 +275,46 @@ 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 = "" 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 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") if isinstance(logs, text_type): @@ -279,8 +326,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