From 0fa7f79dec0fe8f4feb2551d360dadfd9ad8c087 Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Tue, 7 Jan 2014 15:13:01 -0500 Subject: [PATCH 1/5] Allow handlers to control which runner plugins they will load. --- lib/galaxy/jobs/__init__.py | 17 ++++++++++++++--- lib/galaxy/jobs/handler.py | 2 +- 2 files changed, 15 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index eb81d47df39..6b5c30df296 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -86,6 +86,7 @@ class JobConfiguration( object ): self.app = app self.runner_plugins = [] self.handlers = {} + self.handler_runner_plugins = {} self.default_handler_id = None self.destinations = {} self.destination_tags = {} @@ -138,6 +139,10 @@ class JobConfiguration( object ): else: log.debug("Read definition for handler '%s'" % id) self.handlers[id] = (id,) + for plugin in handler.findall('plugin'): + if id not in self.handler_runner_plugins: + self.handler_runner_plugins[id] = [] + self.handler_runner_plugins[id].append( plugin.get('id') ) if handler.get('tags', None) is not None: for tag in [ x.strip() for x in handler.get('tags').split(',') ]: if tag in self.handlers: @@ -420,13 +425,19 @@ class JobConfiguration( object ): """ return self.destinations.get(id_or_tag, None) - def get_job_runner_plugins(self): + def get_job_runner_plugins(self, handler_id): """Load all configured job runner plugins :returns: list of job runner plugins """ rval = {} - for runner in self.runner_plugins: + if handler_id in self.handler_runner_plugins: + plugins_to_load = [ rp for rp in self.runner_plugins if rp['id'] in self.handler_runner_plugins[handler_id] ] + log.info( "Handler '%s' will load specified runner plugins: %s", handler_id, ', '.join( [ rp['id'] for rp in plugins_to_load ] ) ) + else: + plugins_to_load = self.runner_plugins + log.info( "Handler '%s' will load all configured runner plugins", handler_id ) + for runner in plugins_to_load: class_names = [] module = None id = runner['id'] @@ -477,7 +488,7 @@ class JobConfiguration( object ): try: rval[id] = runner_class( self.app, runner[ 'workers' ], **runner.get( 'kwds', {} ) ) except TypeError: - log.warning( "Job runner '%s:%s' has not been converted to a new-style runner" % ( module_name, class_name ) ) + log.exception( "Job runner '%s:%s' has not been converted to a new-style runner or encountered TypeError on load" % ( module_name, class_name ) ) rval[id] = runner_class( self.app ) log.debug( "Loaded job runner '%s:%s' as '%s'" % ( module_name, class_name, id ) ) return rval diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index ae6f2ac3d38..5a262ec4633 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -565,7 +565,7 @@ class JobHandlerStopQueue( object ): class DefaultJobDispatcher( object ): def __init__( self, app ): self.app = app - self.job_runners = self.app.job_config.get_job_runner_plugins() + self.job_runners = self.app.job_config.get_job_runner_plugins( self.app.config.server_name ) # Once plugins are loaded, all job destinations that were created from # URLs can have their URL params converted to the destination's param # dict by the plugin. From 458df6faf2f33f4c4260f47e4c06be4c87cdd826 Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Tue, 7 Jan 2014 15:18:11 -0500 Subject: [PATCH 2/5] Allow runners plugins to accept params, provide some simple load-time validation logic for params. --- lib/galaxy/jobs/runners/__init__.py | 34 +++++++++++++++++++++++++---- 1 file changed, 30 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index 2b412a6db31..2be4f91c405 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -22,13 +22,39 @@ log = logging.getLogger( __name__ ) STOP_SIGNAL = object() + +class RunnerParams( object ): + + def __init__( self, specs = None, params = None ): + self.specs = specs or dict() + self.params = params or dict() + for name, value in self.params.items(): + assert name in self.specs, 'Invalid job runner parameter for this plugin: %s' % name + if 'map' in self.specs[ name ]: + try: + self.params[ name ] = self.specs[ name ][ 'map' ]( value ) + except Exception, e: + raise Exception( 'Job runner parameter "%s" value "%s" could not be converted to the correct type: %s' % ( name, value, e ) ) + if 'valid' in self.specs[ name ]: + assert self.specs[ name ][ 'valid' ]( value ), 'Job runner parameter %s failed validation' % name + + def __getattr__( self, name ): + return self.params.get( name, self.specs[ name ][ 'default' ] ) + + class BaseJobRunner( object ): - def __init__( self, app, nworkers ): + def __init__( self, app, nworkers, **kwargs ): """Start the job runner """ self.app = app self.sa_session = app.model.context self.nworkers = nworkers + runner_param_specs = dict( recheck_missing_job_retries = dict( map = int, valid = lambda x: x >= 0, default = 0 ) ) + if 'runner_param_specs' in kwargs: + runner_param_specs.update( kwargs.pop( 'runner_param_specs' ) ) + if kwargs: + log.debug( 'Loading %s with params: %s', self.runner_name, kwargs ) + self.runner_params = RunnerParams( specs = runner_param_specs, params = kwargs ) def _init_worker_threads(self): """Start ``nworkers`` worker threads. @@ -115,7 +141,7 @@ class BaseJobRunner( object ): job_wrapper.cleanup() return False elif job_state != model.Job.states.QUEUED: - log.info( "(%d) Job is in state %s, skipping execution" % ( job_id, job_state ) ) + log.info( "(%s) Job is in state %s, skipping execution" % ( job_id, job_state ) ) # cleanup may not be safe in all states return False @@ -287,8 +313,8 @@ class AsynchronousJobRunner( BaseJobRunner ): to the correct methods (queue, finish, cleanup) at appropriate times.. """ - def __init__( self, app, nworkers ): - super( AsynchronousJobRunner, self ).__init__( app, nworkers ) + def __init__( self, app, nworkers, **kwargs ): + super( AsynchronousJobRunner, self ).__init__( app, nworkers, **kwargs ) # 'watched' and 'queue' are both used to keep track of jobs to watch. # 'queue' is used to add new watched jobs, and can be called from # any thread (usually by the 'queue_job' method). 'watched' must only From 4a024b882f504aded1db32f1ed056f7c9109d218 Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Tue, 7 Jan 2014 15:21:50 -0500 Subject: [PATCH 3/5] Make job terminal state logic somewhat configurable via runner params, update DRMAA runner plugin changes for configurable terminal state handling. Also allow the drmaa plugin to read the drmaa library path from job_conf.xml rather than the $DRMAA_LIBRARY_PATH environment variable. --- lib/galaxy/jobs/runners/__init__.py | 4 + lib/galaxy/jobs/runners/drmaa.py | 116 +++++++++++++++++++--------- 2 files changed, 84 insertions(+), 36 deletions(-) diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index 2be4f91c405..334b20c4850 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -252,6 +252,10 @@ class BaseJobRunner( object ): options.update(**kwds) return job_script(**options) + def _complete_terminal_job( self, ajs, **kwargs ): + if ajs.job_wrapper.get_state() != model.Job.states.DELETED: + self.work_queue.put( ( self.finish_job, ajs ) ) + class AsynchronousJobState( object ): """ diff --git a/lib/galaxy/jobs/runners/drmaa.py b/lib/galaxy/jobs/runners/drmaa.py index 465d5e9f6ad..9039abc56f3 100644 --- a/lib/galaxy/jobs/runners/drmaa.py +++ b/lib/galaxy/jobs/runners/drmaa.py @@ -16,27 +16,12 @@ from galaxy.jobs import JobDestination from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner eggs.require( "drmaa" ) -# We foolishly named this file the same as the name exported by the drmaa -# library... 'import drmaa' imports itself. -drmaa = __import__( "drmaa" ) log = logging.getLogger( __name__ ) __all__ = [ 'DRMAAJobRunner' ] -drmaa_state = { - drmaa.JobState.UNDETERMINED: 'process status cannot be determined', - drmaa.JobState.QUEUED_ACTIVE: 'job is queued and active', - drmaa.JobState.SYSTEM_ON_HOLD: 'job is queued and in system hold', - drmaa.JobState.USER_ON_HOLD: 'job is queued and in user hold', - drmaa.JobState.USER_SYSTEM_ON_HOLD: 'job is queued and in user and system hold', - drmaa.JobState.RUNNING: 'job is running', - drmaa.JobState.SYSTEM_SUSPENDED: 'job is system suspended', - drmaa.JobState.USER_SUSPENDED: 'job is user suspended', - drmaa.JobState.DONE: 'job finished normally', - drmaa.JobState.FAILED: 'job finished, but failed', -} - +drmaa = None DRMAA_jobTemplate_attributes = [ 'args', 'remoteCommand', 'outputPath', 'errorPath', 'nativeSpecification', 'jobName', 'email', 'project' ] @@ -48,8 +33,50 @@ class DRMAAJobRunner( AsynchronousJobRunner ): """ runner_name = "DRMAARunner" - def __init__( self, app, nworkers ): + def __init__( self, app, nworkers, **kwargs ): """Start the job runner""" + + global drmaa + + runner_param_specs = dict( + drmaa_library_path = dict( map = str, default = os.environ.get( 'DRMAA_LIBRARY_PATH', None ) ), + invalidjobexception_state = dict( map = str, valid = lambda x: x in ( model.Job.states.OK, model.Job.states.ERROR ), default = model.Job.states.OK ), + invalidjobexception_retries = dict( map = int, valid = lambda x: int >= 0, default = 0 ), + internalexception_state = dict( map = str, valid = lambda x: x in ( model.Job.states.OK, model.Job.states.ERROR ), default = model.Job.states.OK ), + internalexception_retries = dict( map = int, valid = lambda x: int >= 0, default = 0 ) ) + + if 'runner_param_specs' not in kwargs: + kwargs[ 'runner_param_specs' ] = dict() + kwargs[ 'runner_param_specs' ].update( runner_param_specs ) + + super( DRMAAJobRunner, self ).__init__( app, nworkers, **kwargs ) + + # This allows multiple drmaa runners (although only one per handler) in the same job config file + if 'drmaa_library_path' in kwargs: + log.info( 'Overriding DRMAA_LIBRARY_PATH due to runner plugin parameter: %s', self.runner_params.drmaa_library_path ) + os.environ['DRMAA_LIBRARY_PATH'] = self.runner_params.drmaa_library_path + + # We foolishly named this file the same as the name exported by the drmaa + # library... 'import drmaa' imports itself. + drmaa = __import__( "drmaa" ) + + # Subclasses may need access to state constants + self.drmaa_job_states = drmaa.JobState + + # Descriptive state strings pulled from the drmaa lib itself + self.drmaa_job_state_strings = { + drmaa.JobState.UNDETERMINED: 'process status cannot be determined', + drmaa.JobState.QUEUED_ACTIVE: 'job is queued and active', + drmaa.JobState.SYSTEM_ON_HOLD: 'job is queued and in system hold', + drmaa.JobState.USER_ON_HOLD: 'job is queued and in user hold', + drmaa.JobState.USER_SYSTEM_ON_HOLD: 'job is queued and in user and system hold', + drmaa.JobState.RUNNING: 'job is running', + drmaa.JobState.SYSTEM_SUSPENDED: 'job is system suspended', + drmaa.JobState.USER_SUSPENDED: 'job is user suspended', + drmaa.JobState.DONE: 'job finished normally', + drmaa.JobState.FAILED: 'job finished, but failed', + } + self.ds = drmaa.Session() self.ds.initialize() @@ -58,7 +85,6 @@ class DRMAAJobRunner( AsynchronousJobRunner ): self.external_killJob_script = app.config.drmaa_external_killjob_script self.userid = None - super( DRMAAJobRunner, self ).__init__( app, nworkers ) self._init_monitor_thread() self._init_worker_threads() @@ -175,6 +201,20 @@ class DRMAAJobRunner( AsynchronousJobRunner ): # Add to our 'queue' of jobs to monitor self.monitor_queue.put( ajs ) + def _complete_terminal_job( self, ajs, drmaa_state, **kwargs ): + """ + Handle a job upon its termination in the DRM. This method is meant to + be overridden by subclasses to improve post-mortem and reporting of + failures. + """ + if drmaa_state == drmaa.JobState.FAILED: + if ajs.job_wrapper.get_state() != model.Job.states.DELETED: + ajs.stop_job = False + ajs.fail_message = "The cluster DRM system terminated this job" + self.work_queue.put( ( self.fail_job, ajs ) ) + elif drmaa_state == drmaa.JobState.DONE: + super( DRMAAJobRunner, self )._complete_terminal_job( ajs ) + def check_watched_items( self ): """ Called by the monitor thread to look at each watched job and deal @@ -188,16 +228,27 @@ class DRMAAJobRunner( AsynchronousJobRunner ): try: assert external_job_id not in ( None, 'None' ), '(%s/%s) Invalid job id' % ( galaxy_id_tag, external_job_id ) state = self.ds.jobStatus( external_job_id ) - # InternalException was reported to be necessary on some DRMs, but - # this could cause failures to be detected as completion! Please - # report if you experience problems with this. - except ( drmaa.InvalidJobException, drmaa.InternalException ), e: - # we should only get here if an orphaned job was put into the queue at app startup - log.info( "(%s/%s) job left DRM queue with following message: %s" % ( galaxy_id_tag, external_job_id, e ) ) - self.work_queue.put( ( self.finish_job, ajs ) ) + except ( drmaa.InternalException, drmaa.InvalidJobException ), e: + ecn = e.__class__.__name__ + retry_param = ecn.lower() + '_retries' + state_param = ecn.lower() + '_state' + retries = getattr( ajs, retry_param, 0 ) + if self.runner_params[ retry_param ] > 0: + if retries < self.runner_params[ retry_param ]: + # will retry check on next iteration + setattr( ajs, retry_param, retries + 1 ) + continue + if self.runner_params[ state_param ] == model.Job.states.OK: + log.info( "(%s/%s) job left DRM queue with following message: %s", galaxy_id_tag, external_job_id, e ) + self.work_queue.put( ( self.finish_job, ajs ) ) + elif self.runner_params[ state_param ] == model.Job.states.ERROR: + log.info( "(%s/%s) job check resulted in %s after %s tries: %s", galaxy_id_tag, external_job_id, ecn, retries, e ) + self.work_queue.put( ( self.fail_job, ajs ) ) + else: + raise Exception( "%s is set to an invalid value (%s), this should not be possible. See galaxy.jobs.drmaa.__init__()", state_param, self.runner_params[ state_param ] ) continue except drmaa.DrmCommunicationException, e: - log.warning( "(%s/%s) unable to communicate with DRM: %s" % ( galaxy_id_tag, external_job_id, e )) + log.warning( "(%s/%s) unable to communicate with DRM: %s", galaxy_id_tag, external_job_id, e ) new_watched.append( ajs ) continue except Exception, e: @@ -208,19 +259,12 @@ class DRMAAJobRunner( AsynchronousJobRunner ): self.work_queue.put( ( self.fail_job, ajs ) ) continue if state != old_state: - log.debug( "(%s/%s) state change: %s" % ( galaxy_id_tag, external_job_id, drmaa_state[state] ) ) + log.debug( "(%s/%s) state change: %s" % ( galaxy_id_tag, external_job_id, self.drmaa_job_state_strings[state] ) ) if state == drmaa.JobState.RUNNING and not ajs.running: ajs.running = True ajs.job_wrapper.change_state( model.Job.states.RUNNING ) - if state == drmaa.JobState.FAILED: - if ajs.job_wrapper.get_state() != model.Job.states.DELETED: - ajs.stop_job = False - ajs.fail_message = "The cluster DRM system terminated this job" - self.work_queue.put( ( self.fail_job, ajs ) ) - continue - if state == drmaa.JobState.DONE: - if ajs.job_wrapper.get_state() != model.Job.states.DELETED: - self.work_queue.put( ( self.finish_job, ajs ) ) + if state in ( drmaa.JobState.FAILED, drmaa.JobState.DONE ): + self._complete_terminal_job( ajs, drmaa_state = state ) continue ajs.old_state = state new_watched.append( ajs ) From 1889852e8f667f9a337279d3e5be2e4bc7f76538 Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Tue, 7 Jan 2014 15:23:11 -0500 Subject: [PATCH 4/5] Add a Slurm runner (subclassed from the drmaa runner), adds some logic for determining/reporting the reason of job failure and requeuing jobs that failed due to node failure. --- lib/galaxy/jobs/runners/slurm.py | 57 ++++++++++++++++++++++++++++++++ 1 file changed, 57 insertions(+) create mode 100644 lib/galaxy/jobs/runners/slurm.py diff --git a/lib/galaxy/jobs/runners/slurm.py b/lib/galaxy/jobs/runners/slurm.py new file mode 100644 index 00000000000..51b7df56c79 --- /dev/null +++ b/lib/galaxy/jobs/runners/slurm.py @@ -0,0 +1,57 @@ +""" +SLURM job control via the DRMAA API. +""" + +import time +import logging +import subprocess + +from galaxy import model +from galaxy.jobs.runners.drmaa import DRMAAJobRunner + +log = logging.getLogger( __name__ ) + +__all__ = [ 'SlurmJobRunner' ] + + +class SlurmJobRunner( DRMAAJobRunner ): + runner_name = "SlurmRunner" + + def _complete_terminal_job( self, ajs, drmaa_state, **kwargs ): + def __get_jobinfo(): + scontrol_out = subprocess.check_output( ( 'scontrol', '-o', 'show', 'job', ajs.job_id ) ) + return dict( [ out_param.split( '=', 1 ) for out_param in scontrol_out.split() ] ) + if drmaa_state == self.drmaa_job_states.FAILED: + try: + job_info = __get_jobinfo() + sleep = 1 + while job_info['JobState'] == 'COMPLETING': + log.debug( '(%s/%s) Waiting %s seconds for failed job to exit COMPLETING state for post-mortem', ajs.job_wrapper.get_id_tag(), ajs.job_id, sleep ) + time.sleep( sleep ) + sleep *= 2 + if sleep > 64: + ajs.fail_message = "This job failed and the system timed out while trying to determine the cause of the failure." + break + job_info = __get_jobinfo() + if job_info['JobState'] == 'TIMEOUT': + ajs.fail_message = "This job was terminated because it ran longer than the maximum allowed job run time." + elif job_info['JobState'] == 'NODE_FAIL': + log.warning( '(%s/%s) Job failed due to node failure, attempting resubmission', ajs.job_wrapper.get_id_tag(), ajs.job_id ) + ajs.job_wrapper.change_state( model.Job.states.QUEUED, info = 'Job was resubmitted due to node failure' ) + try: + self.queue_job( ajs.job_wrapper ) + return + except: + ajs.fail_message = "This job failed due to a cluster node failure, and an attempt to resubmit the job failed." + elif job_info['JobState'] == 'CANCELLED': + ajs.fail_message = "This job failed because it was cancelled by an administrator." + else: + ajs.fail_message = "This job failed for reasons that could not be determined." + ajs.fail_message += '\nPlease click the bug icon to report this problem if you need help.' + ajs.stop_job = False + self.work_queue.put( ( self.fail_job, ajs ) ) + except Exception, e: + log.exception( '(%s/%s) Unable to inspect failed slurm job using scontrol, job will be unconditionally failed: %s', ajs.job_wrapper.get_id_tag(), ajs.job_id, e ) + super( SlurmJobRunner, self )._complete_terminal_job( ajs, drmaa_state = drmaa_state ) + elif drmaa_state == self.drmaa_job_states.DONE: + super( SlurmJobRunner, self )._complete_terminal_job( ajs, drmaa_state = drmaa_state ) From e1913c1b75bf237afdfaa7d97cd9070a147884f8 Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Tue, 7 Jan 2014 15:23:38 -0500 Subject: [PATCH 5/5] Sample job_conf.xml for new features introduced. --- job_conf.xml.sample_advanced | 24 +++++++++++++++++++++++- 1 file changed, 23 insertions(+), 1 deletion(-) diff --git a/job_conf.xml.sample_advanced b/job_conf.xml.sample_advanced index 2637c7ce319..8c9741241bb 100644 --- a/job_conf.xml.sample_advanced +++ b/job_conf.xml.sample_advanced @@ -6,7 +6,19 @@ --> - + + + ok + 0 + ok + 0 + + + + /sge/lib/libdrmaa.so + @@ -14,6 +26,7 @@ + + + + +