diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index a8b762f00dc..1a969cae32d 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -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: diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index ff8296f5c31..db672af32b3 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -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 ) diff --git a/lib/galaxy/jobs/runners/local.py b/lib/galaxy/jobs/runners/local.py index 5fc7dc01dd3..7ec934752a8 100644 --- a/lib/galaxy/jobs/runners/local.py +++ b/lib/galaxy/jobs/runners/local.py @@ -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." ) + diff --git a/lib/galaxy/jobs/runners/pbs.py b/lib/galaxy/jobs/runners/pbs.py index 20d3fc8493a..2d16df94d76 100644 --- a/lib/galaxy/jobs/runners/pbs.py +++ b/lib/galaxy/jobs/runners/pbs.py @@ -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 ) diff --git a/lib/galaxy/jobs/runners/sge.py b/lib/galaxy/jobs/runners/sge.py new file mode 100644 index 00000000000..9273c438bd6 --- /dev/null +++ b/lib/galaxy/jobs/runners/sge.py @@ -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 )