Second pass of the transfer manager and deferred jobs runner. Still fairly alpha.

This commit is contained in:
Nate Coraor
2011-01-11 15:58:17 -05:00
parent a7e168519d
commit 87e58948c4
8 changed files with 178 additions and 98 deletions
+1
View File
@@ -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
+36 -18
View File
@@ -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
+2 -1
View File
@@ -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 ):
+1
View File
@@ -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,
@@ -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 ) )
@@ -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
+7 -1
View File
@@ -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',
+92 -76
View File
@@ -1,38 +1,6 @@
#!/usr/bin/env python
'''
Manage transfers from sequencers
Design:
transfer_manager.py configfile <some identifier>
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 )