diff --git a/lib/galaxy/jobs/runners/drmaa.py b/lib/galaxy/jobs/runners/drmaa.py index abf1c58571a..e64832e7e93 100644 --- a/lib/galaxy/jobs/runners/drmaa.py +++ b/lib/galaxy/jobs/runners/drmaa.py @@ -22,6 +22,8 @@ log = logging.getLogger( __name__ ) __all__ = [ 'DRMAAJobRunner' ] +RETRY_EXCEPTIONS_LOWER = frozenset(['invalidjobexception', 'internalexception']) + class DRMAAJobRunner( AsynchronousJobRunner ): """ @@ -34,12 +36,11 @@ class DRMAAJobRunner( AsynchronousJobRunner ): """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 ) ) + runner_param_specs = { + 'drmaa_library_path': dict( map=str, default=os.environ.get( 'DRMAA_LIBRARY_PATH', None ) ) } + for retry_exception in RETRY_EXCEPTIONS_LOWER: + runner_param_specs[retry_exception + '_state'] = dict( map=str, valid=lambda x: x in ( model.Job.states.OK, model.Job.states.ERROR ), default=model.Job.states.OK ) + runner_param_specs[retry_exception + '_retries'] = dict( map=int, valid=lambda x: int >= 0, default=0 ) if 'runner_param_specs' not in kwargs: kwargs[ 'runner_param_specs' ] = dict() @@ -246,12 +247,15 @@ 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.job_status( external_job_id ) + # Reset exception retries + for retry_exception in RETRY_EXCEPTIONS_LOWER: + setattr( ajs, retry_exception + '_retries', 0) except ( drmaa.InternalException, drmaa.InvalidJobException ) as e: ecn = type(e).__name__ retry_param = ecn.lower() + '_retries' state_param = ecn.lower() + '_state' retries = getattr( ajs, retry_param, 0 ) - log.warning("(%s/%s) unable to check job status because of %s exception for %d tries: %s", galaxy_id_tag, external_job_id, ecn, retries + 1, e) + log.warning("(%s/%s) unable to check job status because of %s exception for %d consecutive tries: %s", galaxy_id_tag, external_job_id, ecn, retries + 1, e) if self.runner_params[ retry_param ] > 0: if retries < self.runner_params[ retry_param ]: # will retry check on next iteration