From bc60e6a7ff7a270fc3068db756ae97eac27c908b Mon Sep 17 00:00:00 2001 From: James Taylor Date: Thu, 14 Dec 2006 21:59:44 +0000 Subject: [PATCH] Fix the bugs I introduced into Nate's code while splitting up job queue (PBS seems to work now in trunk). --- lib/galaxy/config.py | 1 + lib/galaxy/jobs/__init__.py | 4 +- lib/galaxy/jobs/runners/pbs.py | 69 ++++++++++++++++++++-------------- 3 files changed, 43 insertions(+), 31 deletions(-) diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index fb26f193460..c0cfe07699b 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -33,6 +33,7 @@ class Configuration( object ): self.sendmail_path = kwargs.get('sendmail_path',"/usr/sbin/sendmail") self.mailing_join_addr = kwargs.get('mailing_join_addr',"galaxy-user-join@bx.psu.edu") self.use_pbs = kwargs.get('use_pbs', False ) + self.pbs_server = kwargs.get('pbs_server', "" ) self.use_heartbeat = kwargs.get( 'use_heartbeat', False ) def get( self, key, default ): return self.config_dict.get( key, default ) diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index c73deae788e..e2f3354300a 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -179,7 +179,7 @@ class DefaultJobDispatcher( object ): self.local_job_runner.put( job_wrapper ) def dispatch_pbs( self, job_wrapper ): - if "/tools/data_source" not in job_wrapper.get_command_line(): + if "/tools/data_source" in job_wrapper.get_command_line(): self.local_job_runner.put( job_wrapper ) else: self.pbs_job_runner.put( job_wrapper ) @@ -189,4 +189,4 @@ class DefaultJobDispatcher( object ): if self.use_pbs: self.pbs_job_runner.shutdown() - \ No newline at end of file + diff --git a/lib/galaxy/jobs/runners/pbs.py b/lib/galaxy/jobs/runners/pbs.py index 5dae81aaf62..305a3513eaf 100644 --- a/lib/galaxy/jobs/runners/pbs.py +++ b/lib/galaxy/jobs/runners/pbs.py @@ -1,11 +1,15 @@ import logging +import threading +import os +import random from Queue import Queue from galaxy import model import pkg_resources pkg_resources.require( "pbs_python" ) -import pbs, random +pbs = __import__( "pbs" ) + log = logging.getLogger( __name__ ) @@ -21,7 +25,8 @@ class PBSJobState( object ): self.job_wrapper = None self.pbs_job_id = None self.old_state = None - self.running = running + self.running = False + self.job_file = None self.ofile = None self.efile = None @@ -34,17 +39,17 @@ class PBSJobRunner( object ): """Start the job queue with 'nworkers' worker threads""" self.app = app self.queue = Queue() + self.init_pbs() self.monitor_thread = threading.Thread( target=self.monitor ) self.monitor_thread.start() - self.threads.append( worker ) - log.debug( "%d workers ready", nworkers ) + log.debug( "ready" ) def init_pbs( self ): if self.app.config.pbs_server: - pbs_server = self.app.config.pbs_server + self.pbs_server = self.app.config.pbs_server else: - pbs_server = pbs.pbs_default() - if pbs_server is None: + self.pbs_server = pbs.pbs_default() + if self.pbs_server is None: raise Exception( "Could not find torque server" ) def queue_job( self, job_wrapper ): @@ -69,6 +74,10 @@ class PBSJobRunner( object ): bit = random.randint(65,90) job_name += chr(bit) + conn = pbs.pbs_connect( self.pbs_server ) + if conn <= 0: + raise Exception( "Connection to PBS server for submit failed" ) + # set up the job file script = pbs_template % (os.environ['PATH'], os.environ['PYTHONPATH'], os.getcwd(), command_line) job_file = "%s/database/pbs/%s.sh" % (os.getcwd(), job_name) @@ -104,7 +113,7 @@ class PBSJobRunner( object ): # Get initial job state stat_attrl = pbs.new_attrl(1) stat_attrl[0].name = 'job_state' - conn = pbs.pbs_connect(pbs_server) + conn = pbs.pbs_connect( self.pbs_server ) if conn > 0: jobs = pbs.pbs_statjob(conn, job_id, stat_attrl, 'NULL') pbs.pbs_disconnect(conn) @@ -121,10 +130,7 @@ class PBSJobRunner( object ): # ran immediately if old_state == "R": - for data in out_data.values(): - data.state = data.states.RUNNING - data.blurb = "running" - data.flush() + job_wrapper.change_state( "running" ) running = True log.debug("(%s) job is now running" % job_id) @@ -134,10 +140,12 @@ class PBSJobRunner( object ): # Store PBS related state information for job pbs_job_state = PBSJobState() + pbs_job_state.job_wrapper = job_wrapper pbs_job_state.job_name = job_name pbs_job_state.job_id = job_id pbs_job_state.ofile = ofile pbs_job_state.efile = efile + pbs_job_state.job_file = job_file pbs_job_state.old_state = old_state pbs_job_state.running = running @@ -147,7 +155,7 @@ class PBSJobRunner( object ): def monitor( self ): while 1: pbs_job_state = self.queue.get() - if pbs_job_state is STOP_SIGNAL: + if pbs_job_state is self.STOP_SIGNAL: return job_id = pbs_job_state.job_id old_state = pbs_job_state.old_state @@ -159,22 +167,25 @@ class PBSJobRunner( object ): stat_attrl[0].name = 'job_state' jobs = pbs.pbs_statjob(conn, pbs_job_state.job_id, stat_attrl, 'NULL') pbs.pbs_disconnect(conn) - if not jobs: + if len( jobs ) < 1: + log.debug("(%s) job has left queue" % job_id) self.finish_job( pbs_job_state ) - try: - if (jobs[0].attribs[0].name == "job_state"): - state = jobs[0].attribs[0].value - if state != old_state: - log.debug("(%s) job state changed from %s to %s" % ( job_id, old_state, state ) ) - if state == "R" and not running: - running = True - pbs_job_state.job_wrapper.change_state( "running" ) - log.debug("(%s) job is now running" % job_id) - old_state = state - pbs_job_state.old_state = old_state - pbs_job_state.running = running - except: - log.info("(%s) unable to check state" % job_id) + else: + try: + if (jobs[0].attribs[0].name == "job_state"): + state = jobs[0].attribs[0].value + if state != old_state: + log.debug("(%s) job state changed from %s to %s" % ( job_id, old_state, state ) ) + if state == "R" and not running: + running = True + pbs_job_state.job_wrapper.change_state( "running" ) + log.debug("(%s) job is now running" % job_id) + old_state = state + pbs_job_state.old_state = old_state + pbs_job_state.running = running + except: + log.exception("(%s) unable to check state" % job_id) + self.queue.put( pbs_job_state ) def finish_job( self, pbs_job_state ): ofile = pbs_job_state.ofile @@ -209,4 +220,4 @@ class PBSJobRunner( object ): """Attempts to gracefully shut down the worker threads""" log.info( "sending stop signal to worker threads" ) self.queue.put( self.STOP_SIGNAL ) - log.info( "pbs job runner stopped" ) \ No newline at end of file + log.info( "pbs job runner stopped" )