From afc8bdc81d9feb9d2b6e523fd37a945c83e22a2a Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 15 Mar 2016 12:13:41 +0100 Subject: [PATCH 01/56] Initial commit for collaboration on the k8s runner/job configuration. --- .gitignore | 2 +- config/job_conf_k8s.xml | 43 ++++++ lib/galaxy/jobs/runners/kubernetes.py | 209 ++++++++++++++++++++++++++ 3 files changed, 253 insertions(+), 1 deletion(-) create mode 100644 config/job_conf_k8s.xml create mode 100644 lib/galaxy/jobs/runners/kubernetes.py diff --git a/.gitignore b/.gitignore index 8cabb5cbba8..852d568155c 100644 --- a/.gitignore +++ b/.gitignore @@ -63,7 +63,7 @@ shed_data_manager_conf.xml object_store_conf.xml job_metrics_conf.xml workflow_schedulers_conf.xml -config/* +#config/* static/welcome.html.* static/welcome.html diff --git a/config/job_conf_k8s.xml b/config/job_conf_k8s.xml new file mode 100644 index 00000000000..6a21d439a6f --- /dev/null +++ b/config/job_conf_k8s.xml @@ -0,0 +1,43 @@ + + + + + + + + ok + 0 + ok + 0 + + + /path/to/kubeconfig + + + + + + + + + + + + + + + trackster + + + + + + + diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py new file mode 100644 index 00000000000..6991edd5dc6 --- /dev/null +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -0,0 +1,209 @@ +""" +Job control via a command line interface (e.g. qsub/qstat), possibly over a remote connection (e.g. ssh). +""" + +import logging + +from galaxy import model +from galaxy.jobs import JobDestination +from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner +from .util.cli import CliInterface, split_params + +log = logging.getLogger( __name__ ) + +__all__ = [ 'KubernetesJobRunner' ] + + +class KubernetesJobRunner( AsynchronousJobRunner ): + """ + Job runner backed by a finite pool of worker threads. FIFO scheduling + """ + runner_name = "KubernetesRunner" + + def __init__( self, app, nworkers ): + """Start the job runner """ + super( KubernetesJobRunner, self ).__init__( app, nworkers ) + + self.cli_interface = CliInterface() + self._init_monitor_thread() + self._init_worker_threads() + + def get_cli_plugins( self, shell_params, job_params ): + return self.cli_interface.get_plugins( shell_params, job_params ) + + def url_to_destination( self, url ): + params = {} + shell_params, job_params = url.split( '/' )[ 2:4 ] + # split 'foo=bar&baz=quux' into { 'foo' : 'bar', 'baz' : 'quux' } + shell_params = dict( [ ( 'shell_' + k, v ) for k, v in [ kv.split( '=', 1 ) for kv in shell_params.split( '&' ) ] ] ) + job_params = dict( [ ( 'job_' + k, v ) for k, v in [ kv.split( '=', 1 ) for kv in job_params.split( '&' ) ] ] ) + params.update( shell_params ) + params.update( job_params ) + log.debug( "Converted URL '%s' to destination runner=cli, params=%s" % ( url, params ) ) + # Create a dynamic JobDestination + return JobDestination( runner='cli', params=params ) + + def parse_destination_params( self, params ): + return split_params( params ) + + def queue_job( self, job_wrapper ): + """Create job script and submit it to the DRM""" + # prepare the job + if not self.prepare_job( job_wrapper, include_metadata=True ): + return + + # Get shell and job execution interface + job_destination = job_wrapper.job_destination + shell_params, job_params = self.parse_destination_params(job_destination.params) + shell, job_interface = self.get_cli_plugins(shell_params, job_params) + + # wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers. + galaxy_id_tag = job_wrapper.get_id_tag() + + # define job attributes + ajs = AsynchronousJobState( files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper ) + + job_file_kwargs = job_interface.job_script_kwargs(ajs.output_file, ajs.error_file, ajs.job_name) + script = self.get_job_file( + job_wrapper, + exit_code_path=ajs.exit_code_file, + **job_file_kwargs + ) + + try: + self.write_executable_script( ajs.job_file, script ) + except: + log.exception("(%s) failure writing job script" % galaxy_id_tag ) + job_wrapper.fail("failure preparing job script", exception=True) + return + + # job was deleted while we were preparing it + if job_wrapper.get_state() == model.Job.states.DELETED: + log.info("(%s) Job deleted by user before it entered the queue" % galaxy_id_tag ) + if job_wrapper.cleanup_job in ("always", "onsuccess"): + job_wrapper.cleanup() + return + + log.debug( "(%s) submitting file: %s" % ( galaxy_id_tag, ajs.job_file ) ) + + cmd_out = shell.execute(job_interface.submit(ajs.job_file)) + if cmd_out.returncode != 0: + log.error('(%s) submission failed (stdout): %s' % (galaxy_id_tag, cmd_out.stdout)) + log.error('(%s) submission failed (stderr): %s' % (galaxy_id_tag, cmd_out.stderr)) + job_wrapper.fail("failure submitting job") + return + # Some job runners return something like 'Submitted batch job XXXX' + # Strip and split to get job ID. + external_job_id = cmd_out.stdout.strip().split()[-1] + if not external_job_id: + log.error('(%s) submission did not return a job identifier, failing job' % galaxy_id_tag) + job_wrapper.fail("failure submitting job") + return + + log.info("(%s) queued with identifier: %s" % ( galaxy_id_tag, external_job_id ) ) + + # store runner information for tracking if Galaxy restarts + job_wrapper.set_job_destination( job_destination, external_job_id ) + + # Store state information for job + ajs.job_id = external_job_id + ajs.old_state = 'new' + ajs.job_destination = job_destination + + # Add to our 'queue' of jobs to monitor + self.monitor_queue.put( ajs ) + + def check_watched_items( self ): + """ + Called by the monitor thread to look at each watched job and deal + with state changes. + """ + new_watched = [] + + job_states = self.__get_job_states() + + for ajs in self.watched: + external_job_id = ajs.job_id + id_tag = ajs.job_wrapper.get_id_tag() + old_state = ajs.old_state + state = job_states.get(external_job_id, None) + if state is None: + if ajs.job_wrapper.get_state() == model.Job.states.DELETED: + continue + log.debug("(%s/%s) job not found in batch state check" % ( id_tag, external_job_id ) ) + shell_params, job_params = self.parse_destination_params(ajs.job_destination.params) + shell, job_interface = self.get_cli_plugins(shell_params, job_params) + cmd_out = shell.execute(job_interface.get_single_status(external_job_id)) + state = job_interface.parse_single_status(cmd_out.stdout, external_job_id) + if state == model.Job.states.OK: + log.debug('(%s/%s) job execution finished, running job wrapper finish method' % ( id_tag, external_job_id ) ) + self.work_queue.put( ( self.finish_job, ajs ) ) + continue + else: + log.warning('(%s/%s) job not found in batch state check, but found in individual state check' % ( id_tag, external_job_id ) ) + if state != old_state: + ajs.job_wrapper.change_state( state ) + else: + if state != old_state: + log.debug("(%s/%s) state change: %s" % ( id_tag, external_job_id, state ) ) + ajs.job_wrapper.change_state( state ) + if state == model.Job.states.RUNNING and not ajs.running: + ajs.running = True + ajs.job_wrapper.change_state( model.Job.states.RUNNING ) + ajs.old_state = state + new_watched.append( ajs ) + # Replace the watch list with the updated version + self.watched = new_watched + + def __get_job_states(self): + job_destinations = {} + job_states = {} + # unique the list of destinations + for ajs in self.watched: + if ajs.job_destination.id not in job_destinations: + job_destinations[ajs.job_destination.id] = dict( job_destination=ajs.job_destination, job_ids=[ ajs.job_id ] ) + else: + job_destinations[ajs.job_destination.id]['job_ids'].append( ajs.job_id ) + # check each destination for the listed job ids + for job_destination_id, v in job_destinations.items(): + job_destination = v['job_destination'] + job_ids = v['job_ids'] + shell_params, job_params = self.parse_destination_params(job_destination.params) + shell, job_interface = self.get_cli_plugins(shell_params, job_params) + cmd_out = shell.execute(job_interface.get_status(job_ids)) + assert cmd_out.returncode == 0, cmd_out.stderr + job_states.update(job_interface.parse_status(cmd_out.stdout, job_ids)) + return job_states + + def stop_job( self, job ): + """Attempts to delete a dispatched job""" + try: + shell_params, job_params = self.parse_destination_params(job.destination_params) + shell, job_interface = self.get_cli_plugins(shell_params, job_params) + cmd_out = shell.execute(job_interface.delete( job.job_runner_external_id )) + assert cmd_out.returncode == 0, cmd_out.stderr + log.debug( "(%s/%s) Terminated at user's request" % ( job.id, job.job_runner_external_id ) ) + except Exception as e: + log.debug( "(%s/%s) User killed running job, but error encountered during termination: %s" % ( job.id, job.job_runner_external_id, e ) ) + + def recover( self, job, job_wrapper ): + """Recovers jobs stuck in the queued/running state when Galaxy started""" + job_id = job.get_job_runner_external_id() + if job_id is None: + self.put( job_wrapper ) + return + ajs = AsynchronousJobState( files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper ) + ajs.job_id = str( job_id ) + ajs.command_line = job.command_line + ajs.job_wrapper = job_wrapper + ajs.job_destination = job_wrapper.job_destination + if job.state == model.Job.states.RUNNING: + log.debug( "(%s/%s) is still in running state, adding to the runner monitor queue" % ( job.id, job.job_runner_external_id ) ) + ajs.old_state = model.Job.states.RUNNING + ajs.running = True + self.monitor_queue.put( ajs ) + elif job.state == model.Job.states.QUEUED: + log.debug( "(%s/%s) is still in queued state, adding to the runner monitor queue" % ( job.id, job.job_runner_external_id ) ) + ajs.old_state = model.Job.states.QUEUED + ajs.running = False + self.monitor_queue.put( ajs ) From fa3da5a5194cc5317d658a1bab36f824e2b73b93 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 15 Mar 2016 15:32:12 +0100 Subject: [PATCH 02/56] Adds imports for pykube --- lib/galaxy/jobs/runners/kubernetes.py | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 6991edd5dc6..3a8ee10be36 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -9,6 +9,15 @@ from galaxy.jobs import JobDestination from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner from .util.cli import CliInterface, split_params +# pykube imports: +import operator + +from pykube.config import KubeConfig +from pykube.http import HTTPClient +from pykube.objects import Pod + + + log = logging.getLogger( __name__ ) __all__ = [ 'KubernetesJobRunner' ] From eff6d21a516174e619bc3d12b7b78b378dff6526 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 15 Mar 2016 15:37:15 +0100 Subject: [PATCH 03/56] Changes import to handle missing dependency --- 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 3a8ee10be36..7c51a463a68 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -10,11 +10,17 @@ from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner from .util.cli import CliInterface, split_params # pykube imports: -import operator +try: + import operator -from pykube.config import KubeConfig -from pykube.http import HTTPClient -from pykube.objects import Pod + from pykube.config import KubeConfig + from pykube.http import HTTPClient + from pykube.objects import Pod +except ImportError as exc: + operator = None + K8S_IMPORT_MESSAGE = ('The Python pbs-python package is required to use ' + 'this feature, please install it or correct the ' + 'following error:\nImportError %s' % str(exc)) From 9c4ef02e9ff50a173b8531668ab042bf8cbc9ed1 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 15 Mar 2016 16:38:15 +0100 Subject: [PATCH 04/56] Sample config file for k8s based execution. --- config/job_conf_k8s.xml | 17 ++++++----------- 1 file changed, 6 insertions(+), 11 deletions(-) diff --git a/config/job_conf_k8s.xml b/config/job_conf_k8s.xml index 6a21d439a6f..11c4d2575f0 100644 --- a/config/job_conf_k8s.xml +++ b/config/job_conf_k8s.xml @@ -1,17 +1,8 @@ - + - - - ok - 0 - ok - 0 - /path/to/kubeconfig @@ -22,7 +13,10 @@ - + docker-registry.local:6000 + fatcat + blankfilter + latest @@ -39,5 +33,6 @@ and pass to dynamic destination (as resource_params argument). --> + From f2ef8de953930906913dbff88561e98b6800fe90 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 15 Mar 2016 16:42:02 +0100 Subject: [PATCH 05/56] First trial at the __init__ --- lib/galaxy/jobs/runners/kubernetes.py | 16 ++++++++++++---- 1 file changed, 12 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 7c51a463a68..ee09a60044a 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -36,12 +36,20 @@ class KubernetesJobRunner( AsynchronousJobRunner ): runner_name = "KubernetesRunner" def __init__( self, app, nworkers ): - """Start the job runner """ + # Check if pykube was importable, fail if not + assert operator is not None, K8S_IMPORT_MESSAGE + """Start the job runner parent object """ super( KubernetesJobRunner, self ).__init__( app, nworkers ) - self.cli_interface = CliInterface() - self._init_monitor_thread() - self._init_worker_threads() + # self.cli_interface = CliInterface() + + # here we need to fetch the default kubeconfig path from the plugin defined in job_conf... + self._pykube_api = HTTPClient(KubeConfig.from_file(fromJobConfPluginParams)) + # TODO how do we read from the config file the plugin parameters + + # TODO do we need these? + # self._init_monitor_thread() + # self._init_worker_threads() def get_cli_plugins( self, shell_params, job_params ): return self.cli_interface.get_plugins( shell_params, job_params ) From 97af6466cdfc148e4c1f1424d24319a18e29f8bd Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 6 Apr 2016 17:59:09 +0100 Subject: [PATCH 06/56] Configures pykube to use the kubernetes config file path set in the job_conf.xml, imports Job instead of Pod. --- 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 ee09a60044a..f5456be7c02 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -15,7 +15,7 @@ try: from pykube.config import KubeConfig from pykube.http import HTTPClient - from pykube.objects import Pod + from pykube.objects import Job except ImportError as exc: operator = None K8S_IMPORT_MESSAGE = ('The Python pbs-python package is required to use ' @@ -44,8 +44,7 @@ class KubernetesJobRunner( AsynchronousJobRunner ): # self.cli_interface = CliInterface() # here we need to fetch the default kubeconfig path from the plugin defined in job_conf... - self._pykube_api = HTTPClient(KubeConfig.from_file(fromJobConfPluginParams)) - # TODO how do we read from the config file the plugin parameters + self._pykube_api = HTTPClient(KubeConfig.from_file(self.runner_params["k8s_config_path"])) # TODO do we need these? # self._init_monitor_thread() From 63842580dc62070128637ea3bde261040e78e4ac Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Thu, 7 Apr 2016 09:54:14 +0100 Subject: [PATCH 07/56] Adds persistent volume claim name example in config. --- config/job_conf_k8s.xml | 1 + 1 file changed, 1 insertion(+) diff --git a/config/job_conf_k8s.xml b/config/job_conf_k8s.xml index 11c4d2575f0..9cc7da41cb2 100644 --- a/config/job_conf_k8s.xml +++ b/config/job_conf_k8s.xml @@ -5,6 +5,7 @@ /path/to/kubeconfig + galaxy_pvc From d5abfb6eb5182c4aa7e48427c2440edd771c3ff3 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Thu, 7 Apr 2016 09:57:30 +0100 Subject: [PATCH 08/56] Initial commit of queue_job plus required helper methods (not finished). --- lib/galaxy/jobs/runners/kubernetes.py | 124 +++++++++++++++++--------- 1 file changed, 80 insertions(+), 44 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index f5456be7c02..9cefe923729 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -69,71 +69,107 @@ class KubernetesJobRunner( AsynchronousJobRunner ): return split_params( params ) def queue_job( self, job_wrapper ): - """Create job script and submit it to the DRM""" + """Create job script and submit it to Kubernetes cluster""" # prepare the job + # TODO understand weather we need include_metadata and include_work_dir_outputs if not self.prepare_job( job_wrapper, include_metadata=True ): return # Get shell and job execution interface job_destination = job_wrapper.job_destination - shell_params, job_params = self.parse_destination_params(job_destination.params) - shell, job_interface = self.get_cli_plugins(shell_params, job_params) + + # Determine the job's Kubernetes destination (context, namespace) and options from the job destination definition + + k8s_job_obj = { + "apiVersion": "batch/v1", + "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": self.__produce_unique_k8s_job_name(job_wrapper)} + , + "spec": self.__produce_k8s_job_spec(job_wrapper)} + # wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers. galaxy_id_tag = job_wrapper.get_id_tag() + k8s_job = Job(self._pykube_api, k8s_job_obj).create() + + # define job attributes ajs = AsynchronousJobState( files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper ) - job_file_kwargs = job_interface.job_script_kwargs(ajs.output_file, ajs.error_file, ajs.job_name) - script = self.get_job_file( - job_wrapper, - exit_code_path=ajs.exit_code_file, - **job_file_kwargs - ) + # external_runJob_script can be None, in which case it's not used. + external_runjob_script = job_wrapper.get_destination_configuration("drmaa_external_runjob_script", None) - try: - self.write_executable_script( ajs.job_file, script ) - except: - log.exception("(%s) failure writing job script" % galaxy_id_tag ) - job_wrapper.fail("failure preparing job script", exception=True) - return + def __produce_unique_k8s_job_name(self, job_wrapper): + return job_wrapper.get_id_tag() + "-" + - # job was deleted while we were preparing it - if job_wrapper.get_state() == model.Job.states.DELETED: - log.info("(%s) Job deleted by user before it entered the queue" % galaxy_id_tag ) - if job_wrapper.cleanup_job in ("always", "onsuccess"): - job_wrapper.cleanup() - return + 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.""" + k8s_job_spec = {"template": self.__get_k8s_job_spec_template(job_wrapper)} + return k8s_job_spec - log.debug( "(%s) submitting file: %s" % ( galaxy_id_tag, ajs.job_file ) ) + def __get_k8s_job_spec_template(self, job_wrapper): + """The spec template is nothing but a Pod spec, except that it is nested and does not have an apiversion + nor kind. In addition to required fields for a Pod, a pod template in a job must specify appropriate labels + (see pod selector) and an appropriate restart policy.""" + k8s_spec_template = { + "volumes": self.__get_k8s_mountable_volumes(self, job_wrapper), + "containers": self.__get_k8s_containers(self, job_wrapper), + "restartPolicy": self.__get_k8s_restart_policy(self, job_wrapper) + } + # TODO include other relevant elements that people might want to use from + # TODO http://kubernetes.io/docs/api-reference/v1/definitions/#_v1_podspec - cmd_out = shell.execute(job_interface.submit(ajs.job_file)) - if cmd_out.returncode != 0: - log.error('(%s) submission failed (stdout): %s' % (galaxy_id_tag, cmd_out.stdout)) - log.error('(%s) submission failed (stderr): %s' % (galaxy_id_tag, cmd_out.stderr)) - job_wrapper.fail("failure submitting job") - return - # Some job runners return something like 'Submitted batch job XXXX' - # Strip and split to get job ID. - external_job_id = cmd_out.stdout.strip().split()[-1] - if not external_job_id: - log.error('(%s) submission did not return a job identifier, failing job' % galaxy_id_tag) - job_wrapper.fail("failure submitting job") - return + return k8s_spec_template - log.info("(%s) queued with identifier: %s" % ( galaxy_id_tag, external_job_id ) ) + def __get_k8s_restart_policy(self, job_wrapper): + """The default Kubernetes restart policy for Jobs""" + return "Never" - # store runner information for tracking if Galaxy restarts - job_wrapper.set_job_destination( job_destination, external_job_id ) + def __get_k8s_mountable_volumes(self, job_wrapper): + """Provides the required volumes that the containers in the pod should be able to mount. This should be using + the new persistent volumes and persistent volumes claim objects. + """ + - # Store state information for job - ajs.job_id = external_job_id - ajs.old_state = 'new' - ajs.job_destination = job_destination + def __get_k8s_containers(self, job_wrapper): + """Fills in all required for setting up the docker containers to be used.""" + k8s_container = { + "name": self.__get_k8s_container_name(job_wrapper), + "image": self.__assemble_k8s_container_image_name(job_wrapper) + } + + if self.__requires_ports(job_wrapper): + k8s_container['ports'] = self.__get_k8s_containers_ports(job_wrapper) + + return k8s_container + + def __assemble_k8s_container_image_name(self, job_wrapper): + """Assembles the container image name as repo/owner/image:tag, where repo, owner and tag are optional""" + job_destination = job_wrapper.job_destination + + # Determine the job's Kubernetes destination (context, namespace) and options from the job destination definition + repo = "" + owner = "" + if 'repo' in job_destination.params: + repo = job_destination.params['repo']+"/" + if 'owner' in job_destination.params: + owner = job_destination.params['owner']+"/" + + k8s_cont_image = repo + owner + job_destination.params['image'] + + if 'tag' in job_destination.params: + k8s_cont_image += ":" + job_destination.params['tag'] + + return k8s_cont_image + + def __get_k8s_container_name(self, job_wrapper): + # TODO check if this is correct + return job_wrapper.job_destination.id - # Add to our 'queue' of jobs to monitor - self.monitor_queue.put( ajs ) def check_watched_items( self ): """ From 2c485d768057ca6e80eb1650fe93fad5b7e0e371 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Thu, 7 Apr 2016 14:08:34 +0100 Subject: [PATCH 09/56] Adds pieces for PersistentVolume and PersistentVolumesClaims --- config/job_conf_k8s.xml | 4 ++++ lib/galaxy/jobs/runners/kubernetes.py | 19 ++++++++++++++++--- 2 files changed, 20 insertions(+), 3 deletions(-) diff --git a/config/job_conf_k8s.xml b/config/job_conf_k8s.xml index 9cc7da41cb2..a3f44f16bd5 100644 --- a/config/job_conf_k8s.xml +++ b/config/job_conf_k8s.xml @@ -6,6 +6,10 @@ /path/to/kubeconfig galaxy_pvc + + /mnt/glusterfs diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 9cefe923729..dfa98a33710 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -45,6 +45,7 @@ class KubernetesJobRunner( AsynchronousJobRunner ): # here we need to fetch the default kubeconfig path from the plugin defined in job_conf... self._pykube_api = HTTPClient(KubeConfig.from_file(self.runner_params["k8s_config_path"])) + self._galaxy_vol_name = "pvc_galaxy" # TODO do we need these? # self._init_monitor_thread() @@ -131,15 +132,27 @@ class KubernetesJobRunner( AsynchronousJobRunner ): def __get_k8s_mountable_volumes(self, job_wrapper): """Provides the required volumes that the containers in the pod should be able to mount. This should be using - the new persistent volumes and persistent volumes claim objects. + the new persistent volumes and persistent volumes claim objects. This requires that both a PersistentVolume and + a PersistentVolumeClaim are created before starting galaxy (starting a k8s job). """ - + # TODO on this initial version we only support a single volume to be mounted. + k8s_mountable_volume = { + "name": self._galaxy_vol_name, + "persistentVolumeClaim": { + "claimName": self.runner_params['k8s_persistent_volume_claim_name'] + } + } + return k8s_mountable_volume def __get_k8s_containers(self, job_wrapper): """Fills in all required for setting up the docker containers to be used.""" k8s_container = { "name": self.__get_k8s_container_name(job_wrapper), - "image": self.__assemble_k8s_container_image_name(job_wrapper) + "image": self.__assemble_k8s_container_image_name(job_wrapper), + "volumeMounts": { + "mountPath": self.runner_params['k8s_persistent_volume_claim_mount_path'], + "name": self._galaxy_vol_name + } } if self.__requires_ports(job_wrapper): From 2607be1810a5b95cab863d715efe8e889918c2fd Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 13 Apr 2016 15:53:51 +0100 Subject: [PATCH 10/56] Adds namespace and port proposal for job_conf file. --- config/job_conf_k8s.xml | 2 ++ 1 file changed, 2 insertions(+) diff --git a/config/job_conf_k8s.xml b/config/job_conf_k8s.xml index a3f44f16bd5..0fdb42c53fa 100644 --- a/config/job_conf_k8s.xml +++ b/config/job_conf_k8s.xml @@ -10,6 +10,7 @@ set in universe_wsgi.ini (or equivalent general galaxy config file). --> /mnt/glusterfs + galaxy-instanceA @@ -22,6 +23,7 @@ fatcat blankfilter latest + 8080 From ac2a82c47cffceb4eeb18d201f4f48887828ca55 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 13 Apr 2016 15:56:56 +0100 Subject: [PATCH 11/56] Minor formatting --- lib/galaxy/jobs/runners/kubernetes.py | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index dfa98a33710..242b6b619dc 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -39,7 +39,7 @@ class KubernetesJobRunner( AsynchronousJobRunner ): # Check if pykube was importable, fail if not assert operator is not None, K8S_IMPORT_MESSAGE """Start the job runner parent object """ - super( KubernetesJobRunner, self ).__init__( app, nworkers ) + super(KubernetesJobRunner, self).__init__(app, nworkers) # self.cli_interface = CliInterface() @@ -51,12 +51,13 @@ class KubernetesJobRunner( AsynchronousJobRunner ): # self._init_monitor_thread() # self._init_worker_threads() - def get_cli_plugins( self, shell_params, job_params ): + def get_cli_plugins(self, shell_params, job_params ): return self.cli_interface.get_plugins( shell_params, job_params ) def url_to_destination( self, url ): + # TODO apparently needs to be implemented for pykube-k8s params = {} - shell_params, job_params = url.split( '/' )[ 2:4 ] + shell_params, job_params = url.split('/')[2:4] # split 'foo=bar&baz=quux' into { 'foo' : 'bar', 'baz' : 'quux' } shell_params = dict( [ ( 'shell_' + k, v ) for k, v in [ kv.split( '=', 1 ) for kv in shell_params.split( '&' ) ] ] ) job_params = dict( [ ( 'job_' + k, v ) for k, v in [ kv.split( '=', 1 ) for kv in job_params.split( '&' ) ] ] ) @@ -67,13 +68,14 @@ class KubernetesJobRunner( AsynchronousJobRunner ): return JobDestination( runner='cli', params=params ) def parse_destination_params( self, params ): + # TODO apparently no need to re-implement, can be deleted. return split_params( params ) def queue_job( self, job_wrapper ): """Create job script and submit it to Kubernetes cluster""" # prepare the job - # TODO understand weather we need include_metadata and include_work_dir_outputs - if not self.prepare_job( job_wrapper, include_metadata=True ): + # TODO understand whether we need to include_metadata and include_work_dir_outputs + if not self.prepare_job(job_wrapper, include_metadata=True ): return # Get shell and job execution interface From af03b96ac13038aab81f2301f213a2a710c7520d Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 13 Apr 2016 15:57:36 +0100 Subject: [PATCH 12/56] Job name and correct spec. --- lib/galaxy/jobs/runners/kubernetes.py | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 242b6b619dc..835ef122ebe 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -83,15 +83,18 @@ class KubernetesJobRunner( AsynchronousJobRunner ): # Determine the job's Kubernetes destination (context, namespace) and options from the job destination definition + # 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_obj = { "apiVersion": "batch/v1", "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": self.__produce_unique_k8s_job_name(job_wrapper)} + {"name": k8s_job_name} , - "spec": self.__produce_k8s_job_spec(job_wrapper)} + "spec": self.__get_k8s_job_spec(job_wrapper) + } # wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers. From c16fe1618fe535a92cb4ee97432e4675010741fa Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 13 Apr 2016 15:59:03 +0100 Subject: [PATCH 13/56] Adds async jobs to monitor queue --- lib/galaxy/jobs/runners/kubernetes.py | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 835ef122ebe..d3707949c81 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -103,11 +103,13 @@ class KubernetesJobRunner( AsynchronousJobRunner ): k8s_job = Job(self._pykube_api, k8s_job_obj).create() - # define job attributes - ajs = AsynchronousJobState( files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper ) + # define job attributes in the AsyncronousJobState for follow-up + ajs = AsynchronousJobState(files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper, + job_id=k8s_job_name, job_destination=job_destination) + self.monitor_queue.put(ajs) # external_runJob_script can be None, in which case it's not used. - external_runjob_script = job_wrapper.get_destination_configuration("drmaa_external_runjob_script", None) + external_runjob_script = None def __produce_unique_k8s_job_name(self, job_wrapper): return job_wrapper.get_id_tag() + "-" + From f6dc05f57cda5f510eaaec6ec978360b918fd73b Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 13 Apr 2016 16:00:05 +0100 Subject: [PATCH 14/56] Improves job name, prepending string to make it URLable as required by kubernetes --- 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 d3707949c81..73d1a69ce6d 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -112,7 +112,7 @@ class KubernetesJobRunner( AsynchronousJobRunner ): external_runjob_script = None def __produce_unique_k8s_job_name(self, job_wrapper): - return job_wrapper.get_id_tag() + "-" + + return "galaxy-" + job_wrapper.get_id_tag() 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.""" From 0d21b52d52469300917fe7ea83d6f1100987878e Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 13 Apr 2016 16:01:23 +0100 Subject: [PATCH 15/56] Comments ports usage (we leave this feature for later) and other minor formatting. --- lib/galaxy/jobs/runners/kubernetes.py | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 73d1a69ce6d..3877fe4bf04 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -162,16 +162,22 @@ class KubernetesJobRunner( AsynchronousJobRunner ): } } - if self.__requires_ports(job_wrapper): - k8s_container['ports'] = self.__get_k8s_containers_ports(job_wrapper) + #if self.__requires_ports(job_wrapper): + # k8s_container['ports'] = self.__get_k8s_containers_ports(job_wrapper) return k8s_container + #def __get_k8s_containers_ports(self, job_wrapper): + + # for k,v self.runner_params: + # if k.startswith("container_port_"): + def __assemble_k8s_container_image_name(self, job_wrapper): """Assembles the container image name as repo/owner/image:tag, where repo, owner and tag are optional""" job_destination = job_wrapper.job_destination - # Determine the job's Kubernetes destination (context, namespace) and options from the job destination definition + # Determine the job's Kubernetes destination (context, namespace) and options from the job destination + # definition repo = "" owner = "" if 'repo' in job_destination.params: From 6cbbbfe161003c4a595eba4c10dc358a920a26ab Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 13 Apr 2016 16:03:25 +0100 Subject: [PATCH 16/56] changes way that 'command' option is passed to k8s on the Job object to cater for multi line commands. --- lib/galaxy/jobs/runners/kubernetes.py | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 3877fe4bf04..a00c2a3942c 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -156,6 +156,11 @@ class KubernetesJobRunner( AsynchronousJobRunner ): k8s_container = { "name": self.__get_k8s_container_name(job_wrapper), "image": self.__assemble_k8s_container_image_name(job_wrapper), + # this form of command overrides the entrypoint and allows multi command + # command line execution, separated by ;, which is what Galaxy does + # to assemble the command. + # TODO possibly shell needs to be set by job_wrapper + "command": "[\"/bin/bash\",\"-c\",\""+job_wrapper.runner_command_line+"\"]", "volumeMounts": { "mountPath": self.runner_params['k8s_persistent_volume_claim_mount_path'], "name": self._galaxy_vol_name From 96d86de6cf519b39c3c85cfd9d05b52894897f78 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 13 Apr 2016 16:49:43 +0100 Subject: [PATCH 17/56] Implements stop job. --- lib/galaxy/jobs/runners/kubernetes.py | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index a00c2a3942c..4cb4f10f43f 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -265,12 +265,14 @@ class KubernetesJobRunner( AsynchronousJobRunner ): return job_states def stop_job( self, job ): - """Attempts to delete a dispatched job""" + """Attempts to delete a dispatched job to the k8s cluster""" try: - shell_params, job_params = self.parse_destination_params(job.destination_params) - shell, job_interface = self.get_cli_plugins(shell_params, job_params) - cmd_out = shell.execute(job_interface.delete( job.job_runner_external_id )) - assert cmd_out.returncode == 0, cmd_out.stderr + jobs = Job.objects(self._pykube_api).filter(selector="app="+job.job_runner_external_id) + if jobs.response['items'].len() >= 0: + job_to_delete = Job(self._pykube_api, jobs.response['items'][0]) + if job_to_delete.exists(): + job_to_delete.delete() + assert not job_to_delete.exists(), "Could not delete job,"+job.job_runner_external_id+" it still exists" log.debug( "(%s/%s) Terminated at user's request" % ( job.id, job.job_runner_external_id ) ) except Exception as e: log.debug( "(%s/%s) User killed running job, but error encountered during termination: %s" % ( job.id, job.job_runner_external_id, e ) ) From 820d33c7645c546cd341054fc18981fcf80a59c6 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 15 Apr 2016 17:23:15 +0100 Subject: [PATCH 18/56] Minor correction of message. --- lib/galaxy/jobs/runners/kubernetes.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 4cb4f10f43f..6eb4147820b 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -16,9 +16,10 @@ try: from pykube.config import KubeConfig from pykube.http import HTTPClient from pykube.objects import Job + from pykube.objects import Pod except ImportError as exc: operator = None - K8S_IMPORT_MESSAGE = ('The Python pbs-python package is required to use ' + 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)) From 820913c2d641c38d279903e07d2008cd864390bc Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 15 Apr 2016 17:24:38 +0100 Subject: [PATCH 19/56] Adds some todo's --- lib/galaxy/jobs/runners/kubernetes.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 6eb4147820b..9df840ad3db 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -267,6 +267,7 @@ class KubernetesJobRunner( AsynchronousJobRunner ): def stop_job( self, job ): """Attempts to delete a dispatched job to the k8s cluster""" + # TODO rescale instead of delete the job, as in kubectl scale --replicas=0 jobs/myjob try: jobs = Job.objects(self._pykube_api).filter(selector="app="+job.job_runner_external_id) if jobs.response['items'].len() >= 0: @@ -280,6 +281,7 @@ class KubernetesJobRunner( AsynchronousJobRunner ): def recover( self, job, job_wrapper ): """Recovers jobs stuck in the queued/running state when Galaxy started""" + # TODO this needs to be implemented to override unimplemented base method job_id = job.get_job_runner_external_id() if job_id is None: self.put( job_wrapper ) From ef70f0edfad324a0887d53084d38414d110ac1b4 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Mon, 18 Apr 2016 14:27:27 +0100 Subject: [PATCH 20/56] Implements watch_item and replaces the watched_items, to use the parent logic. --- lib/galaxy/jobs/runners/kubernetes.py | 72 ++++++++++++--------------- 1 file changed, 33 insertions(+), 39 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 9df840ad3db..fa953645d9b 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -203,47 +203,41 @@ class KubernetesJobRunner( AsynchronousJobRunner ): return job_wrapper.job_destination.id - def check_watched_items( self ): - """ - Called by the monitor thread to look at each watched job and deal - with state changes. - """ - new_watched = [] + def check_watched_item(self, job_state): + """Checks the state of a job already submitted on k8s""" + 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]) + succeeded = 0 + active = 0 + failed = 0 + if 'succeeded' in job.obj['status']: + succeeded = job.obj['status']['succeeded'] + if 'active' in job.obj['status']: + active = job.obj['status']['active'] + if 'failed' in job.obj['status']: + failed = job.obj['status']['failed'] - job_states = self.__get_job_states() + # This assumes jobs dependent on a single pod. + if succeeded > 0: + self.mark_as_finished(job_state) + elif active > 0 or succeeded + active + failed == 0: + self.mark_as_queued(job_state) + elif failed > job_state.job_destination.params['max_pod_retrials']: + self.mark_as_failed(job_state) + job.scale(replicas=0) + + elif len(jobs.response['items']) == 0: + # there is no job responding to this job_id, it is either lost or something happened. + self.mark_as_failed(job_state) + return job_state + else: + # TODO: possibly some warning or message should be provided here to stderr + # TODO: of the job + # there is more than one job associated to the expected unique job id used as selector. + self.mark_as_failed(job_state) + return job_state - for ajs in self.watched: - external_job_id = ajs.job_id - id_tag = ajs.job_wrapper.get_id_tag() - old_state = ajs.old_state - state = job_states.get(external_job_id, None) - if state is None: - if ajs.job_wrapper.get_state() == model.Job.states.DELETED: - continue - log.debug("(%s/%s) job not found in batch state check" % ( id_tag, external_job_id ) ) - shell_params, job_params = self.parse_destination_params(ajs.job_destination.params) - shell, job_interface = self.get_cli_plugins(shell_params, job_params) - cmd_out = shell.execute(job_interface.get_single_status(external_job_id)) - state = job_interface.parse_single_status(cmd_out.stdout, external_job_id) - if state == model.Job.states.OK: - log.debug('(%s/%s) job execution finished, running job wrapper finish method' % ( id_tag, external_job_id ) ) - self.work_queue.put( ( self.finish_job, ajs ) ) - continue - else: - log.warning('(%s/%s) job not found in batch state check, but found in individual state check' % ( id_tag, external_job_id ) ) - if state != old_state: - ajs.job_wrapper.change_state( state ) - else: - if state != old_state: - log.debug("(%s/%s) state change: %s" % ( id_tag, external_job_id, state ) ) - ajs.job_wrapper.change_state( state ) - if state == model.Job.states.RUNNING and not ajs.running: - ajs.running = True - ajs.job_wrapper.change_state( model.Job.states.RUNNING ) - ajs.old_state = state - new_watched.append( ajs ) - # Replace the watch list with the updated version - self.watched = new_watched def __get_job_states(self): job_destinations = {} From f80242069b4e4eed455e5ac2a4f37f3e2404c6a0 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Mon, 18 Apr 2016 14:33:49 +0100 Subject: [PATCH 21/56] Changes stop job to scale job to zero replicas instead of deleting the job (for logging purposes) --- lib/galaxy/jobs/runners/kubernetes.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index fa953645d9b..5333428169d 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -261,14 +261,13 @@ class KubernetesJobRunner( AsynchronousJobRunner ): def stop_job( self, job ): """Attempts to delete a dispatched job to the k8s cluster""" - # TODO rescale instead of delete the job, as in kubectl scale --replicas=0 jobs/myjob try: jobs = Job.objects(self._pykube_api).filter(selector="app="+job.job_runner_external_id) if jobs.response['items'].len() >= 0: job_to_delete = Job(self._pykube_api, jobs.response['items'][0]) - if job_to_delete.exists(): - job_to_delete.delete() - assert not job_to_delete.exists(), "Could not delete job,"+job.job_runner_external_id+" it still exists" + job_to_delete.scale(replicas=0) + # TODO assert whether job parallelism == 0 + # assert not job_to_delete.exists(), "Could not delete job,"+job.job_runner_external_id+" it still exists" log.debug( "(%s/%s) Terminated at user's request" % ( job.id, job.job_runner_external_id ) ) except Exception as e: log.debug( "(%s/%s) User killed running job, but error encountered during termination: %s" % ( job.id, job.job_runner_external_id, e ) ) From 000259512057592909b459ff609742d1ae26a35e Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Mon, 18 Apr 2016 14:34:46 +0100 Subject: [PATCH 22/56] Minor documentation improvement. --- lib/galaxy/jobs/runners/kubernetes.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 5333428169d..4af89167691 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -121,7 +121,7 @@ class KubernetesJobRunner( AsynchronousJobRunner ): return k8s_job_spec def __get_k8s_job_spec_template(self, job_wrapper): - """The spec template is nothing but a Pod spec, except that it is nested and does not have an apiversion + """The k8s spec template is nothing but a Pod spec, except that it is nested and does not have an apiversion nor kind. In addition to required fields for a Pod, a pod template in a job must specify appropriate labels (see pod selector) and an appropriate restart policy.""" k8s_spec_template = { @@ -240,6 +240,8 @@ class KubernetesJobRunner( AsynchronousJobRunner ): def __get_job_states(self): + """Get the states of all jobs submitted by this Galaxy runner + to the Kubernetes cluster""" job_destinations = {} job_states = {} # unique the list of destinations From 25f422290e2b2a5a48b84b872efa3d1f15d50a66 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Mon, 18 Apr 2016 14:35:43 +0100 Subject: [PATCH 23/56] Adds pods retrials for per installation and individual docker destination. --- config/job_conf_k8s.xml | 3 +++ 1 file changed, 3 insertions(+) diff --git a/config/job_conf_k8s.xml b/config/job_conf_k8s.xml index 0fdb42c53fa..4e3c861b6c4 100644 --- a/config/job_conf_k8s.xml +++ b/config/job_conf_k8s.xml @@ -11,6 +11,8 @@ --> /mnt/glusterfs galaxy-instanceA + + 4 @@ -24,6 +26,7 @@ blankfilter latest 8080 + 3 From bdf283edb0a4ef49349f58a6b2df4014d0e3ae00 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Mon, 18 Apr 2016 18:03:11 +0100 Subject: [PATCH 24/56] Removes unnecessary __get_job_states --- lib/galaxy/jobs/runners/kubernetes.py | 23 ----------------------- 1 file changed, 23 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 4af89167691..481bde1386c 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -238,29 +238,6 @@ class KubernetesJobRunner( AsynchronousJobRunner ): self.mark_as_failed(job_state) return job_state - - def __get_job_states(self): - """Get the states of all jobs submitted by this Galaxy runner - to the Kubernetes cluster""" - job_destinations = {} - job_states = {} - # unique the list of destinations - for ajs in self.watched: - if ajs.job_destination.id not in job_destinations: - job_destinations[ajs.job_destination.id] = dict( job_destination=ajs.job_destination, job_ids=[ ajs.job_id ] ) - else: - job_destinations[ajs.job_destination.id]['job_ids'].append( ajs.job_id ) - # check each destination for the listed job ids - for job_destination_id, v in job_destinations.items(): - job_destination = v['job_destination'] - job_ids = v['job_ids'] - shell_params, job_params = self.parse_destination_params(job_destination.params) - shell, job_interface = self.get_cli_plugins(shell_params, job_params) - cmd_out = shell.execute(job_interface.get_status(job_ids)) - assert cmd_out.returncode == 0, cmd_out.stderr - job_states.update(job_interface.parse_status(cmd_out.stdout, job_ids)) - return job_states - def stop_job( self, job ): """Attempts to delete a dispatched job to the k8s cluster""" try: From 8ff8c1c4f591d2ae1d61f8ce76b4758f47cff305 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Mon, 18 Apr 2016 18:04:31 +0100 Subject: [PATCH 25/56] Uses k8s pod log as stdout/stderr. --- lib/galaxy/jobs/runners/kubernetes.py | 22 +++++++++++++++++++++- 1 file changed, 21 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 481bde1386c..6e9092460b2 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -8,6 +8,7 @@ from galaxy import model from galaxy.jobs import JobDestination from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner from .util.cli import CliInterface, split_params +from os import sep as os_sep # pykube imports: try: @@ -218,9 +219,13 @@ class KubernetesJobRunner( AsynchronousJobRunner ): if 'failed' in job.obj['status']: failed = job.obj['status']['failed'] - # This assumes jobs dependent on a single pod. + # This assumes jobs dependent on a single pod, single container if succeeded > 0: + logs_file_path = self.__produce_log_file(job_state) + job_state.output_file = logs_file_path + self.mark_as_finished(job_state) + elif active > 0 or succeeded + active + failed == 0: self.mark_as_queued(job_state) elif failed > job_state.job_destination.params['max_pod_retrials']: @@ -235,9 +240,24 @@ class KubernetesJobRunner( AsynchronousJobRunner ): # TODO: possibly some warning or message should be provided here to stderr # TODO: of the job # there is more than one job associated to the expected unique job id used as selector. + job_state.error_file = self.__produce_log_file(job_state) self.mark_as_failed(job_state) return job_state + 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.get_logs(timestamps=True) + logs += "\n\n==== Pod " + pod.name + " log end ====" + logs_file_path = job_state.files_dir + os_sep + pod.name + '.log' + logs_file = open(logs_file_path) + logs_file.write(logs) + logs_file.close() + return logs_file_path + def stop_job( self, job ): """Attempts to delete a dispatched job to the k8s cluster""" try: From 9fbf03dd7d45e3c06b66b98a562063284e28b8ed Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 19 Apr 2016 13:20:59 +0100 Subject: [PATCH 26/56] Removes unnecessary get CLI plugins method. --- lib/galaxy/jobs/runners/kubernetes.py | 3 --- 1 file changed, 3 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 6e9092460b2..150aed8a1c9 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -53,9 +53,6 @@ class KubernetesJobRunner( AsynchronousJobRunner ): # self._init_monitor_thread() # self._init_worker_threads() - def get_cli_plugins(self, shell_params, job_params ): - return self.cli_interface.get_plugins( shell_params, job_params ) - def url_to_destination( self, url ): # TODO apparently needs to be implemented for pykube-k8s params = {} From f686253ca7715e509c74778f3e78ab96358ccc92 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 19 Apr 2016 13:21:40 +0100 Subject: [PATCH 27/56] PEP8 compliant formatting --- lib/galaxy/jobs/runners/kubernetes.py | 84 +++++++++++++-------------- 1 file changed, 41 insertions(+), 43 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 150aed8a1c9..160af5e578a 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -24,20 +24,18 @@ except ImportError as exc: 'this feature, please install it or correct the ' 'following error:\nImportError %s' % str(exc)) +log = logging.getLogger(__name__) + +__all__ = ['KubernetesJobRunner'] -log = logging.getLogger( __name__ ) - -__all__ = [ 'KubernetesJobRunner' ] - - -class KubernetesJobRunner( AsynchronousJobRunner ): +class KubernetesJobRunner(AsynchronousJobRunner): """ Job runner backed by a finite pool of worker threads. FIFO scheduling """ runner_name = "KubernetesRunner" - def __init__( self, app, nworkers ): + def __init__(self, app, nworkers): # Check if pykube was importable, fail if not assert operator is not None, K8S_IMPORT_MESSAGE """Start the job runner parent object """ @@ -53,28 +51,28 @@ class KubernetesJobRunner( AsynchronousJobRunner ): # self._init_monitor_thread() # self._init_worker_threads() - def url_to_destination( self, url ): + def url_to_destination(self, url): # TODO apparently needs to be implemented for pykube-k8s params = {} shell_params, job_params = url.split('/')[2:4] # split 'foo=bar&baz=quux' into { 'foo' : 'bar', 'baz' : 'quux' } - shell_params = dict( [ ( 'shell_' + k, v ) for k, v in [ kv.split( '=', 1 ) for kv in shell_params.split( '&' ) ] ] ) - job_params = dict( [ ( 'job_' + k, v ) for k, v in [ kv.split( '=', 1 ) for kv in job_params.split( '&' ) ] ] ) - params.update( shell_params ) - params.update( job_params ) - log.debug( "Converted URL '%s' to destination runner=cli, params=%s" % ( url, params ) ) + shell_params = dict([('shell_' + k, v) for k, v in [kv.split('=', 1) for kv in shell_params.split('&')]]) + job_params = dict([('job_' + k, v) for k, v in [kv.split('=', 1) for kv in job_params.split('&')]]) + params.update(shell_params) + params.update(job_params) + log.debug("Converted URL '%s' to destination runner=cli, params=%s" % (url, params)) # Create a dynamic JobDestination - return JobDestination( runner='cli', params=params ) + return JobDestination(runner='cli', params=params) - def parse_destination_params( self, params ): + def parse_destination_params(self, params): # TODO apparently no need to re-implement, can be deleted. - return split_params( params ) + return split_params(params) - def queue_job( self, job_wrapper ): + def queue_job(self, job_wrapper): """Create job script and submit it to Kubernetes cluster""" # prepare the job # TODO understand whether we need to include_metadata and include_work_dir_outputs - if not self.prepare_job(job_wrapper, include_metadata=True ): + if not self.prepare_job(job_wrapper, include_metadata=True): return # Get shell and job execution interface @@ -88,20 +86,18 @@ class KubernetesJobRunner( AsynchronousJobRunner ): "apiVersion": "batch/v1", "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} + # 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} , "spec": self.__get_k8s_job_spec(job_wrapper) } - # wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers. galaxy_id_tag = job_wrapper.get_id_tag() k8s_job = Job(self._pykube_api, k8s_job_obj).create() - # define job attributes in the AsyncronousJobState for follow-up ajs = AsynchronousJobState(files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper, job_id=k8s_job_name, job_destination=job_destination) @@ -126,7 +122,7 @@ class KubernetesJobRunner( AsynchronousJobRunner ): "volumes": self.__get_k8s_mountable_volumes(self, job_wrapper), "containers": self.__get_k8s_containers(self, job_wrapper), "restartPolicy": self.__get_k8s_restart_policy(self, job_wrapper) - } + } # TODO include other relevant elements that people might want to use from # TODO http://kubernetes.io/docs/api-reference/v1/definitions/#_v1_podspec @@ -159,19 +155,19 @@ class KubernetesJobRunner( AsynchronousJobRunner ): # command line execution, separated by ;, which is what Galaxy does # to assemble the command. # TODO possibly shell needs to be set by job_wrapper - "command": "[\"/bin/bash\",\"-c\",\""+job_wrapper.runner_command_line+"\"]", + "command": "[\"/bin/bash\",\"-c\",\"" + job_wrapper.runner_command_line + "\"]", "volumeMounts": { "mountPath": self.runner_params['k8s_persistent_volume_claim_mount_path'], "name": self._galaxy_vol_name } } - #if self.__requires_ports(job_wrapper): + # if self.__requires_ports(job_wrapper): # k8s_container['ports'] = self.__get_k8s_containers_ports(job_wrapper) return k8s_container - #def __get_k8s_containers_ports(self, job_wrapper): + # def __get_k8s_containers_ports(self, job_wrapper): # for k,v self.runner_params: # if k.startswith("container_port_"): @@ -185,9 +181,9 @@ class KubernetesJobRunner( AsynchronousJobRunner ): repo = "" owner = "" if 'repo' in job_destination.params: - repo = job_destination.params['repo']+"/" + repo = job_destination.params['repo'] + "/" if 'owner' in job_destination.params: - owner = job_destination.params['owner']+"/" + owner = job_destination.params['owner'] + "/" k8s_cont_image = repo + owner + job_destination.params['image'] @@ -200,10 +196,9 @@ class KubernetesJobRunner( AsynchronousJobRunner ): # TODO check if this is correct return job_wrapper.job_destination.id - def check_watched_item(self, job_state): """Checks the state of a job already submitted on k8s""" - jobs = Job.objects(self._pykube_api).filter(selector="app="+job_state.job_id) + 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]) succeeded = 0 @@ -255,38 +250,41 @@ class KubernetesJobRunner( AsynchronousJobRunner ): logs_file.close() return logs_file_path - def stop_job( self, job ): + 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) + jobs = Job.objects(self._pykube_api).filter(selector="app=" + job.job_runner_external_id) if jobs.response['items'].len() >= 0: job_to_delete = Job(self._pykube_api, jobs.response['items'][0]) job_to_delete.scale(replicas=0) # TODO assert whether job parallelism == 0 # assert not job_to_delete.exists(), "Could not delete job,"+job.job_runner_external_id+" it still exists" - log.debug( "(%s/%s) Terminated at user's request" % ( job.id, job.job_runner_external_id ) ) + log.debug("(%s/%s) Terminated at user's request" % (job.id, job.job_runner_external_id)) except Exception as e: - log.debug( "(%s/%s) User killed running job, but error encountered during termination: %s" % ( job.id, job.job_runner_external_id, e ) ) + log.debug("(%s/%s) User killed running job, but error encountered during termination: %s" % ( + job.id, job.job_runner_external_id, e)) - def recover( self, job, job_wrapper ): + def recover(self, job, job_wrapper): """Recovers jobs stuck in the queued/running state when Galaxy started""" # TODO this needs to be implemented to override unimplemented base method job_id = job.get_job_runner_external_id() if job_id is None: - self.put( job_wrapper ) + self.put(job_wrapper) return - ajs = AsynchronousJobState( files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper ) - ajs.job_id = str( job_id ) + ajs = AsynchronousJobState(files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper) + ajs.job_id = str(job_id) ajs.command_line = job.command_line ajs.job_wrapper = job_wrapper ajs.job_destination = job_wrapper.job_destination if job.state == model.Job.states.RUNNING: - log.debug( "(%s/%s) is still in running state, adding to the runner monitor queue" % ( job.id, job.job_runner_external_id ) ) + log.debug("(%s/%s) is still in running state, adding to the runner monitor queue" % ( + job.id, job.job_runner_external_id)) ajs.old_state = model.Job.states.RUNNING ajs.running = True - self.monitor_queue.put( ajs ) + self.monitor_queue.put(ajs) elif job.state == model.Job.states.QUEUED: - log.debug( "(%s/%s) is still in queued state, adding to the runner monitor queue" % ( job.id, job.job_runner_external_id ) ) + log.debug("(%s/%s) is still in queued state, adding to the runner monitor queue" % ( + job.id, job.job_runner_external_id)) ajs.old_state = model.Job.states.QUEUED ajs.running = False - self.monitor_queue.put( ajs ) + self.monitor_queue.put(ajs) From 58c969723264fb8fec5e837c9483c7ad401d3fae Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 19 Apr 2016 13:57:05 +0100 Subject: [PATCH 28/56] Clean up for queue method. --- lib/galaxy/jobs/runners/kubernetes.py | 11 ++++------- 1 file changed, 4 insertions(+), 7 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 160af5e578a..b3f10b94b55 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -71,15 +71,13 @@ class KubernetesJobRunner(AsynchronousJobRunner): def queue_job(self, job_wrapper): """Create job script and submit it to Kubernetes cluster""" # prepare the job - # TODO understand whether we need to include_metadata and include_work_dir_outputs - if not self.prepare_job(job_wrapper, include_metadata=True): + # We currently don't need to include_metadata or include_work_dir_outputs, as working directory is the same + # were galaxy will expect results. + if not self.prepare_job(job_wrapper, include_metadata=False, include_work_dir_outputs=False): return - # Get shell and job execution interface job_destination = job_wrapper.job_destination - # Determine the job's Kubernetes destination (context, namespace) and options from the job destination definition - # 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_obj = { @@ -93,8 +91,6 @@ class KubernetesJobRunner(AsynchronousJobRunner): "spec": self.__get_k8s_job_spec(job_wrapper) } - # wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers. - galaxy_id_tag = job_wrapper.get_id_tag() k8s_job = Job(self._pykube_api, k8s_job_obj).create() @@ -107,6 +103,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): external_runjob_script = None def __produce_unique_k8s_job_name(self, job_wrapper): + # wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers. return "galaxy-" + job_wrapper.get_id_tag() def __get_k8s_job_spec(self, job_wrapper): From e50ed67ced41554856498246f5d5fa3004e40b63 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 19 Apr 2016 13:59:39 +0100 Subject: [PATCH 29/56] Removes apparently unnecessary parse_destination_params... --- lib/galaxy/jobs/runners/kubernetes.py | 4 ---- 1 file changed, 4 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index b3f10b94b55..51770beab10 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -64,10 +64,6 @@ class KubernetesJobRunner(AsynchronousJobRunner): # Create a dynamic JobDestination return JobDestination(runner='cli', params=params) - def parse_destination_params(self, params): - # TODO apparently no need to re-implement, can be deleted. - return split_params(params) - def queue_job(self, job_wrapper): """Create job script and submit it to Kubernetes cluster""" # prepare the job From daf8ccbf0473872cb288267f23150cdcd2f80e84 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 19 Apr 2016 14:03:58 +0100 Subject: [PATCH 30/56] Removes url_to_destination, as this is only for legacy destinations (which I guess our's isn't, wouldn't have any URL to transform either) --- lib/galaxy/jobs/runners/kubernetes.py | 13 ------------- 1 file changed, 13 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 51770beab10..943b06a9b1b 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -51,19 +51,6 @@ class KubernetesJobRunner(AsynchronousJobRunner): # self._init_monitor_thread() # self._init_worker_threads() - def url_to_destination(self, url): - # TODO apparently needs to be implemented for pykube-k8s - params = {} - shell_params, job_params = url.split('/')[2:4] - # split 'foo=bar&baz=quux' into { 'foo' : 'bar', 'baz' : 'quux' } - shell_params = dict([('shell_' + k, v) for k, v in [kv.split('=', 1) for kv in shell_params.split('&')]]) - job_params = dict([('job_' + k, v) for k, v in [kv.split('=', 1) for kv in job_params.split('&')]]) - params.update(shell_params) - params.update(job_params) - log.debug("Converted URL '%s' to destination runner=cli, params=%s" % (url, params)) - # Create a dynamic JobDestination - return JobDestination(runner='cli', params=params) - def queue_job(self, job_wrapper): """Create job script and submit it to Kubernetes cluster""" # prepare the job From 258c2812b9c15be5101856ae0b4a33d51f087015 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 19 Apr 2016 15:39:02 +0100 Subject: [PATCH 31/56] Adds some logging plus job state changes. --- lib/galaxy/jobs/runners/kubernetes.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 943b06a9b1b..6446c9da52e 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -177,7 +177,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): return job_wrapper.job_destination.id def check_watched_item(self, job_state): - """Checks the state of a job already submitted on k8s""" + """Checks the state of a job already submitted on k8s. Job state is a AsynchronousJobState""" 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]) @@ -195,23 +195,24 @@ class KubernetesJobRunner(AsynchronousJobRunner): if succeeded > 0: logs_file_path = self.__produce_log_file(job_state) job_state.output_file = logs_file_path - + job_state.running = False self.mark_as_finished(job_state) elif active > 0 or succeeded + active + failed == 0: self.mark_as_queued(job_state) + job_state.running = True elif failed > job_state.job_destination.params['max_pod_retrials']: self.mark_as_failed(job_state) job.scale(replicas=0) elif len(jobs.response['items']) == 0: # there is no job responding to this job_id, it is either lost or something happened. + log.error("No Jobs are available under expected selector app="+job_state.job_id) self.mark_as_failed(job_state) return job_state else: - # TODO: possibly some warning or message should be provided here to stderr - # TODO: of the job # there is more than one job associated to the expected unique job id used as selector. + log.error("There is more than one Kubernetes Job associated to job id "+job_state.job_id) job_state.error_file = self.__produce_log_file(job_state) self.mark_as_failed(job_state) return job_state From 0c08228349a88f92d044bc903056cf3ad7b4e0cc Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 20 Apr 2016 17:48:03 +0100 Subject: [PATCH 32/56] Fixes initial handling of runner params. --- lib/galaxy/jobs/runners/kubernetes.py | 17 ++++++++++++++--- 1 file changed, 14 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 6446c9da52e..26b21548ece 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -8,7 +8,7 @@ from galaxy import model from galaxy.jobs import JobDestination from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner from .util.cli import CliInterface, split_params -from os import sep as os_sep +from os import sep as os_sep, environ as os_environ # pykube imports: try: @@ -35,11 +35,22 @@ class KubernetesJobRunner(AsynchronousJobRunner): """ runner_name = "KubernetesRunner" - def __init__(self, app, nworkers): + def __init__(self, app, nworkers, **kwargs): # Check if pykube was importable, fail if not assert operator is not None, K8S_IMPORT_MESSAGE + runner_param_specs = dict( + k8s_config_path=dict(map=str, default=os_environ.get('KUBECONFIG', None)), + k8s_persistent_volume_claim_name=dict(map=str), + k8s_persistent_volume_claim_mount_path=dict(map=str), + k8s_namespace=dict(map=str, default="default"), + k8s_pod_retrials=dict(map=int, valid=lambda x: int > 0, default=3)) + + if 'runner_param_specs' not in kwargs: + kwargs['runner_param_specs'] = dict() + kwargs['runner_param_specs'].update(runner_param_specs) + """Start the job runner parent object """ - super(KubernetesJobRunner, self).__init__(app, nworkers) + super(KubernetesJobRunner, self).__init__(app, nworkers, **kwargs) # self.cli_interface = CliInterface() From 43ff13d005aa8c7192837df2e33c559fbdf8d42d Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 20 Apr 2016 17:49:15 +0100 Subject: [PATCH 33/56] Starts monitor thread and worker threads. --- lib/galaxy/jobs/runners/kubernetes.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 26b21548ece..bef96c6ab59 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -58,9 +58,9 @@ class KubernetesJobRunner(AsynchronousJobRunner): self._pykube_api = HTTPClient(KubeConfig.from_file(self.runner_params["k8s_config_path"])) self._galaxy_vol_name = "pvc_galaxy" - # TODO do we need these? - # self._init_monitor_thread() - # self._init_worker_threads() + + self._init_monitor_thread() + self._init_worker_threads() def queue_job(self, job_wrapper): """Create job script and submit it to Kubernetes cluster""" From e3d239a484c4e04395a29dc8c030c30856a1f942 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 22 Apr 2016 12:12:28 +0100 Subject: [PATCH 34/56] Changes k8s job api path to extensions/v1beta1 as pykube doesn't support the other one yet. --- 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 bef96c6ab59..79d1ae5196a 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -75,7 +75,7 @@ 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) k8s_job_obj = { - "apiVersion": "batch/v1", + "apiVersion": "extensions/v1beta1", "kind": "Job", "metadata": # metadata.name is the name of the pod resource created, and must be unique From ffdc08bbc1cef98b3e10130b5763b8bd794ca9c8 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 22 Apr 2016 12:14:40 +0100 Subject: [PATCH 35/56] Fixes labels and faulty spec hierarchy. --- lib/galaxy/jobs/runners/kubernetes.py | 17 +++++++++++++---- 1 file changed, 13 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 79d1ae5196a..35ccfefb347 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -80,7 +80,11 @@ class KubernetesJobRunner(AsynchronousJobRunner): "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} + { + "name": k8s_job_name, + "namespace": "default", # TODO this should be set + "labels": {"app": k8s_job_name}, + } , "spec": self.__get_k8s_job_spec(job_wrapper) } @@ -110,9 +114,14 @@ class KubernetesJobRunner(AsynchronousJobRunner): nor kind. In addition to required fields for a Pod, a pod template in a job must specify appropriate labels (see pod selector) and an appropriate restart policy.""" k8s_spec_template = { - "volumes": self.__get_k8s_mountable_volumes(self, job_wrapper), - "containers": self.__get_k8s_containers(self, job_wrapper), - "restartPolicy": self.__get_k8s_restart_policy(self, job_wrapper) + "metadata" : { + "labels": { "app": self.__produce_unique_k8s_job_name(job_wrapper) } + }, + "spec" : { + "volumes": self.__get_k8s_mountable_volumes(job_wrapper), + "restartPolicy": self.__get_k8s_restart_policy(job_wrapper), + "containers": self.__get_k8s_containers(job_wrapper) + } } # TODO include other relevant elements that people might want to use from # TODO http://kubernetes.io/docs/api-reference/v1/definitions/#_v1_podspec From cabfd83ad2a747382d613585c0e040e4506b217b Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 22 Apr 2016 12:16:08 +0100 Subject: [PATCH 36/56] Fixes volumes and containers to be as expected json arrays. --- lib/galaxy/jobs/runners/kubernetes.py | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 35ccfefb347..f9db863925a 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -144,7 +144,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): "claimName": self.runner_params['k8s_persistent_volume_claim_name'] } } - return k8s_mountable_volume + return [ k8s_mountable_volume ] def __get_k8s_containers(self, job_wrapper): """Fills in all required for setting up the docker containers to be used.""" @@ -155,17 +155,17 @@ class KubernetesJobRunner(AsynchronousJobRunner): # command line execution, separated by ;, which is what Galaxy does # to assemble the command. # TODO possibly shell needs to be set by job_wrapper - "command": "[\"/bin/bash\",\"-c\",\"" + job_wrapper.runner_command_line + "\"]", - "volumeMounts": { + "command": ["/bin/bash", "-c", job_wrapper.runner_command_line], + "volumeMounts": [{ "mountPath": self.runner_params['k8s_persistent_volume_claim_mount_path'], "name": self._galaxy_vol_name - } + }] } # if self.__requires_ports(job_wrapper): # k8s_container['ports'] = self.__get_k8s_containers_ports(job_wrapper) - return k8s_container + return [ k8s_container ] # def __get_k8s_containers_ports(self, job_wrapper): From ec97296ba8990bda593bd78bbd9e79f02e62d31b Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 22 Apr 2016 12:24:00 +0100 Subject: [PATCH 37/56] Fixes non-DNS persistent volume claim. --- 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 f9db863925a..daab7ed41d6 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -56,7 +56,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): # here we need to fetch the default kubeconfig path from the plugin defined in job_conf... self._pykube_api = HTTPClient(KubeConfig.from_file(self.runner_params["k8s_config_path"])) - self._galaxy_vol_name = "pvc_galaxy" + self._galaxy_vol_name = "pvc-galaxy" self._init_monitor_thread() From 943fb128f1a8445ce8e79784b32ec71a625ec24a Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Mon, 25 Apr 2016 12:48:14 +0100 Subject: [PATCH 38/56] Improves formatting (PEP suggestions). --- lib/galaxy/jobs/runners/kubernetes.py | 32 +++++++++++++-------------- 1 file changed, 16 insertions(+), 16 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index daab7ed41d6..1a45cf08735 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -81,10 +81,10 @@ class KubernetesJobRunner(AsynchronousJobRunner): # 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": "default", # TODO this should be set - "labels": {"app": k8s_job_name}, - } + "name": k8s_job_name, + "namespace": "default", # TODO this should be set + "labels": {"app": k8s_job_name}, + } , "spec": self.__get_k8s_job_spec(job_wrapper) } @@ -114,14 +114,14 @@ class KubernetesJobRunner(AsynchronousJobRunner): nor kind. In addition to required fields for a Pod, a pod template in a job must specify appropriate labels (see pod selector) and an appropriate restart policy.""" k8s_spec_template = { - "metadata" : { - "labels": { "app": self.__produce_unique_k8s_job_name(job_wrapper) } + "metadata": { + "labels": {"app": self.__produce_unique_k8s_job_name(job_wrapper)} }, - "spec" : { - "volumes": self.__get_k8s_mountable_volumes(job_wrapper), - "restartPolicy": self.__get_k8s_restart_policy(job_wrapper), - "containers": self.__get_k8s_containers(job_wrapper) - } + "spec": { + "volumes": self.__get_k8s_mountable_volumes(job_wrapper), + "restartPolicy": self.__get_k8s_restart_policy(job_wrapper), + "containers": self.__get_k8s_containers(job_wrapper) + } } # TODO include other relevant elements that people might want to use from # TODO http://kubernetes.io/docs/api-reference/v1/definitions/#_v1_podspec @@ -144,7 +144,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): "claimName": self.runner_params['k8s_persistent_volume_claim_name'] } } - return [ k8s_mountable_volume ] + return [k8s_mountable_volume] def __get_k8s_containers(self, job_wrapper): """Fills in all required for setting up the docker containers to be used.""" @@ -165,7 +165,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): # if self.__requires_ports(job_wrapper): # k8s_container['ports'] = self.__get_k8s_containers_ports(job_wrapper) - return [ k8s_container ] + return [k8s_container] # def __get_k8s_containers_ports(self, job_wrapper): @@ -263,7 +263,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): log.debug("(%s/%s) Terminated at user's request" % (job.id, job.job_runner_external_id)) except Exception as e: log.debug("(%s/%s) User killed running job, but error encountered during termination: %s" % ( - job.id, job.job_runner_external_id, e)) + job.id, job.job_runner_external_id, e)) def recover(self, job, job_wrapper): """Recovers jobs stuck in the queued/running state when Galaxy started""" @@ -279,13 +279,13 @@ class KubernetesJobRunner(AsynchronousJobRunner): ajs.job_destination = job_wrapper.job_destination if job.state == model.Job.states.RUNNING: log.debug("(%s/%s) is still in running state, adding to the runner monitor queue" % ( - job.id, job.job_runner_external_id)) + job.id, job.job_runner_external_id)) ajs.old_state = model.Job.states.RUNNING ajs.running = True self.monitor_queue.put(ajs) elif job.state == model.Job.states.QUEUED: log.debug("(%s/%s) is still in queued state, adding to the runner monitor queue" % ( - job.id, job.job_runner_external_id)) + job.id, job.job_runner_external_id)) ajs.old_state = model.Job.states.QUEUED ajs.running = False self.monitor_queue.put(ajs) From a55afa51aa41198b027dfe757b9037890da97e9a Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Mon, 25 Apr 2016 12:49:35 +0100 Subject: [PATCH 39/56] Handles correctly error and output files. --- lib/galaxy/jobs/runners/kubernetes.py | 36 +++++++++++++++++++++------ 1 file changed, 28 insertions(+), 8 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 1a45cf08735..91d891998f8 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -213,27 +213,47 @@ class KubernetesJobRunner(AsynchronousJobRunner): # This assumes jobs dependent on a single pod, single container if succeeded > 0: - logs_file_path = self.__produce_log_file(job_state) - job_state.output_file = logs_file_path + self.__produce_log_file(job_state) + error_file = open(job_state.error_file, 'w') + error_file.write("") + error_file.close() job_state.running = False self.mark_as_finished(job_state) + return None elif active > 0 or succeeded + active + failed == 0: - self.mark_as_queued(job_state) job_state.running = True + return job_state elif failed > job_state.job_destination.params['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 self.mark_as_failed(job_state) job.scale(replicas=0) + return None + + # We should not get here + log.debug( + "Reaching unexpected point for Kubernetes job, where it is not classified as succ., active nor failed.") + return job_state elif len(jobs.response['items']) == 0: # there is no job responding to this job_id, it is either lost or something happened. - log.error("No Jobs are available under expected selector app="+job_state.job_id) + log.error("No Jobs are available under expected selector app=" + job_state.job_id) + error_file = open(job_state.error_file, 'w') + error_file.write("No Kubernetes Jobs are available under expected selector app=" + job_state.job_id + "\n") + error_file.close() self.mark_as_failed(job_state) return job_state else: # there is more than one job associated to the expected unique job id used as selector. - log.error("There is more than one Kubernetes Job associated to job id "+job_state.job_id) - job_state.error_file = self.__produce_log_file(job_state) + log.error("There is more than one Kubernetes Job associated to job id " + job_state.job_id) + self.__produce_log_file(job_state) + error_file = open(job_state.error_file, 'w') + error_file.write("There is more than one Kubernetes Job associated to job id " + job_state.job_id + "\n") + error_file.close() self.mark_as_failed(job_state) return job_state @@ -245,8 +265,8 @@ class KubernetesJobRunner(AsynchronousJobRunner): logs += "\n\n==== Pod " + pod.name + " log start ====\n\n" logs += pod.get_logs(timestamps=True) logs += "\n\n==== Pod " + pod.name + " log end ====" - logs_file_path = job_state.files_dir + os_sep + pod.name + '.log' - logs_file = open(logs_file_path) + logs_file_path = job_state.output_file + logs_file = open(logs_file_path, mode="w") logs_file.write(logs) logs_file.close() return logs_file_path From 74a4772e77dadddd4da7bbbce8164e5b5bbdd41b Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Mon, 25 Apr 2016 13:44:41 +0100 Subject: [PATCH 40/56] Removes some unneeded imports. --- lib/galaxy/jobs/runners/kubernetes.py | 3 --- 1 file changed, 3 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 91d891998f8..5da363a7d88 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -5,10 +5,7 @@ Job control via a command line interface (e.g. qsub/qstat), possibly over a remo import logging from galaxy import model -from galaxy.jobs import JobDestination from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner -from .util.cli import CliInterface, split_params -from os import sep as os_sep, environ as os_environ # pykube imports: try: From 94f38e79572ee44a2b5c7c4108d18a0aa2716bd0 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Tue, 3 May 2016 15:25:33 +0100 Subject: [PATCH 41/56] Missing os environment functionality import (forgot to include in commit). --- lib/galaxy/jobs/runners/kubernetes.py | 1 + 1 file changed, 1 insertion(+) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 5da363a7d88..61e2cf08da3 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -6,6 +6,7 @@ import logging from galaxy import model from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner +from os import environ as os_environ # pykube imports: try: From 3347713198940111c553676e3982bb96338707b5 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Wed, 4 May 2016 17:29:17 +0100 Subject: [PATCH 42/56] Adds job_conf documentation for Kubernetes runner on the job_conf.xml.sample_advanced file. --- config/job_conf.xml.sample_advanced | 107 ++++++++++++++++++++++++++++ 1 file changed, 107 insertions(+) diff --git a/config/job_conf.xml.sample_advanced b/config/job_conf.xml.sample_advanced index 92dbe4e3410..96f3b465e68 100644 --- a/config/job_conf.xml.sample_advanced +++ b/config/job_conf.xml.sample_advanced @@ -117,6 +117,83 @@ --> + + + + /path/to/kubeconfig + <--- This is the path to the kube config file, which is normally on ~/.kube/config, but that will depend on + your installation. This is the file that tells the plugin where the k8s cluster is, access credentials, + etc. --> + + galaxy_pvc + + + /scratch1/galaxy_data + + + galaxy-instanceA + + + 4 + + @@ -475,6 +552,31 @@ false + + + + + my-docker-registry.org + superbioinfo + my-tool + latest + 3 + + + + @@ -506,6 +608,11 @@ and pass to dynamic destination (as resource_params argument). --> + + + - - - - - /path/to/kubeconfig - galaxy_pvc - - /mnt/glusterfs - galaxy-instanceA - - 4 - - - - - - - - - docker-registry.local:6000 - fatcat - blankfilter - latest - 8080 - 3 - - - - - - trackster - - - - - - - - From b2cf8b5a9f7ca330e2c2b10051e0d19416534be6 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Thu, 5 May 2016 14:36:25 +0100 Subject: [PATCH 45/56] Adds temporary dependency on my version of pykube, until changes on pykube are pulled and a new release is available. --- lib/galaxy/dependencies/pinned-requirements.txt | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/lib/galaxy/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index 728f2526a87..457f9eb6a6f 100644 --- a/lib/galaxy/dependencies/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pinned-requirements.txt @@ -67,3 +67,7 @@ ecdsa==0.13 # Flexible BAM index naming pysam==0.8.4+gx1 + +# Kubernetes runner pykube requirement +-e git+https://github.com/pcm32/pykube.git@feature/allMergedFeatures#egg=pykube + From 33109d192b870e89769792ceeafbae3563c14f70 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Thu, 5 May 2016 16:55:24 +0100 Subject: [PATCH 46/56] Fixes typo on comment of runner documentation. --- config/job_conf.xml.sample_advanced | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/config/job_conf.xml.sample_advanced b/config/job_conf.xml.sample_advanced index a6b889fa6dc..041495e7d5b 100644 --- a/config/job_conf.xml.sample_advanced +++ b/config/job_conf.xml.sample_advanced @@ -168,7 +168,7 @@ --> /path/to/kubeconfig - <--- This is the path to the kube config file, which is normally on ~/.kube/config, but that will depend on + From 8cb13de22fa0876f53471c433e987e44c3b53e8b Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Thu, 5 May 2016 16:55:52 +0100 Subject: [PATCH 47/56] Fixes other errors spotted by travis. --- lib/galaxy/jobs/runners/kubernetes.py | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 61e2cf08da3..c6da17fc184 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -1,5 +1,5 @@ """ -Job control via a command line interface (e.g. qsub/qstat), possibly over a remote connection (e.g. ssh). +Offload jobs to a Kubernetes cluster. """ import logging @@ -54,8 +54,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): # here we need to fetch the default kubeconfig path from the plugin defined in job_conf... self._pykube_api = HTTPClient(KubeConfig.from_file(self.runner_params["k8s_config_path"])) - self._galaxy_vol_name = "pvc-galaxy" - + self._galaxy_vol_name = "pvc-galaxy" # TODO this needs to be read from params!! self._init_monitor_thread() self._init_worker_threads() @@ -87,8 +86,8 @@ class KubernetesJobRunner(AsynchronousJobRunner): "spec": self.__get_k8s_job_spec(job_wrapper) } - - k8s_job = Job(self._pykube_api, k8s_job_obj).create() + # Creates the Kubernetes Job + Job(self._pykube_api, k8s_job_obj).create() # define job attributes in the AsyncronousJobState for follow-up ajs = AsynchronousJobState(files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper, @@ -97,6 +96,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): # external_runJob_script can be None, in which case it's not used. external_runjob_script = None + return external_runjob_script def __produce_unique_k8s_job_name(self, job_wrapper): # wrapper.get_id_tag() instead of job_id for compatibility with TaskWrappers. From b46b4915bbd93fd75a025c745a7ac2047fe7c507 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 6 May 2016 10:02:29 +0100 Subject: [PATCH 48/56] Small typos/issues fixed for travis. --- config/job_conf.xml.sample_advanced | 2 +- lib/galaxy/jobs/runners/kubernetes.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/config/job_conf.xml.sample_advanced b/config/job_conf.xml.sample_advanced index 041495e7d5b..9141b51499b 100644 --- a/config/job_conf.xml.sample_advanced +++ b/config/job_conf.xml.sample_advanced @@ -624,7 +624,7 @@ + --> diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index c6da17fc184..59ecc27705c 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -54,7 +54,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): # here we need to fetch the default kubeconfig path from the plugin defined in job_conf... self._pykube_api = HTTPClient(KubeConfig.from_file(self.runner_params["k8s_config_path"])) - self._galaxy_vol_name = "pvc-galaxy" # TODO this needs to be read from params!! + self._galaxy_vol_name = "pvc-galaxy" # TODO this needs to be read from params!! self._init_monitor_thread() self._init_worker_threads() From 0016c51cfe146361401beb5922665d1e356894d4 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 13 May 2016 11:57:18 +0100 Subject: [PATCH 49/56] Adheres to pykube new logs method instead of previous get_logs --- 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 59ecc27705c..762aa31ad32 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -261,7 +261,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): 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.get_logs(timestamps=True) + logs += pod.logs(timestamps=True) logs += "\n\n==== Pod " + pod.name + " log end ====" logs_file_path = job_state.output_file logs_file = open(logs_file_path, mode="w") From 46fb489cb4165fbeb198e2305cf27d4faa35367e Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 20 May 2016 17:32:43 +0100 Subject: [PATCH 50/56] Adds the ability to build container ids for pieces for the centralised galaxy containers class and takes into account the default and override modes, suggested by jmchilton. --- lib/galaxy/tools/deps/containers.py | 26 ++++++++++++++++++++++++-- 1 file changed, 24 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/tools/deps/containers.py b/lib/galaxy/tools/deps/containers.py index 948c5b851b4..60ae73886dd 100644 --- a/lib/galaxy/tools/deps/containers.py +++ b/lib/galaxy/tools/deps/containers.py @@ -102,7 +102,25 @@ class ContainerFinder(object): def __overridden_container_id(self, container_type, destination_info): if not self.__container_type_enabled(container_type, destination_info): return None - return destination_info.get("%s_container_id_override" % container_type) + if "%s_container_id_override" % container_type in destination_info: + return destination_info.get("%s_container_id_override" % container_type) + if "%s_image_override" % container_type in destination_info: + return self.__build_container_id_from_parts(container_type, destination_info, mode="override") + + def __build_container_id_from_parts(self, container_type, destination_info, mode): + repo = "" + owner = "" + repo_key = "%s_repo_%s" % container_type, mode + owner_key = "%s_owner_%s" % container_type, mode + if repo_key in destination_info: + repo = destination_info[repo_key] + "/" + if owner_key in destination_info: + owner = destination_info[owner_key] + "/" + cont_id = repo + owner + destination_info["%s_image_%s" % container_type, mode] + tag_key = "%s_tag_%s" % container_type, mode + if tag_key in destination_info: + cont_id += ":" + destination_info[tag_key] + return cont_id def __default_container_id(self, container_type, destination_info): if not self.__container_type_enabled(container_type, destination_info): @@ -111,7 +129,11 @@ class ContainerFinder(object): # Also allow docker_image... if key not in destination_info: key = "%s_image" % container_type - return destination_info.get(key) + if key in destination_info: + return destination_info.get(key) + elif "%s_image_default" in destination_info: + return self.__build_container_id_from_parts(container_type, destination_info, mode="default") + return None def __destination_container(self, container_id, container_type, tool_info, destination_info, job_info): # TODO: ensure destination_info is dict-like From 8e176d4318ee934cbfce924a5f65dbefc196b4a0 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 20 May 2016 17:38:55 +0100 Subject: [PATCH 51/56] Moves the kubernetes runner to use the centralised galaxy way of picking containers, allowing destinations to override a tool, use a tool, or have a default container if the tool doesn't provide one. --- config/job_conf.xml.sample_advanced | 22 ++++++++++++++++++---- lib/galaxy/jobs/runners/kubernetes.py | 2 +- 2 files changed, 19 insertions(+), 5 deletions(-) diff --git a/config/job_conf.xml.sample_advanced b/config/job_conf.xml.sample_advanced index 2680bfa8ae7..5522ecbd2e1 100644 --- a/config/job_conf.xml.sample_advanced +++ b/config/job_conf.xml.sample_advanced @@ -592,10 +592,24 @@ - my-docker-registry.org - superbioinfo - my-tool - latest + my-docker-registry.org + superbioinfo + my-tool + latest + + 3 + + "true" From 84a9ae36e3061d8911fcdfea63773ff31bbc57d3 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 20 May 2016 18:13:36 +0100 Subject: [PATCH 53/56] Removes quotes to docker_enabled paramter value in documentation of destination associated with Kubernetes. --- config/job_conf.xml.sample_advanced | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/config/job_conf.xml.sample_advanced b/config/job_conf.xml.sample_advanced index 85ed4b268b4..306ab63e97e 100644 --- a/config/job_conf.xml.sample_advanced +++ b/config/job_conf.xml.sample_advanced @@ -623,7 +623,7 @@ --> - "true" + true From c484ab37eefa6906be549bbe559fe17754d0b7ec Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 20 May 2016 18:14:31 +0100 Subject: [PATCH 54/56] Adds missing parenthesis for string formatting. --- lib/galaxy/tools/deps/containers.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/tools/deps/containers.py b/lib/galaxy/tools/deps/containers.py index 60ae73886dd..bae23af4ad1 100644 --- a/lib/galaxy/tools/deps/containers.py +++ b/lib/galaxy/tools/deps/containers.py @@ -110,14 +110,14 @@ class ContainerFinder(object): def __build_container_id_from_parts(self, container_type, destination_info, mode): repo = "" owner = "" - repo_key = "%s_repo_%s" % container_type, mode - owner_key = "%s_owner_%s" % container_type, mode + repo_key = "%s_repo_%s" % (container_type, mode) + owner_key = "%s_owner_%s" % (container_type, mode) if repo_key in destination_info: repo = destination_info[repo_key] + "/" if owner_key in destination_info: owner = destination_info[owner_key] + "/" - cont_id = repo + owner + destination_info["%s_image_%s" % container_type, mode] - tag_key = "%s_tag_%s" % container_type, mode + cont_id = repo + owner + destination_info["%s_image_%s" % (container_type, mode)] + tag_key = "%s_tag_%s" % (container_type, mode) if tag_key in destination_info: cont_id += ":" + destination_info[tag_key] return cont_id From 4a908141e3a954582676d5e924fb3c4b41032356 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 20 May 2016 19:17:12 +0100 Subject: [PATCH 55/56] Fixes use of container object returned by underlying galaxy methods. --- 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 69111cc8abc..4b658953acb 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -148,7 +148,7 @@ class KubernetesJobRunner(AsynchronousJobRunner): """Fills in all required for setting up the docker containers to be used.""" k8s_container = { "name": self.__get_k8s_container_name(job_wrapper), - "image": self._find_container(job_wrapper), + "image": self._find_container(job_wrapper).container_id, # this form of command overrides the entrypoint and allows multi command # command line execution, separated by ;, which is what Galaxy does # to assemble the command. From c7929f882e07a9d9d42e0250dfee2bb487f7ee94 Mon Sep 17 00:00:00 2001 From: Pablo Moreno Date: Fri, 20 May 2016 19:20:45 +0100 Subject: [PATCH 56/56] Introduces additional parameter to avoid having the command-line modified for kubernetes. --- lib/galaxy/jobs/command_factory.py | 7 ++++--- lib/galaxy/jobs/runners/__init__.py | 8 ++++++-- lib/galaxy/jobs/runners/kubernetes.py | 4 +++- 3 files changed, 13 insertions(+), 6 deletions(-) diff --git a/lib/galaxy/jobs/command_factory.py b/lib/galaxy/jobs/command_factory.py index 62b5384088f..2bdbb2052f3 100644 --- a/lib/galaxy/jobs/command_factory.py +++ b/lib/galaxy/jobs/command_factory.py @@ -20,6 +20,7 @@ def build_command( runner, job_wrapper, container=None, + modify_command_for_container=True, include_metadata=False, include_work_dir_outputs=True, create_tool_working_directory=True, @@ -57,14 +58,14 @@ def build_command( if not container: __handle_dependency_resolution(commands_builder, job_wrapper, remote_command_params) - if container or job_wrapper.commands_in_new_shell: - if container: + if (container and modify_command_for_container) or job_wrapper.commands_in_new_shell: + if container and modify_command_for_container: # Many Docker containers do not have /bin/bash. external_command_shell = "/bin/sh" else: external_command_shell = shell externalized_commands = __externalize_commands(job_wrapper, external_command_shell, commands_builder, remote_command_params) - if container: + if container and modify_command_for_container: # Stop now and build command before handling metadata and copying # working directory files back. These should always happen outside # of docker container - no security implications when generating diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index d129d46e3b7..9861afc4f10 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -145,7 +145,8 @@ class BaseJobRunner( object ): """ raise NotImplementedError() - def prepare_job(self, job_wrapper, include_metadata=False, include_work_dir_outputs=True): + def prepare_job(self, job_wrapper, include_metadata=False, include_work_dir_outputs=True, + modify_command_for_container=True): """Some sanity checks that all runners' queue_job() methods are likely to want to do """ job_id = job_wrapper.get_id_tag() @@ -171,6 +172,7 @@ class BaseJobRunner( object ): job_wrapper, include_metadata=include_metadata, include_work_dir_outputs=include_work_dir_outputs, + modify_command_for_container=modify_command_for_container ) except Exception as e: log.exception("(%s) Failure preparing job" % job_id) @@ -193,13 +195,15 @@ class BaseJobRunner( object ): def recover(self, job, job_wrapper): raise NotImplementedError() - def build_command_line( self, job_wrapper, include_metadata=False, include_work_dir_outputs=True ): + def build_command_line( self, job_wrapper, include_metadata=False, include_work_dir_outputs=True, + modify_command_for_container=True ): container = self._find_container( job_wrapper ) return build_command( self, job_wrapper, include_metadata=include_metadata, include_work_dir_outputs=include_work_dir_outputs, + modify_command_for_container=modify_command_for_container, container=container ) diff --git a/lib/galaxy/jobs/runners/kubernetes.py b/lib/galaxy/jobs/runners/kubernetes.py index 4b658953acb..3829ed8da7f 100644 --- a/lib/galaxy/jobs/runners/kubernetes.py +++ b/lib/galaxy/jobs/runners/kubernetes.py @@ -64,7 +64,9 @@ class KubernetesJobRunner(AsynchronousJobRunner): # prepare the job # We currently don't need to include_metadata or include_work_dir_outputs, as working directory is the same # were galaxy will expect results. - if not self.prepare_job(job_wrapper, include_metadata=False, include_work_dir_outputs=False): + log.debug("Starting queue_job for job " + job_wrapper.get_id_tag()) + if not self.prepare_job(job_wrapper, include_metadata=False, include_work_dir_outputs=False, + modify_command_for_container=False): return job_destination = job_wrapper.job_destination