sample tracking data transfer mechanism streamlined.

- eliminated the transfer_datasets.ini config file
- eliminated the data transfer user
- now the admin user initiating the data transfer is used as the data transfer user.
- the admin user is provided the add_library_item permission before data transfer if they dont have it
- sample update & data transfer amqp messages now includes api_key
This commit is contained in:
Ramkrishna Chakrabarty
2010-11-07 19:46:59 -05:00
parent bb665f6b5e
commit 3d72267f44
10 changed files with 188 additions and 223 deletions
@@ -19,13 +19,15 @@ xml = \
''' <sample>
<barcode>%(BARCODE)s</barcode>
<state>%(STATE)s</state>
<api_key>%(API_KEY)s</api_key>
</sample>'''
def handle_scan(states, amqp_config, barcode):
if states.get(barcode[:2], None):
values = dict( BARCODE=barcode[2:],
STATE=states.get(barcode[:2]) )
STATE=states.get(barcode[:2]),
API_KEY=amqp_config['api_key'] )
print values
data = xml % values
print data
@@ -68,6 +70,10 @@ def main():
states = {}
for option in config.options("galaxy:amqp"):
amqp_config[option] = config.get("galaxy:amqp", option)
# abort if api_key is not set in the config file
if not amqp_config['api_key']:
print 'Error: Set the api_key config variable in the config file before starting the amqp_publisher script.'
sys.exit( 1 )
count = 1
while True:
section = 'scanner%i' % count
@@ -14,6 +14,7 @@
#queue = galaxy_queue
#exchange = galaxy_exchange
#routing_key = bar_code_scanner
#api_key =
# The following section(s) 'scanner#' is for specifying the state of the
# sample this scanner represents. This state name should be one of the
@@ -15,6 +15,9 @@ import optparse
import xml.dom.minidom
import subprocess
import urllib2
from xml_helper import get_value, get_value_index
from galaxydb_interface import GalaxyDbInterface
api_path = [ os.path.join( os.getcwd(), "scripts/api" ) ]
@@ -44,77 +47,58 @@ log.addHandler(fh)
# data transfer script
data_transfer_script = os.path.join( os.getcwd(),
"scripts/galaxy_messaging/server/data_transfer.py" )
global dbconnstr
global webconfig
global config
global config_file_name
global http_server_section
def get_value(dom, tag_name):
'''
This method extracts the tag value from the xml message
'''
nodelist = dom.getElementsByTagName(tag_name)[0].childNodes
rc = ""
for node in nodelist:
if node.nodeType == node.TEXT_NODE:
rc = rc + node.data
return rc
def get_value_index(dom, tag_name, index):
'''
This method extracts the tag value from the xml message
'''
try:
nodelist = dom.getElementsByTagName(tag_name)[index].childNodes
except:
return None
rc = ""
for node in nodelist:
if node.nodeType == node.TEXT_NODE:
rc = rc + node.data
return rc
def start_data_transfer(message):
def start_data_transfer( message ):
# fork a new process to transfer datasets
cmd = '%s "%s" "%s" "%s"' % ( "python",
data_transfer_script,
message.body,
sys.argv[1] ) # Galaxy config file name
config_file_name ) # Galaxy config file name
pid = subprocess.Popen(cmd, shell=True).pid
log.debug('Started process (%i): %s' % (pid, str(cmd)))
def update_sample_state(message):
def update_sample_state( message ):
dom = xml.dom.minidom.parseString(message.body)
barcode = get_value(dom, 'barcode')
state = get_value(dom, 'state')
api_key = get_value(dom, 'api_key')
log.debug('Barcode: ' + barcode)
log.debug('State: ' + state)
log.debug('api_key: ' + api_key)
# validate
if not barcode or not state or not api_key:
log.debug( 'Incomplete sample_state_update message received. Sample barcode, desired state and user API key is required.' )
return
# update the sample state in Galaxy db
galaxydb = GalaxyDbInterface(dbconnstr)
sample_id = galaxydb.get_sample_id(field_name='bar_code', value=barcode)
dbconnstr = config.get("app:main", "database_connection")
galaxydb = GalaxyDbInterface( dbconnstr )
sample_id = galaxydb.get_sample_id( field_name='bar_code', value=barcode )
if sample_id == -1:
log.debug('Invalid barcode.')
log.debug( 'Invalid barcode.' )
return
galaxydb.change_state(sample_id, state)
# after updating the sample state, update request status
request_id = galaxydb.get_request_id(sample_id)
update_request( request_id )
update_request( api_key, request_id )
def update_request( request_id ):
http_server_section = webconfig.get( "universe_wsgi_config", "http_server_section" )
def update_request( api_key, request_id ):
encoded_request_id = api.encode_id( config.get( "app:main", "id_secret" ), request_id )
api_key = webconfig.get( "data_transfer_user_login_info", "api_key" )
data = dict( update_type=RequestsController.update_types.REQUEST )
url = "http://%s:%s/api/requests/%s" % ( config.get(http_server_section, "host"),
config.get(http_server_section, "port"),
encoded_request_id )
log.debug( 'Updating request %i' % request_id )
try:
api.update( api_key, url, data, return_formatted=False )
retval = api.update( api_key, url, data, return_formatted=False )
log.debug( str( retval ) )
except urllib2.URLError, e:
log.debug( 'ERROR(update_request (%s)): %s' % ( str((self.api_key, url, data)), str(e) ) )
def recv_callback(message):
def recv_callback( message ):
# check the meesage type.
msg_type = message.properties['application_headers'].get('msg_type')
log.debug( 'MESSAGE RECVD: ' + str( msg_type ) )
@@ -126,22 +110,25 @@ def recv_callback(message):
update_sample_state( message )
def main():
if len(sys.argv) < 2:
print 'Usage: python amqp_consumer.py <Galaxy configuration file>'
return
parser = optparse.OptionParser()
parser.add_option('-c', '--config-file', help='Galaxy configuration file',
dest='config_file', action='store')
parser.add_option('-s', '--http-server-section', help='Name of the HTTP server section in the Galaxy configuration file',
dest='http_server_section', action='store')
(opts, args) = parser.parse_args()
log.debug( "GALAXY LISTENER PID: " + str(os.getpid()) + " - " + str( opts ) )
# read the Galaxy config file
global config_file_name
config_file_name = opts.config_file
global config
config = ConfigParser.ConfigParser()
config.read( sys.argv[1] )
global dbconnstr
dbconnstr = config.get("app:main", "database_connection")
config.read( opts.config_file )
global http_server_section
http_server_section = opts.http_server_section
amqp_config = {}
for option in config.options("galaxy_amqp"):
amqp_config[option] = config.get("galaxy_amqp", option)
log.debug("PID: " + str(os.getpid()) + ", " + str(amqp_config))
# web server config
global webconfig
webconfig = ConfigParser.ConfigParser()
webconfig.read('transfer_datasets.ini')
log.debug( str( amqp_config ) )
# connect
conn = amqp.Connection(host=amqp_config['host']+":"+amqp_config['port'],
userid=amqp_config['userid'],
@@ -19,6 +19,8 @@ import urllib,urllib2, cookielib, shutil
import logging, time, datetime
import xml.dom.minidom
from xml_helper import get_value, get_value_index
log = logging.getLogger( "datatx_" + str( os.getpid() ) )
log.setLevel( logging.DEBUG )
fh = logging.FileHandler( "data_transfer.log" )
@@ -36,35 +38,38 @@ new_path = [ os.path.join( os.getcwd(), "lib" ) ]
new_path.extend( sys.path[1:] ) # remove scripts/ from the path
sys.path = new_path
from galaxy import eggs
from galaxy.util.json import from_json_string, to_json_string
from galaxy.model import SampleDataset
from galaxy.web.api.requests import RequestsController
from galaxy import eggs
import pkg_resources
pkg_resources.require( "pexpect" )
import pexpect
pkg_resources.require( "simplejson" )
import simplejson
log.debug(str(dir(api)))
class DataTransfer( object ):
def __init__( self, msg, config_file ):
log.info( msg )
self.dom = xml.dom.minidom.parseString( msg )
self.sequencer_host = self.get_value( self.dom, 'data_host' )
self.sequencer_username = self.get_value( self.dom, 'data_user' )
self.sequencer_password = self.get_value( self.dom, 'data_password' )
self.request_id = self.get_value( self.dom, 'request_id' )
self.sample_id = self.get_value( self.dom, 'sample_id' )
self.library_id = self.get_value( self.dom, 'library_id' )
self.folder_id = self.get_value( self.dom, 'folder_id' )
self.galaxy_host = get_value( self.dom, 'galaxy_host' )
self.api_key = get_value( self.dom, 'api_key' )
self.sequencer_host = get_value( self.dom, 'data_host' )
self.sequencer_username = get_value( self.dom, 'data_user' )
self.sequencer_password = get_value( self.dom, 'data_password' )
self.request_id = get_value( self.dom, 'request_id' )
self.sample_id = get_value( self.dom, 'sample_id' )
self.library_id = get_value( self.dom, 'library_id' )
self.folder_id = get_value( self.dom, 'folder_id' )
self.dataset_files = []
count=0
while True:
dataset_id = self.get_value_index( self.dom, 'dataset_id', count )
file = self.get_value_index( self.dom, 'file', count )
name = self.get_value_index( self.dom, 'name', count )
dataset_id = get_value_index( self.dom, 'dataset_id', count )
file = get_value_index( self.dom, 'file', count )
name = get_value_index( self.dom, 'name', count )
if file:
self.dataset_files.append( dict( name=name,
dataset_id=int( dataset_id ),
@@ -72,18 +77,6 @@ class DataTransfer( object ):
else:
break
count=count+1
try:
# Retrieve the upload user login information from the config file
transfer_datasets_config = ConfigParser.ConfigParser( )
transfer_datasets_config.read( 'transfer_datasets.ini' )
self.data_transfer_user_email = transfer_datasets_config.get( "data_transfer_user_login_info", "email" )
self.data_transfer_user_password = transfer_datasets_config.get( "data_transfer_user_login_info", "password" )
self.api_key = transfer_datasets_config.get( "data_transfer_user_login_info", "api_key" )
self.http_server_section = transfer_datasets_config.get( "universe_wsgi_config", "http_server_section" )
except:
log.error( traceback.format_exc() )
log.error( 'ERROR reading config values from transfer_datasets.ini.' )
sys.exit(1)
# read config variables
config = ConfigParser.ConfigParser()
retval = config.read( config_file )
@@ -91,14 +84,6 @@ class DataTransfer( object ):
error_msg = 'FATAL ERROR: Unable to open config file %s.' % config_file
log.error( error_msg )
sys.exit(1)
try:
self.server_host = config.get( self.http_server_section, "host" )
except ConfigParser.NoOptionError,e:
self.server_host = '127.0.0.1'
try:
self.server_port = config.get( self.http_server_section, "port" )
except ConfigParser.NoOptionError,e:
self.server_port = '8080'
try:
self.config_id_secret = config.get( "app:main", "id_secret" )
except ConfigParser.NoOptionError,e:
@@ -198,9 +183,8 @@ class DataTransfer( object ):
data[ 'dbkey' ] = ''
data[ 'upload_option' ] = 'upload_directory'
data[ 'create_type' ] = 'file'
url = "http://%s:%s/api/libraries/%s/contents" % ( self.server_host,
self.server_port,
api.encode_id( self.config_id_secret, self.library_id ) )
url = "http://%s/api/libraries/%s/contents" % ( self.galaxy_host,
api.encode_id( self.config_id_secret, self.library_id ) )
log.debug( str( ( self.api_key, url, data ) ) )
retval = api.submit( self.api_key, url, data, return_formatted=False )
log.debug( str( retval ) )
@@ -225,9 +209,8 @@ class DataTransfer( object ):
data[ 'sample_dataset_ids' ] = sample_dataset_ids
data[ 'new_status' ] = status
data[ 'error_msg' ] = msg
url = "http://%s:%s/api/requests/%s" % ( self.server_host,
self.server_port,
api.encode_id( self.config_id_secret, self.request_id ) )
url = "http://%s/api/requests/%s" % ( self.galaxy_host,
api.encode_id( self.config_id_secret, self.request_id ) )
log.debug( str( ( self.api_key, url, data)))
retval = api.update( self.api_key, url, data, return_formatted=False )
log.debug( str( retval ) )
@@ -239,31 +222,6 @@ class DataTransfer( object ):
log.error( 'FATAL ERROR' )
sys.exit( 1 )
def get_value( self, dom, tag_name ):
'''
This method extracts the tag value from the xml message
'''
nodelist = dom.getElementsByTagName( tag_name )[ 0 ].childNodes
rc = ""
for node in nodelist:
if node.nodeType == node.TEXT_NODE:
rc = rc + node.data
return rc
def get_value_index( self, dom, tag_name, dataset_id ):
'''
This method extracts the tag value from the xml message
'''
try:
nodelist = dom.getElementsByTagName( tag_name )[ dataset_id ].childNodes
except:
return None
rc = ""
for node in nodelist:
if node.nodeType == node.TEXT_NODE:
rc = rc + node.data
return rc
if __name__ == '__main__':
log.info( 'STARTING %i %s' % ( os.getpid(), str( sys.argv ) ) )
#
@@ -0,0 +1,28 @@
#======= XML helper methods ====================================================
import xml.dom.minidom
def get_value( dom, tag_name ):
'''
This method extracts the tag value from the xml message
'''
nodelist = dom.getElementsByTagName( tag_name )[ 0 ].childNodes
rc = ""
for node in nodelist:
if node.nodeType == node.TEXT_NODE:
rc = rc + node.data
return rc
def get_value_index( dom, tag_name, dataset_id ):
'''
This method extracts the tag value from the xml message
'''
try:
nodelist = dom.getElementsByTagName( tag_name )[ dataset_id ].childNodes
except:
return None
rc = ""
for node in nodelist:
if node.nodeType == node.TEXT_NODE:
rc = rc + node.data
return rc