diff --git a/eggs.ini b/eggs.ini index 05225c80abe..d5f5a0748bd 100644 --- a/eggs.ini +++ b/eggs.ini @@ -46,6 +46,7 @@ PasteDeploy = 1.3.3 PasteScript = 1.7.3 pexpect = 2.4 python_openid = 2.2.5 +python_daemon = 1.5.5 Routes = 1.12.3 SQLAlchemy = 0.5.6 sqlalchemy_migrate = 0.5.4 diff --git a/lib/galaxy/jobs/transfer_manager.py b/lib/galaxy/jobs/transfer_manager.py index cd3f081677c..cd214478cd0 100644 --- a/lib/galaxy/jobs/transfer_manager.py +++ b/lib/galaxy/jobs/transfer_manager.py @@ -2,10 +2,12 @@ Client interface to the Galaxy Transfer Manager, which is a standalone, lightweight on-demand daemon. """ -import subprocess, socket +import subprocess, socket, logging from galaxy import eggs -from galaxy.util import listify +from galaxy.util import listify, json + +log = logging.getLogger( __name__ ) class TransferManager( object ): def __init__( self, app ): @@ -31,31 +33,47 @@ class TransferManager( object ): running daemon, so it should be fairly quick to return. """ spaced_flag = ' %s ' % self.tm_transfer_job_id_flag - cmd = '%s %s %s' % ( self.tm_command, self.tm_transfer_job_id_flag, spaced_flag.join( [ tj.id for tj in transfer_jobs ] ) ) - p = subprocess.Popen( cmd, shell=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE ) + cmd = '%s %s %s' % ( self.tm_command, self.tm_transfer_job_id_flag, spaced_flag.join( [ str( tj.id ) for tj in transfer_jobs ] ) ) + log.debug( 'Initiating Transfer Job(s): %s' % ', '.join( [ str( tj.id ) for tj in transfer_jobs ] ) ) + p = subprocess.Popen( cmd, shell=True, stdout=subprocess.PIPE, stderr=subprocess.STDOUT ) p.wait() - return p.stderr.read() + return p.stdout.read() def status( self, transfer_jobs ): - transfer_jobs = util.listify( transfer_jobs ) - # TODO: timeout, error handling + transfer_jobs = listify( transfer_jobs ) rval = [] sock = socket.socket( socket.AF_INET, socket.SOCK_STREAM ) + sock.settimeout( 10 ) try: sock.connect( ( 'localhost', self.tm_port ) ) - except: - log.exception( 'sock.connect for status update failed:' ) - for transfer_job in transfer_jobs: - self.sa_session.refresh( transfer_job ) - if transfer_job.state in [ self.app.model.TransferJob.states.DONE, \ - self.app.model.TransferJob.states.ERROR ]: - rval.append( dict( status=transfer_job.state ) ) - else: - raise Exception( 'transfer manager not running or responding and transfer job (id: %s) state is non-terminal: %s' % ( transfer_job.id, transfer_job.state ) ) - sock.send( json.to_json_string( dict( status=[ t.id for t in transfer_jobs ] ) ) + '\n' ) + except Exception, e: + log.warning( 'sock.connect for status update of Transfer Jobs %s failed (this is okay if all jobs have finished): %s' % ( ', '.join( [ str( tj.id ) for tj in transfer_jobs ] ), str( e ) ) ) + [ self.sa_session.refresh( tj ) for tj in transfer_jobs ] + new_jobs = filter( lambda x: x.state == self.app.model.TransferJob.states.NEW, transfer_jobs ) + #terminal_jobs = filter( lambda x: x.state in [ self.app.model.TransferJob.states.DONE, \ + # self.app.model.TransferJob.states.ERROR ], transfer_jobs ) + if new_jobs: + # This could be a bad idea if the transfer manager daemon is misbehaving. + output = self.run( new_jobs ) + for tj in transfer_jobs: + if tj.state == tj.states.DONE: + log.debug( 'Transfer Job %s is complete' % tj.id ) + rval.append( dict( transfer_job_id=tj.id, state=tj.state ) ) + if len( rval ) == 1: + return rval[0] + return rval + sock.send( json.to_json_string( dict( state_transfer_job_ids=[ t.id for t in transfer_jobs ] ) ) + '\n' ) resp = sock.recv( 8192 ) for line in resp.splitlines(): status = json.from_json_string( line ) - rval.append( status ) + # TODO: need a bunch for this + if status['state'] == 'unknown': + transfer_job = [ tj for tj in transfer_jobs if tj.id == int( status['transfer_job_id'] ) ][0] + self.sa_session.refresh( tj ) + rval.append( dict( transfer_job_id=tj.id, state=tj.state ) ) + else: + if status['state'] == 'progress' and 'percent' in status: + log.debug( 'Transfer Job %s is %s complete' % ( status['transfer_job_id'], status['percent'] ) ) + rval.append( status ) if len( rval ) == 1: return rval[0] return rval diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 6bead604d83..84b7f93304a 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -2160,9 +2160,10 @@ class TransferJob( object ): RUNNING = 'running', ERROR = 'error', DONE = 'done' ) - def __init__( self, state=None, path=None, params=None ): + def __init__( self, state=None, path=None, info=None, params=None ): self.state = state self.path = path + self.info = info self.params = params class Tag ( object ): diff --git a/lib/galaxy/model/mapping.py b/lib/galaxy/model/mapping.py index 05c008d6e7d..bc3c67f48d6 100644 --- a/lib/galaxy/model/mapping.py +++ b/lib/galaxy/model/mapping.py @@ -444,6 +444,7 @@ TransferJob.table = Table( "transfer_job", metadata, Column( "update_time", DateTime, default=now, onupdate=now ), Column( "state", String( 64 ), index=True ), Column( "path", String( 1024 ) ), + Column( "info", TEXT ), Column( "params", JSONType ) ) Event.table = Table( "event", metadata, diff --git a/lib/galaxy/model/migrate/versions/0070_add_info_column_to_deferred_job_table.py b/lib/galaxy/model/migrate/versions/0070_add_info_column_to_deferred_job_table.py new file mode 100644 index 00000000000..8e121394269 --- /dev/null +++ b/lib/galaxy/model/migrate/versions/0070_add_info_column_to_deferred_job_table.py @@ -0,0 +1,35 @@ +""" +Migration script to add 'info' column to the transfer_job table. +""" + +from sqlalchemy import * +from sqlalchemy.orm import * +from migrate import * +from migrate.changeset import * + +import logging +log = logging.getLogger( __name__ ) + +metadata = MetaData( migrate_engine ) +db_session = scoped_session( sessionmaker( bind=migrate_engine, autoflush=False, autocommit=True ) ) + +def upgrade(): + print __doc__ + metadata.reflect() + try: + TransferJob_table = Table( "transfer_job", metadata, autoload=True ) + c = Column( "info", TEXT ) + c.create( TransferJob_table ) + assert c is TransferJob_table.c.info + except Exception, e: + print "Adding info column to transfer_job table failed: %s" % str( e ) + log.debug( "Adding info column to transfer_job table failed: %s" % str( e ) ) + +def downgrade(): + metadata.reflect() + try: + TransferJob_table = Table( "transfer_job", metadata, autoload=True ) + TransferJob_table.c.info.drop() + except Exception, e: + print "Dropping info column from transfer_job table failed: %s" % str( e ) + log.debug( "Dropping info column from transfer_job table failed: %s" % str( e ) ) diff --git a/lib/galaxy/sample_tracking/external_service_types.py b/lib/galaxy/sample_tracking/external_service_types.py index 8d8a5477004..939da3f120f 100644 --- a/lib/galaxy/sample_tracking/external_service_types.py +++ b/lib/galaxy/sample_tracking/external_service_types.py @@ -98,10 +98,12 @@ class ExternalServiceType( object ): results_elem = run_details_elem.find( 'results' ) if results_elem: # get the list of resulting datatypes - self.run_details[ 'results' ] = self.parse_run_details_results( results_elem ) + self.run_details[ 'results' ], self.run_details[ 'results_urls' ] = self.parse_run_details_results( results_elem ) def parse_run_details_results( self, root ): datatypes_dict = {} + urls_dict = {} for datatype_elem in root.findall( "dataset" ): name = datatype_elem.get( 'name' ) datatypes_dict[ name ] = datatype_elem.get( 'datatype' ) - return datatypes_dict + urls_dict[ name ] = datatype_elem.get( 'url', None ) + return datatypes_dict, urls_dict diff --git a/lib/galaxy/web/api/samples.py b/lib/galaxy/web/api/samples.py index 2872c18f7bc..4f5374d91a7 100644 --- a/lib/galaxy/web/api/samples.py +++ b/lib/galaxy/web/api/samples.py @@ -69,12 +69,18 @@ class SamplesAPIController( BaseController ): sample = None if not sample: trans.response.status = 400 - return "Invalid request id ( %s ) specified." % str( request_id ) + return "Invalid sample id ( %s ) specified." % str( sample_id ) if not trans.user_is_admin(): trans.response.status = 403 return "You are not authorized to update samples." requests_admin_controller = trans.webapp.controllers[ 'requests_admin' ] if update_type == 'run_details': + deferred_plugin = payload.pop( 'deferred_plugin', None ) + if deferred_plugin: + try: + trans.app.job_manager.deferred_job_queue.plugins[deferred_plugin].create_job( trans, sample=sample, **payload ) + except: + log.exception( 'update() called with a deferred job plugin (%s) but creating the deferred job failed:' % deferred_plugin ) status, output = requests_admin_controller.edit_template_info( trans, cntrller='api', item_type='sample', diff --git a/transfer_manager.py b/transfer_manager.py index 090d2b83722..5787e5f07e2 100644 --- a/transfer_manager.py +++ b/transfer_manager.py @@ -1,38 +1,6 @@ #!/usr/bin/env python ''' -Manage transfers from sequencers - -Design: - - transfer_manager.py configfile - -Initialization - - try: - Open socket - except: - try: - Submit to socket - except: - Instruct caller to resubmit (avoids other manager shutting down race condition) - -Run - - Needs 2 threads? One to monitor socket, one to monitor queue. - Watch the race condition of put into transfer queue after transfer queue is exiting? - Store queue to db to restart/resume on premature manager termination? - - while len(queue) > 0: - msg = socket.read() - if msg == STATUS_REQUEST: - return state - elif msg == NEW_DOWNLOAD: - queue.put( msg ) - - while queue: - spawn transfer - read lines from transfer for progress - +Downloads files to temp locations. ''' import os, sys, optparse, ConfigParser, socket, SocketServer, errno, Queue, threading, subprocess @@ -46,6 +14,9 @@ from galaxy import eggs import galaxy.model.mapping from galaxy.util import json, bunch +eggs.require( 'python_daemon' ) +from daemon import DaemonContext + class ArgHandler( object ): """ Collect command line flags. @@ -56,10 +27,10 @@ class ArgHandler( object ): self.parser.add_option( '-i', '--transfer-job-id', action='append', dest='transfer_job_ids', help='Initiate management of the specified TransferJob id' ) self.parser.add_option( '--do', dest='initiate_transfer_job_id', help='Used by this script when it calls itself to actually initiate the download' ) self.parser.add_option( '-s', '--state-transfer-job-id', action='append', dest='state_transfer_job_ids', help='Report the state of the specified TransferJob id' ) - def parse(): - opts, args = parser.parse_args() - if opts.initiate_transfer_job_id is not None: - opts.initiate_transfer_job_id + self.parser.add_option( '-d', '--debug', action='store_true', dest='debug', help="Debug (don't detach)" ) + self.opts = None + def parse( self ): + self.opts, args = self.parser.parse_args() class GalaxyApp( object ): """ @@ -69,7 +40,8 @@ class GalaxyApp( object ): def __init__( self, config_file='universe_wsgi.ini' ): self.config = ConfigParser.ConfigParser( dict( database_file = 'database/universe.sqlite', file_path = 'database/files', - transfer_manager_port = 8163 ) ) + transfer_manager_port = '8163', + transfer_manager_log = 'transfer_manager.log' ) ) self.config.read( config_file ) self.model = None @property @@ -101,6 +73,7 @@ class ListenerRequestHandler( SocketServer.BaseRequestHandler ): """ def handle( self ): if not self.server.transfer_manager.accepting: + # TODO: does the submitter handle this condition? i'm sure it doesn't... self.request.send( 'Manager shutting down\n' ) return data = '' @@ -118,11 +91,16 @@ class ListenerRequestHandler( SocketServer.BaseRequestHandler ): data = json.from_json_string( data ) if 'transfer_job_ids' in data: # Get all of the TransferJob objects and stick them on the queue. + print 'Adding transfer job ids to transfer queue: %s' % data['transfer_job_ids'] [ self.server.transfer_manager.transfer_queue.put( self.server.app.get_transfer_job( transfer_job_id ) ) for transfer_job_id in data['transfer_job_ids'] ] - elif 'state' in data: - for transfer_id in data['state']: - self.request.send( '%s %s\n' % ( transfer_id, self.server.transfer_manager.get_state( int( transfer_id ) ) ) ) - self.request.send( 'OK\n' ) + self.request.send( "Added jobs to transfer queue: %s\n" % ', '.join( data['transfer_job_ids'] ) ) + elif 'state_transfer_job_ids' in data: + print 'Servicing state request for transfer job ids: %s' % data['state_transfer_job_ids'] + for state_transfer_job_id in data['state_transfer_job_ids']: + state = self.server.transfer_manager.get_state( int( state_transfer_job_id ) ) + state['transfer_job_id'] = state_transfer_job_id + print 'State of transfer job id %s is: %s' % ( state_transfer_job_id, state ) + self.request.send( json.to_json_string( state ) ) class Transfer( object ): """ @@ -133,22 +111,24 @@ class Transfer( object ): UNKNOWN = 'unknown', STARTED = 'started', PROGRESS = 'progress', - DONE = 'done' ) - def __init__( self, transfer ): - self.transfer = transfer + DONE = 'done', + ERROR = 'error' ) + def __init__( self, transfer_job ): + self.transfer_job = transfer_job self.state = dict( state = self.states.NEW ) self.done = False def run( self ): - cmd = '%s -u %s -d %s' % ( sys.executable, os.path.abspath( __file__ ), self.transfer.id ) - self.p = subprocess.Popen( cmd, bufsize=0, shell=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE ) + cmd = '%s -u %s --do %s' % ( sys.executable, os.path.abspath( __file__ ), self.transfer_job.id ) + print cmd + self.p = subprocess.Popen( cmd, bufsize=0, shell=True, stdout=subprocess.PIPE, stderr=subprocess.STDOUT ) line = self.p.stdout.readline() while line: try: self.state = json.from_json_string( line ) assert 'state' in self.state except: - print 'Received unknown state from transfer (id: %s): %s' % ( transfer.id, line ) - self.state = dict( state=self.states.UNKNOWN + print 'Received unknown state from transfer (transfer job id: %s): %s' % ( self.transfer_job.id, line ) + self.state = dict( state = self.states.UNKNOWN, info=line ) line = self.p.stdout.readline() self.p.wait() self.done = True @@ -158,11 +138,14 @@ class TransferManager( object ): Manage the queue of transfers, handle the setup of new transfers and completion of finished transfers. """ - def __init__( self, app, transfer_job_ids ): + def __init__( self, app, transfer_job_ids, debug=False ): self.app = app + self.daemon_context = None self.sa_session = app.sa_session self.transfer_job_ids = transfer_job_ids - self.port = app.config.get( 'app:main', 'transfer_manager_port' ) + self.debug = debug + self.port = int( app.config.get( 'app:main', 'transfer_manager_port' ) ) + self.log_file_name = app.config.get( 'app:main', 'transfer_manager_log' ) self.transfer_queue = Queue.Queue() self.accepting = True self.watchlist = [] @@ -180,19 +163,34 @@ class TransferManager( object ): if e[0] == errno.EADDRINUSE: self.submit() sys.exit() + + # Daemonize + if not self.debug: + log_file = open( self.log_file_name, 'a+' ) + self.daemon_context = DaemonContext( files_preserve=[ self.listener_server.fileno() ], + working_directory=os.getcwd(), + stdout=log_file, + stderr=log_file ) + self.daemon_context.open() + + # Listen for more stuff on the socket self.listener = threading.Thread( target=self.listener_server.serve_forever ) self.listener.start() # Put all the URLs into the queue, additional URLs received before this # instance of the manager terminates will also be handled. + print 'Started up with transfer job ids: %s' % self.transfer_job_ids [ self.transfer_queue.put( app.get_transfer_job( tj_id ) ) for tj_id in self.transfer_job_ids ] while True: # TODO: max transfer limit try: transfer_job = self.transfer_queue.get_nowait() - print 'fetching transfer job:', transfer_job.id, 'url is:', transfer_job.params['url'] + print 'Fetching transfer job:', transfer_job.id, 'URL is:', transfer_job.params['url'] t = Transfer( transfer_job ) + transfer_job.state = app.model.TransferJob.states.RUNNING + self.sa_session.add( transfer_job ) + self.sa_session.flush() tt = threading.Thread( target=t.run ) tt.start() self.watchlist.append( t ) @@ -211,11 +209,16 @@ class TransferManager( object ): # TODO: handle failure if transfer.done: if transfer.state['state'] == transfer.states.DONE: - # TODO: make a bunch for real TransferJob object states. transfer.transfer_job.state = app.model.TransferJob.states.DONE transfer.transfer_job.path = transfer.state['path'] + print 'Transfer of job %s ended successfully, output path is: %s' % ( transfer.transfer_job.id, transfer.transfer_job.path ) + elif transfer.state['state'] == transfer.states.ERROR: + transfer.transfer_job.info = transfer.state.get( 'info', None ) + transfer.transfer_job.state = app.model.TransferJob.states.ERROR + print 'Transfer of job %s ended in error: %s' % ( transfer.transfer_job.id, transfer.transfer_job.info ) else: - # TODO: transfer.transfer.info = output or transfer.transfer.params['info'] = output + print 'Unknown state received for transfer job %s: %s' % ( transfer.transfer_job.id, transfer.state['state'] ) + transfer.transfer_job.info = 'Unknown error encountered in transfer manager' transfer.transfer_job.state = app.model.TransferJob.states.ERROR self.sa_session.add( transfer.transfer_job ) self.sa_session.flush() @@ -228,19 +231,24 @@ class TransferManager( object ): self.listener_server.shutdown() def get_state( self, transfer_job_id ): - rval = None + rval = {} for transfer in self.watchlist: - if transfer.transfer_job.id == transfer_id_job: + if transfer.transfer_job.id == transfer_job_id: rval = transfer.state break + else: + rval['state'] = Transfer.states.UNKNOWN # should be DONE? return rval def submit( self ): # TODO: can fail if shutdown occurs between failure to bind and submission + # Needs error handling. + # This may not work at all right now. + print "Submitting jobs to running transfer manager: %s" % ', '.join( self.transfer_job_ids ) sock = socket.socket( socket.AF_INET, socket.SOCK_STREAM ) - sock.connect( ( 'localhost', self.options.port ) ) - sock.send( json.to_json_string( dict( transfer_job_ids=self.transfer_job_ids ] ) ) + '\n' ) - print sock.recv( 1024 ), + sock.connect( ( 'localhost', self.port ) ) + sock.send( json.to_json_string( dict( transfer_job_ids=self.transfer_job_ids ) ) + '\n' ) + print sock.recv( 8192 ), sock.close() sys.exit() @@ -261,11 +269,16 @@ def do_http_transfer( transfer_job ): "Plugin" for handling http(s) transfers. """ url = transfer_job.params['url'] - f = urllib2.urlopen( url ) + try: + f = urllib2.urlopen( url ) + except urllib2.URLError, e: + print json.to_json_string( dict( state=Transfer.states.ERROR, + info=str( e ) ) ) + return size = f.info().getheader( 'Content-Length' ) if size is not None: size = int( size ) - chunksize = 1024 + chunksize = 1024 * 1024 read = 0 last = 0 fh, fn = tempfile.mkstemp() @@ -276,39 +289,42 @@ def do_http_transfer( transfer_job ): os.write( fh, chunk ) if read == 0 and size is None: print json.to_json_string( dict( state=Transfer.states.STARTED, - size=None ) ) + '\n' + size=None ) ) #+ '\n' elif read == 0: print json.to_json_string( dict( state=Transfer.states.STARTED, - size=size ) ) + '\n' + size=size ) ) #+ '\n' read += chunksize if size is not None and read < size: percent = int( float( read ) / size * 100 ) if percent != last: print json.to_json_string( dict( state=Transfer.states.PROGRESS, - percent='%s%%' % percent ) ) + '\n' + read=read, + percent='%s%%' % percent ) ) #+ '\n' last = percent - time.sleep( 1 ) + elif size is None: + print json.to_json_string( dict( state=Transfer.states.PROGRESS, + read=read ) ) os.close( fh ) - print json.to_json_string( dict( state=Transfer.states.DONE, path=fn ) ) + '\n' + print json.to_json_string( dict( state=Transfer.states.DONE, path=fn ) ) #+ '\n' -def request_state( transfer_job_ids, port ): +def request_state( state_transfer_job_ids, port ): sock = socket.socket( socket.AF_INET, socket.SOCK_STREAM ) sock.connect( ( 'localhost', port ) ) - sock.send( json.to_json_string( dict( state=transfer_job_ids ) ) + '\n' ) + sock.send( json.to_json_string( dict( state_transfer_job_ids=state_transfer_job_ids ) ) + '\n' ) print sock.recv( 1024 ), sock.close() sys.exit() if __name__ == '__main__': arg_handler = ArgHandler() - app = GalaxyApp() - opts, args = arg_handler.parse() - if opts.initiate_transfer_job_id is not None: - do_transfer( app.get_transfer_job( opts.initiate_transfer_job_id ) ) - elif opts.state_transfer_job_ids: - request_state( opts.state_transfer_job_ids, app.config.get( 'app:main', 'transfer_manager_port' ) ) - elif opts.transfer_job_ids: - transfer_manager = TransferManager( app, opts.transfer_job_ids ) + arg_handler.parse() + app = GalaxyApp( config_file=arg_handler.opts.config ) + if arg_handler.opts.initiate_transfer_job_id is not None: + do_transfer( app.get_transfer_job( arg_handler.opts.initiate_transfer_job_id ) ) + elif arg_handler.opts.state_transfer_job_ids: + request_state( arg_handler.opts.state_transfer_job_ids, int( app.config.get( 'app:main', 'transfer_manager_port' ) ) ) + elif arg_handler.opts.transfer_job_ids: + transfer_manager = TransferManager( app, arg_handler.opts.transfer_job_ids, arg_handler.opts.debug ) else: arg_handler.parser.print_usage( sys.stderr ) sys.exit( 1 )