Merge pull request #2487 from nsoranzo/release_16.04_reset_exception_retries

[16.04] Reset exception retries after successful job status check
This commit is contained in:
Martin Cech
2016-06-13 11:51:50 -04:00
committed by GitHub
+11 -7
View File
@@ -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