Sun Grid Engine support

This commit is contained in:
Nate Coraor
2008-04-15 22:04:09 +00:00
parent 8056ea2e94
commit 9bca1b9db2
5 changed files with 405 additions and 99 deletions
+3 -2
View File
@@ -31,8 +31,9 @@ class UniverseApplication( object ):
#Load datatype converters
self.datatypes_registry.load_datatype_converters(self.config.datatype_converters_config, self.config.datatype_converters_path, self.toolbox)
# Start the job queue
self.job_queue = jobs.JobQueue( self )
self.job_stop_queue = jobs.JobStopQueue( self )
job_dispatcher = jobs.DefaultJobDispatcher( self )
self.job_queue = jobs.JobQueue( self, job_dispatcher )
self.job_stop_queue = jobs.JobStopQueue( self, job_dispatcher )
self.heartbeat = None
# Start the heartbeat process if configured and available
if self.config.use_heartbeat:
+34 -61
View File
@@ -39,7 +39,7 @@ class JobQueue( object ):
a JobRunner.
"""
STOP_SIGNAL = object()
def __init__( self, app ):
def __init__( self, app, dispatcher ):
"""Start the job manager"""
self.app = app
# Should we use IPC to communicate (needed if forking)
@@ -82,7 +82,7 @@ class JobQueue( object ):
# Helper for interruptable sleep
self.sleeper = Sleeper()
self.running = True
self.dispatcher = DefaultJobDispatcher( app )
self.dispatcher = dispatcher
self.monitor_thread = threading.Thread( target=self.monitor )
self.monitor_thread.start()
log.info( "job manager started" )
@@ -97,50 +97,14 @@ class JobQueue( object ):
model = self.app.model
# Jobs in the NEW state won't be requeued unless we're tracking in the database
if not self.track_jobs_in_database:
for j in model.Job.select( model.Job.c.state == model.Job.states.NEW ):
log.debug( "no runner: %s is still in new state, adding to the jobs queue" %j.id )
self.queue.put( ( j.id, j.tool_id ) )
for j in model.Job.select( (model.Job.c.state == model.Job.states.RUNNING)
for job in model.Job.select( model.Job.c.state == model.Job.states.NEW ):
log.debug( "no runner: %s is still in new state, adding to the jobs queue" %job.id )
self.queue.put( ( job.id, job.tool_id ) )
for job in model.Job.select( (model.Job.c.state == model.Job.states.RUNNING)
| (model.Job.c.state == model.Job.states.QUEUED) ):
job_wrapper = JobWrapper( j.id, self.app.toolbox.tools_by_id[ j.tool_id ], self )
if j.job_runner_name is None:
# should only happen with stuff in the database prior to job_runner_name column
log.error( "job %s was queued but has no runner, this should not be possible" % j.id )
job_wrapper.change_state( model.Job.states.ERROR, info = "This job was killed when Galaxy was restarted. Please retry the job." )
elif j.job_runner_name.startswith('local://'):
# Local jobs are never in the QUEUED state
log.debug( "local runner: %s is still in running state, setting to error" %j.id )
job_wrapper.change_state( model.Job.states.ERROR, info = "This job was killed when Galaxy was restarted. Please retry the job." )
elif j.job_runner_name.startswith('pbs://'):
pbs_job_state = runners.pbs.PBSJobState()
pbs_job_state.ofile = "%s/database/pbs/%s.o" % (os.getcwd(), j.id)
pbs_job_state.efile = "%s/database/pbs/%s.e" % (os.getcwd(), j.id)
pbs_job_state.job_file = "%s/database/pbs/%s.sh" % (os.getcwd(), j.id)
pbs_job_state.job_id = str( j.job_runner_external_id )
pbs_job_state.runner_url = job_wrapper.tool.job_runner
job_wrapper.command_line = j.command_line
pbs_job_state.job_wrapper = job_wrapper
# The job was in PBS and actually running
if j.state == model.Job.states.RUNNING:
log.debug( "pbs runner: %s is still in running state, adding to the PBS queue" %j.id )
pbs_job_state.old_state = 'R'
pbs_job_state.running = True
self.dispatcher.job_runners["pbs"].queue.put( pbs_job_state )
elif j.state == model.Job.states.QUEUED:
# The job was in PBS but not yet running
if j.job_runner_external_id:
log.debug( "pbs runner: %s is still in PBS queued state, adding to the PBS queue" %j.id )
pbs_job_state.old_state = 'Q'
pbs_job_state.running = False
self.dispatcher.job_runners["pbs"].queue.put( pbs_job_state )
# The job had reached the pbs queue_job method, but not the PBS queue itself.
else:
log.debug( "pbs runner: %s is still in jobs queued state, adding to the jobs queue" %j.id )
job_wrapper.change_state( model.Job.states.NEW )
# is there a way we can do this without blocking startup?
self.dispatcher.put( job_wrapper )
else:
log.error( "job %s has an unknown job runner: %s" %( j.id, j.job_runner_name ) )
# why are we passing the queue to the wrapper?
job_wrapper = JobWrapper( job.id, self.app.toolbox.tools_by_id[ job.tool_id ], self )
self.dispatcher.recover( job, job_wrapper )
def monitor( self ):
"""
@@ -341,6 +305,10 @@ class JobWrapper( object ):
"""
job = model.Job.get( self.job_id )
job.refresh()
# if the job was deleted, don't fail it
if job.state == job.states.DELETED:
self.cleanup()
return
for dataset_assoc in job.output_datasets:
dataset = dataset_assoc.dataset
dataset.refresh()
@@ -423,6 +391,7 @@ class JobWrapper( object ):
# default post job setup
mapping.context.current.clear()
job = model.Job.get( self.job_id )
# TODO: change PBS to use fail() instead of finish()
# if the job was deleted, don't finish it
if job.state == job.states.DELETED:
self.cleanup()
@@ -518,6 +487,9 @@ class DefaultJobDispatcher( object ):
elif runner_name == "pbs":
import runners.pbs
self.job_runners[runner_name] = runners.pbs.PBSJobRunner( app )
elif runner_name == "sge":
import runners.sge
self.job_runners[runner_name] = runners.sge.SGEJobRunner( app )
else:
log.error( "Unable to start unknown job runner: %s" %runner_name )
@@ -526,6 +498,16 @@ class DefaultJobDispatcher( object ):
log.debug( "dispatching job %d to %s runner" %( job_wrapper.job_id, runner_name ) )
self.job_runners[runner_name].put( job_wrapper )
def stop( self, job ):
runner_name = ( job.job_runner_name.split(":", 1) )[0]
log.debug( "stopping job %d in %s runner" %( job.id, runner_name ) )
self.job_runners[runner_name].stop_job( job )
def recover( self, job, job_wrapper ):
runner_name = ( job.job_runner_name.split(":", 1) )[0]
log.debug( "recovering job %d in %s runner" %( job.id, runner_name ) )
self.job_runners[runner_name].recover( job, job_wrapper )
def shutdown( self ):
for runner in self.job_runners.itervalues():
runner.shutdown()
@@ -535,8 +517,9 @@ class JobStopQueue( object ):
A queue for jobs which need to be terminated prematurely.
"""
STOP_SIGNAL = object()
def __init__( self, app ):
def __init__( self, app, dispatcher ):
self.app = app
self.dispatcher = dispatcher
# Keep track of the pid that started the job manager, only it
# has valid threads
@@ -585,28 +568,18 @@ class JobStopQueue( object ):
pass
for job in jobs:
# only handle jobs queued or running
# jobs in a non queued/running/new state do not need to be stopped
if job.state not in [ model.Job.states.QUEUED, model.Job.states.RUNNING, model.Job.states.NEW ]:
return
# job has multiple datasets that aren't parent/child and not all of them are deleted.
if not self.check_if_output_datasets_deleted( job.id ):
# job has multiple datasets that aren't parent/child and not all of them are deleted.
return
self.mark_deleted( job.id )
# job is in JobQueue or FooJobRunner, will be dequeued due to state change above
if job.job_runner_name is None:
# job is in JobQueue or PBSJobRunner, will be dequeued due to state change above
return
elif job.job_runner_name.startswith( 'local://' ):
stop_job = runners.local.stop_job
elif job.job_runner_name.startswith( 'pbs://' ):
stop_job = runners.pbs.stop_job
else:
return
# stop
try:
log.debug( "User deleted running job %s, attempting to stop" % job.id )
stop_job( job )
except Exception, e:
log.error( "Unable to stop job %s (error in runner's stop method): %s" % ( job.id, e ) )
# tell the dispatcher to stop the job
self.dispatcher.stop( job )
def check_if_output_datasets_deleted( self, job_id ):
job = model.Job.get( job_id )
+30 -25
View File
@@ -92,31 +92,36 @@ class LocalJobRunner( object ):
self.queue.put( self.STOP_SIGNAL )
log.info( "local job runner stopped" )
def check_pid( pid ):
try:
os.kill( pid, 0 )
return True
except OSError, e:
if e.errno == errno.ESRCH:
log.debug( "local.check_pid(): PID %d is dead" % pid )
else:
log.warning( "local.check_pid(): Got errno %s when attempting to check PID %d: %s" %( errno.errorcode[e.errno], pid, e.strerror ) )
return False
def stop_job( job ):
pid = int( job.job_runner_external_id )
if not check_pid( pid ):
log.warning( "local.stop_job(): %s: PID %d was already dead or can't be signaled" %job.id )
return
for sig in [ 15, 9 ]:
def check_pid( self, pid ):
try:
os.killpg( pid, sig )
os.kill( pid, 0 )
return True
except OSError, e:
log.warning( "local.stop_job(): %s: Got errno %s when attempting to signal %d to PID %d: %s" % ( job.id, errno.errorcode[e.errno], sig, pid, e.strerror ) )
return # give up
sleep( 2 )
if not check_pid( pid ):
log.debug( "local.stop_job(): %s: PID %d successfully killed with signal %d" %( job.id, pid, sig ) )
if e.errno == errno.ESRCH:
log.debug( "check_pid(): PID %d is dead" % pid )
else:
log.warning( "check_pid(): Got errno %s when attempting to check PID %d: %s" %( errno.errorcode[e.errno], pid, e.strerror ) )
return False
def stop_job( self, job ):
pid = int( job.job_runner_external_id )
if not self.check_pid( pid ):
log.warning( "stop_job(): %s: PID %d was already dead or can't be signaled" %job.id )
return
else:
log.warning( "local.stop_job(): %s: PID %d refuses to die after signaling TERM/KILL" %( job.id, pid ) )
for sig in [ 15, 9 ]:
try:
os.killpg( pid, sig )
except OSError, e:
log.warning( "stop_job(): %s: Got errno %s when attempting to signal %d to PID %d: %s" % ( job.id, errno.errorcode[e.errno], sig, pid, e.strerror ) )
return # give up
sleep( 2 )
if not self.check_pid( pid ):
log.debug( "stop_job(): %s: PID %d successfully killed with signal %d" %( job.id, pid, sig ) )
return
else:
log.warning( "stop_job(): %s: PID %d refuses to die after signaling TERM/KILL" %( job.id, pid ) )
def recover( self, job, job_wrapper ):
# local jobs can't be recovered
job_wrapper.change_state( model.Job.states.ERROR, info = "This job was killed when Galaxy was restarted. Please retry the job." )
+32 -11
View File
@@ -391,14 +391,35 @@ class PBSJobRunner( object ):
stage += "%s@%s:%s" % (stage_name, self.app.config.pbs_dataset_server, fname)
return stage
def stop_job( job ):
"""Attempts to delete a job from the PBS queue"""
pbs_server = PBSServer()
pbs_server_name = pbs_server.determine_pbs_server( str( job.job_runner_name ) )
conn = pbs.pbs_connect( pbs_server_name )
if conn <= 0:
log.debug("(%s/%s) Connection to PBS server for job delete failed" % ( job.id, job.job_runner_external_id ) )
return
pbs.pbs_deljob( conn, str( job.job_runner_external_id ), 'NULL' )
pbs.pbs_disconnect( conn )
log.debug( "(%s/%s) Removed from PBS queue at user's request" % ( job.id, job.job_runner_external_id ) )
def stop_job( self, job ):
"""Attempts to delete a job from the PBS queue"""
#pbs_server = PBSServer()
pbs_server_name = self.pbs_server.determine_pbs_server( str( job.job_runner_name ) )
conn = pbs.pbs_connect( pbs_server_name )
if conn <= 0:
log.debug("(%s/%s) Connection to PBS server for job delete failed" % ( job.id, job.job_runner_external_id ) )
return
pbs.pbs_deljob( conn, str( job.job_runner_external_id ), 'NULL' )
pbs.pbs_disconnect( conn )
log.debug( "(%s/%s) Removed from PBS queue at user's request" % ( job.id, job.job_runner_external_id ) )
def recover( self, job, job_wrapper ):
"""Recovers jobs stuck in the queued/running state when Galaxy started"""
pbs_job_state = PBSJobState()
pbs_job_state.ofile = "%s/database/pbs/%s.o" % (os.getcwd(), job.id)
pbs_job_state.efile = "%s/database/pbs/%s.e" % (os.getcwd(), job.id)
pbs_job_state.job_file = "%s/database/pbs/%s.sh" % (os.getcwd(), job.id)
pbs_job_state.job_id = str( job.job_runner_external_id )
pbs_job_state.runner_url = job_wrapper.tool.job_runner
job_wrapper.command_line = job.command_line
pbs_job_state.job_wrapper = job_wrapper
if job.state == model.Job.states.RUNNING:
log.debug( "(%s/%s) is still in running state, adding to the PBS queue" % ( job.id, job.job_runner_external_id ) )
pbs_job_state.old_state = 'R'
pbs_job_state.running = True
self.queue.put( pbs_job_state )
elif job.state == model.Job.states.QUEUED:
log.debug( "(%s/%s) is still in PBS queued state, adding to the PBS queue" % ( job.id, job.job_runner_external_id ) )
pbs_job_state.old_state = 'Q'
pbs_job_state.running = False
self.queue.put( pbs_job_state )
+306
View File
@@ -0,0 +1,306 @@
import os, logging, threading, time
from Queue import Queue, Empty
from galaxy import model
from paste.deploy.converters import asbool
import pkg_resources
try:
pkg_resources.require( "DRMAA_python" )
DRMAA = __import__( "DRMAA" )
except:
DRMAA = None
log = logging.getLogger( __name__ )
DRMAA_state = {
DRMAA.Session.UNDETERMINED: 'process status cannot be determined',
DRMAA.Session.QUEUED_ACTIVE: 'job is queued and waiting to be scheduled',
DRMAA.Session.SYSTEM_ON_HOLD: 'job is queued and in system hold',
DRMAA.Session.USER_ON_HOLD: 'job is queued and in user hold',
DRMAA.Session.USER_SYSTEM_ON_HOLD: 'job is queued and in user and system hold',
DRMAA.Session.RUNNING: 'job is running',
DRMAA.Session.SYSTEM_SUSPENDED: 'job is system suspended',
DRMAA.Session.USER_SUSPENDED: 'job is user suspended',
DRMAA.Session.DONE: 'job finished normally',
DRMAA.Session.FAILED: 'job finished, but failed',
}
sge_template = """#!/bin/sh
#$ -S /bin/sh
GALAXY_LIB="%s"
if [ "$GALAXY_LIB" != "None" ]; then
if [ -n "$PYTHONPATH" ]; then
export PYTHONPATH="$GALAXY_LIB:$PYTHONPATH"
else
export PYTHONPATH="$GALAXY_LIB"
fi
fi
cd %s
%s
"""
class SGEJobState( object ):
def __init__( self ):
"""
Encapsulates state related to a job that is being run via SGE and
that we need to monitor.
"""
self.job_wrapper = None
self.job_id = None
self.old_state = None
self.running = False
self.job_file = None
self.ofile = None
self.efile = None
self.runner_url = None
class SGEJobRunner( object ):
"""
Job runner backed by a finite pool of worker threads. FIFO scheduling
"""
STOP_SIGNAL = object()
def __init__( self, app ):
"""Initialize this job runner and start the monitor thread"""
# Check if SGE was importable, fail if not
if DRMAA is None:
raise Exception( "SGEJobRunner requires DRMAA_python which was not found" )
self.app = app
# '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
# be modified by the monitor thread, which will move items from 'queue'
# to 'watched' and then manage the watched jobs.
self.watched = []
self.queue = Queue()
self.default_cell = self.determine_sge_cell( self.app.config.default_cluster_job_runner )
self.ds = DRMAA.Session()
self.ds.init( self.default_cell )
self.monitor_thread = threading.Thread( target=self.monitor )
self.monitor_thread.start()
log.debug( "ready" )
def determine_sge_cell( self, url ):
"""Determine what SGE cell we are using"""
url_split = url.split("/")
if url_split[0] == 'sge:':
return url_split[2]
# this could happen if sge is started, but is not the default runner
else:
return ''
def determine_sge_queue( self, url ):
"""Determine what SGE queue we are submitting to"""
url_split = url.split("/")
queue = url_split[3]
if queue == "":
# None == server's default queue
queue = None
return queue
def queue_job( self, job_wrapper ):
"""Create SGE script for a job and submit it to the SGE queue"""
job_wrapper.prepare()
command_line = job_wrapper.get_command_line()
runner_url = job_wrapper.tool.job_runner
# This is silly, why would we queue a job with no command line?
if not command_line:
job_wrapper.finish( '', '' )
return
# Change to queued state immediately
job_wrapper.change_state( 'queued' )
if self.determine_sge_cell( runner_url ) != self.default_cell:
# TODO: support multiple cells
log.warning( "(%s) Using multiple SGE cells is not supported. This job will be submitted to the default cell." % job_wrapper.job_id )
sge_queue_name = self.determine_sge_queue( runner_url )
# define job attributes
ofile = "%s/database/pbs/%s.o" % (os.getcwd(), job_wrapper.job_id)
efile = "%s/database/pbs/%s.e" % (os.getcwd(), job_wrapper.job_id)
jt = self.ds.createJobTemplate()
jt.remoteCommand = "%s/database/pbs/galaxy_%s.sh" % (os.getcwd(), job_wrapper.job_id)
jt.outputPath = ":%s" % ofile
jt.errorPath = ":%s" % efile
if sge_queue_name is not None:
jt.setNativeSpecification( "-q %s" % sge_queue_name )
script = sge_template % (job_wrapper.galaxy_lib_dir, os.getcwd(), command_line)
fh = file( jt.remoteCommand, "w" )
fh.write( script )
fh.close()
os.chmod( jt.remoteCommand, 0750 )
# job was deleted while we were preparing it
if job_wrapper.get_state() == 'deleted':
log.debug( "Job %s deleted by user before it entered the SGE queue" % job_wrapper.job_id )
self.cleanup( ( ofile, efile, jt.remoteCommand ) )
job_wrapper.cleanup()
return
galaxy_job_id = job_wrapper.job_id
log.debug("(%s) submitting file %s" % ( galaxy_job_id, jt.remoteCommand ) )
log.debug("(%s) command is: %s" % ( galaxy_job_id, command_line ) )
# runJob will raise if there's a submit problem
job_id = self.ds.runJob(jt)
if sge_queue_name is None:
log.debug("(%s) queued in default queue as %s" % (galaxy_job_id, job_id) )
else:
log.debug("(%s) queued in %s queue as %s" % (galaxy_job_id, sge_queue_name, job_id) )
# store runner information for tracking if Galaxy restarts
job_wrapper.set_runner( runner_url, job_id )
# Store SGE related state information for job
sge_job_state = SGEJobState()
sge_job_state.job_wrapper = job_wrapper
sge_job_state.job_id = job_id
sge_job_state.ofile = ofile
sge_job_state.efile = efile
sge_job_state.job_file = jt.remoteCommand
sge_job_state.old_state = 'new'
sge_job_state.running = False
sge_job_state.runner_url = runner_url
# delete the job template
self.ds.deleteJobTemplate( jt )
# Add to our 'queue' of jobs to monitor
self.queue.put( sge_job_state )
def monitor( self ):
"""
Watches jobs currently in the PBS queue and deals with state changes
(queued to running) and job completion
"""
while 1:
# Take any new watched jobs and put them on the monitor list
try:
while 1:
sge_job_state = self.queue.get_nowait()
if sge_job_state is self.STOP_SIGNAL:
# TODO: This is where any cleanup would occur
self.ds.exit()
return
self.watched.append( sge_job_state )
except Empty:
pass
# Iterate over the list of watched jobs and check state
self.check_watched_items()
# Sleep a bit before the next state check
time.sleep( 1 )
def check_watched_items( self ):
"""
Called by the monitor thread to look at each watched job and deal
with state changes.
"""
new_watched = []
for sge_job_state in self.watched:
job_id = sge_job_state.job_id
galaxy_job_id = sge_job_state.job_wrapper.job_id
old_state = sge_job_state.old_state
try:
state = self.ds.getJobProgramStatus( job_id )
except DRMAA.InvalidJobError:
# we should only get here if an orphaned job was put into the queue at app startup
log.debug("(%s/%s) job left SGE queue" % ( galaxy_job_id, job_id ) )
self.finish_job( sge_job_state )
continue
except Exception, e:
# so we don't kill the monitor thread
log.exception("(%s/%s) Unable to check job status" % ( galaxy_job_id, job_id ) )
log.warning("(%s/%s) job will now be errored" % ( galaxy_job_id, job_id ) )
sge_job_state.job_wrapper.fail( "Cluster could not complete job" )
continue
if state != old_state:
log.debug("(%s/%s) state change: %s" % ( galaxy_job_id, job_id, DRMAA_state[state] ) )
if state == DRMAA.Session.RUNNING and not sge_job_state.running:
sge_job_state.running = True
sge_job_state.job_wrapper.change_state( "running" )
if state == DRMAA.Session.DONE:
self.finish_job( sge_job_state )
continue
if state == DRMAA.Session.FAILED:
sge_job_state.job_wrapper.fail( "Cluster could not complete job" )
sge_job_state.job_wrapper.cleanup()
continue
sge_job_state.old_state = state
new_watched.append( sge_job_state )
# Replace the watch list with the updated version
self.watched = new_watched
def finish_job( self, sge_job_state ):
"""
Get the output/error for a finished job, pass to `job_wrapper.finish`
and cleanup all the SGE temporary files.
"""
ofile = sge_job_state.ofile
efile = sge_job_state.efile
job_file = sge_job_state.job_file
# collect the output
try:
ofh = file(ofile, "r")
efh = file(efile, "r")
stdout = ofh.read()
stderr = efh.read()
except:
stdout = ''
stderr = 'Job output not returned from cluster'
log.debug(stderr)
try:
sge_job_state.job_wrapper.finish( stdout, stderr )
except:
log.exception("Job wrapper finish method failed")
# clean up the sge files
self.cleanup( ( ofile, efile, job_file ) )
def cleanup( self, files ):
if not asbool( self.app.config.get( 'debug', False ) ):
for file in files:
if os.access( file, os.R_OK ):
os.unlink( file )
def put( self, job_wrapper ):
"""Add a job to the queue (by job identifier)"""
self.queue_job( job_wrapper )
def shutdown( self ):
"""Attempts to gracefully shut down the monitor thread"""
log.info( "sending stop signal to worker threads" )
self.queue.put( self.STOP_SIGNAL )
log.info( "sge job runner stopped" )
def stop_job( self, job ):
"""Attempts to delete a job from the SGE queue"""
try:
self.ds.control( job.job_runner_external_id, DRMAA.Session.TERMINATE )
log.debug( "(%s/%s) Removed from SGE queue at user's request" % ( job.id, job.job_runner_external_id ) )
except DRMAA.InvalidJobError:
log.debug( "(%s/%s) User killed running job, but it was already dead" ( job.id, job.job_runner_external_id ) )
def recover( self, job, job_wrapper ):
"""Recovers jobs stuck in the queued/running state when Galaxy started"""
sge_job_state = SGEJobState()
sge_job_state.ofile = "%s/database/pbs/%s.o" % (os.getcwd(), job.id)
sge_job_state.efile = "%s/database/pbs/%s.e" % (os.getcwd(), job.id)
sge_job_state.job_file = "%s/database/pbs/galaxy_%s.sh" % (os.getcwd(), job.id)
sge_job_state.job_id = str( job.job_runner_external_id )
sge_job_state.runner_url = job_wrapper.tool.job_runner
job_wrapper.command_line = job.command_line
sge_job_state.job_wrapper = job_wrapper
if job.state == model.Job.states.RUNNING:
log.debug( "(%s/%s) is still in running state, adding to the SGE queue" % ( job.id, job.job_runner_external_id ) )
sge_job_state.old_state = DRMAA.Session.RUNNING
sge_job_state.running = True
self.queue.put( sge_job_state )
elif job.state == model.Job.states.QUEUED:
log.debug( "(%s/%s) is still in SGE queued state, adding to the SGE queue" % ( job.id, job.job_runner_external_id ) )
sge_job_state.old_state = DRMAA.Session.QUEUED
sge_job_state.running = False
self.queue.put( sge_job_state )