Merge pull request #3217 from interactomix/dev

Per user total walltime limit.
This commit is contained in:
Björn Grüning
2016-12-03 14:11:39 +01:00
committed by GitHub
4 changed files with 66 additions and 3 deletions
+7
View File
@@ -762,6 +762,13 @@
will be terminated by Galaxy.
-->
<limit type="walltime">24:00:00</limit>
<!-- total_walltime:
Total walltime that jobs may not exceed during a set period.
If total walltime of finished jobs exceeds this value, any
new jobs are paused. `window` is a number in days,
representing the period.
-->
<limit type="total_walltime" window="30">24:00:00</limit>
<!-- output_size:
Size that any defined tool output can grow to before the job
will be terminated. This does not include temporary files
+17
View File
@@ -265,12 +265,14 @@ class JobConfiguration( object ):
types = dict(registered_user_concurrent_jobs=int,
anonymous_user_concurrent_jobs=int,
walltime=str,
total_walltime=str,
output_size=util.size_to_bytes)
self.limits = Bunch(registered_user_concurrent_jobs=None,
anonymous_user_concurrent_jobs=None,
walltime=None,
walltime_delta=None,
total_walltime={},
output_size=None,
destination_user_concurrent_jobs={},
destination_total_concurrent_jobs={})
@@ -287,6 +289,13 @@ class JobConfiguration( object ):
self.limits.destination_total_concurrent_jobs[id] = int(limit.text)
else:
self.limits.destination_user_concurrent_jobs[id] = int(limit.text)
elif type == 'total_walltime':
self.limits.total_walltime["window"] = (
int( limit.get('window') ) or 30
)
self.limits.total_walltime["raw"] = (
types.get(type, str)(limit.text)
)
elif limit.text:
self.limits.__dict__[type] = types.get(type, str)(limit.text)
@@ -294,6 +303,13 @@ class JobConfiguration( object ):
h, m, s = [ int( v ) for v in self.limits.walltime.split( ':' ) ]
self.limits.walltime_delta = datetime.timedelta( 0, s, 0, 0, m, h )
if "raw" in self.limits.total_walltime:
h, m, s = [ int( v ) for v in
self.limits.total_walltime["raw"].split( ':' ) ]
self.limits.total_walltime["delta"] = datetime.timedelta(
0, s, 0, 0, m, h
)
log.debug('Done loading job configuration')
def __parse_job_conf_legacy(self):
@@ -356,6 +372,7 @@ class JobConfiguration( object ):
anonymous_user_concurrent_jobs=self.app.config.anonymous_user_job_limit,
walltime=self.app.config.job_walltime,
walltime_delta=self.app.config.job_walltime_delta,
total_walltime={},
output_size=self.app.config.output_size_limit,
destination_user_concurrent_jobs={},
destination_total_concurrent_jobs={})
+39 -3
View File
@@ -2,6 +2,7 @@
Galaxy job handler, prepares, runs, tracks, and finishes Galaxy jobs
"""
import datetime
import os
import time
import logging
@@ -18,7 +19,7 @@ from galaxy.jobs.mapper import JobNotReadyException
log = logging.getLogger( __name__ )
# States for running a job. These are NOT the same as data states
JOB_WAIT, JOB_ERROR, JOB_INPUT_ERROR, JOB_INPUT_DELETED, JOB_READY, JOB_DELETED, JOB_ADMIN_DELETED, JOB_USER_OVER_QUOTA = 'wait', 'error', 'input_error', 'input_deleted', 'ready', 'deleted', 'admin_deleted', 'user_over_quota'
JOB_WAIT, JOB_ERROR, JOB_INPUT_ERROR, JOB_INPUT_DELETED, JOB_READY, JOB_DELETED, JOB_ADMIN_DELETED, JOB_USER_OVER_QUOTA, JOB_USER_OVER_TOTAL_WALLTIME = 'wait', 'error', 'input_error', 'input_deleted', 'ready', 'deleted', 'admin_deleted', 'user_over_quota', 'user_over_total_walltime'
DEFAULT_JOB_PUT_FAILURE_MESSAGE = 'Unable to run job due to a misconfiguration of the Galaxy job running system. Please contact a site administrator.'
@@ -284,8 +285,13 @@ class JobHandlerQueue( object ):
log.info( "(%d) Job deleted by user while still queued" % job.id )
elif job_state == JOB_ADMIN_DELETED:
log.info( "(%d) Job deleted by admin while still queued" % job.id )
elif job_state == JOB_USER_OVER_QUOTA:
log.info( "(%d) User (%s) is over quota: job paused" % ( job.id, job.user_id ) )
elif job_state in ( JOB_USER_OVER_QUOTA,
JOB_USER_OVER_TOTAL_WALLTIME ):
if job_state == JOB_USER_OVER_QUOTA:
log.info( "(%d) User (%s) is over quota: job paused" % ( job.id, job.user_id ) )
else:
log.info( "(%d) User (%s) is over total walltime limit: job paused" % ( job.id, job.user_id ) )
job.set_state( model.Job.states.PAUSED )
for dataset_assoc in job.output_datasets + job.output_library_datasets:
dataset_assoc.dataset.dataset.state = model.Dataset.states.PAUSED
@@ -375,6 +381,7 @@ class JobHandlerQueue( object ):
# job is ready to run, check limits
# TODO: these checks should be refactored to minimize duplication and made more modular/pluggable
state = self.__check_destination_jobs( job, job_wrapper )
if state == JOB_READY:
state = self.__check_user_jobs( job, job_wrapper )
if state == JOB_READY and self.app.config.enable_quotas:
@@ -386,6 +393,35 @@ class JobHandlerQueue( object ):
return JOB_USER_OVER_QUOTA, job_destination
except AssertionError as e:
pass # No history, should not happen with an anon user
# Check total walltime limits
if ( state == JOB_READY and
"delta" in self.app.job_config.limits.total_walltime ):
jobs_to_check = self.sa_session.query( model.Job ).filter(
model.Job.user_id == job.user.id,
model.Job.update_time >= datetime.datetime.now() -
datetime.timedelta(
self.app.job_config.limits.total_walltime["window"]
),
model.Job.state == 'ok'
).all()
time_spent = datetime.timedelta(0)
for job in jobs_to_check:
# History is job.state_history
started = None
finished = None
for history in sorted(
job.state_history,
key=lambda history: history.update_time ):
if history.state == "running":
started = history.create_time
elif history.state == "ok":
finished = history.create_time
time_spent += finished - started
if time_spent > self.app.job_config.limits.total_walltime["delta"]:
return JOB_USER_OVER_TOTAL_WALLTIME, job_destination
return state, job_destination
def __verify_in_memory_job_inputs( self, job ):
+3
View File
@@ -100,6 +100,7 @@ class JobConfXmlParserTestCase( unittest.TestCase ):
assert limits.anonymous_user_concurrent_jobs is None
assert limits.walltime is None
assert limits.walltime_delta is None
assert limits.total_walltime == {}
assert limits.output_size is None
assert limits.destination_user_concurrent_jobs == {}
assert limits.destination_total_concurrent_jobs == {}
@@ -113,6 +114,8 @@ class JobConfXmlParserTestCase( unittest.TestCase ):
assert limits.destination_user_concurrent_jobs[ "mycluster" ] == 2
assert limits.destination_user_concurrent_jobs[ "longjobs" ] == 1
assert limits.walltime_delta == datetime.timedelta( 0, 0, 0, 0, 0, 24 )
assert limits.total_walltime["delta"] == datetime.timedelta( 0, 0, 0, 0, 0, 24)
assert limits.total_walltime["window"] == 30
def test_env_parsing( self ):
self.__with_advanced_config()