AMQP messaging server and client files.

This commit is contained in:
Ramkrishna Chakrabarty
2009-11-04 10:04:51 -05:00
parent 77afaab66d
commit f4bedf8e5c
9 changed files with 307 additions and 19 deletions
@@ -0,0 +1,87 @@
'''
This script gets barcode data from a barcode scanner using serial communication
and sends the state representated by the barcode scanner & the barcode string
to the Galaxy LIMS RabbitMQ server. The message is sent in XML which has 2 tags,
barcode & state. The state of the scanner should be set in the galaxy_amq.ini
file as a configuration variable.
'''
from amqplib import client_0_8 as amqp
import ConfigParser
import sys, os
import serial
import array
import time
import optparse
xml = \
''' <sample>
<barcode>%(BARCODE)s</barcode>
<state>%(STATE)s</state>
</sample>'''
def handle_scan(states, amqp_config, barcode):
if states.get(barcode[:2], None):
values = dict( BARCODE=barcode[2:],
STATE=states.get(barcode[:2]) )
print values
data = xml % values
print data
conn = amqp.Connection(host=amqp_config['host']+":"+amqp_config['port'],
userid=amqp_config['userid'],
password=amqp_config['password'],
virtual_host=amqp_config['virtual_host'],
insist=False)
chan = conn.channel()
msg = amqp.Message(data)
msg.properties["delivery_mode"] = 2
chan.basic_publish(msg,
exchange=amqp_config['exchange'],
routing_key=amqp_config['routing_key'])
chan.close()
conn.close()
def recv_data(states, amqp_config, s):
while True:
bytes = s.inWaiting()
if bytes:
print '%i bytes recvd' % bytes
msg = s.read(bytes)
print msg
handle_scan(states, amqp_config, msg.strip())
def main():
parser = optparse.OptionParser()
parser.add_option('-c', '--config-file', help='config file with all the AMQP config parameters',
dest='config_file', action='store')
parser.add_option('-p', '--port', help='Name of the port where the scanner is connected',
dest='port', action='store')
(opts, args) = parser.parse_args()
config = ConfigParser.ConfigParser()
config.read(opts.config_file)
amqp_config = {}
states = {}
for option in config.options("galaxy:amqp"):
amqp_config[option] = config.get("galaxy:amqp", option)
count = 1
while True:
section = 'scanner%i' % count
if config.has_section(section):
states[config.get(section, 'prefix')] = config.get(section, 'state')
count = count + 1
else:
break
print amqp_config
print states
s = serial.Serial(int(opts.port))
print 'Port %s is open: %s' %( opts.port, s.isOpen())
recv_data(states, amqp_config, s)
s.close()
print 'Port %s is open: %s' %( opts.port, s.isOpen())
if __name__ == '__main__':
main()
@@ -0,0 +1,32 @@
# Galaxy Message Queue
# Galaxy uses AMQ protocol to receive messages from external sources like
# bar code scanners. Galaxy has been tested against RabbitMQ AMQP implementation.
# For Galaxy to receive messages from a message queue the RabbitMQ server has
# to be set up with a user account and other parameters listed below. The 'host'
# and 'port' fields should point to where the RabbitMQ server is running.
#[galaxy:amqp]
#host = 127.0.0.1
#port = 5672
#userid = galaxy
#password = galaxy
#virtual_host = galaxy_messaging_engine
#queue = galaxy_queue
#exchange = galaxy_exchange
#routing_key = bar_code_scanner
# The following section(s) 'scanner#' is for specifying the state of the
# sample this scanner represents. This state name should be one of the
# possible states created for this request type in Galaxy
# If there multiple scanners attached to this host the add as many "scanner#"
# sections below each with the name & prefix of the bar code scanner and
# the state it represents
#[scanner1]
#name =
#state =
#prefix =
#[scanner2]
#name =
#state =
#prefix =
@@ -0,0 +1 @@
python scanner.py -p 2 -c galaxy_amq.ini -r
@@ -0,0 +1 @@
python amqp_publisher.py -p 2 -c galaxy_amq.ini
@@ -0,0 +1 @@
python amqp_publisher.py -p 3 -c galaxy_amq.ini
@@ -0,0 +1,92 @@
import sys, os
import serial
import array
import time
import optparse
import ConfigParser, logging
from scanner_interface import ScannerInterface
logging.basicConfig(level=logging.DEBUG)
log = logging.getLogger( 'Scanner' )
# command prefix: SYN M CR
cmd = [22, 77, 13]
response = { 6: 'ACK', 5: 'ENQ', 21: 'NAK' }
image_scanner_report = 'RPTSCN.'
get_prefix1 = 'PREBK2?.'
get_prefix2 = ':4820:PREBK2?.'
set_prefix = 'PREBK2995859.'
clear_prefix = 'PRECA2.'
def get_prefix_cmd(name):
return ':' + name + ':' + 'PREBK2?.'
def set_prefix_cmd(name, prefix):
prefix_str = ''
for c in prefix:
prefix_str = prefix_str + hex(ord(c))[2:]
return ':' + name + ':' + 'PREBK299' + prefix_str + '!'
def read_config_file(config_file):
config = ConfigParser.ConfigParser()
config.read(config_file)
count = 1
scanners_list = []
while True:
section = 'scanner%i' % count
if config.has_section(section):
scanner = dict(name=config.get(section, 'name'),
prefix=config.get(section, 'prefix'),
state=config.get(section, 'state'))
scanners_list.append(scanner)
count = count + 1
else:
return scanners_list
def main():
usage = "python %s -p PORT -c CONFIG_FILE [ OPTION ]" % sys.argv[0]
parser = optparse.OptionParser(usage=usage)
parser.add_option('-p', '--port', help='Name of the port where the scanner is connected',
dest='port', action='store')
parser.add_option('-c', '--config-file', help='config file with all the AMQP config parameters',
dest='config_file', action='store')
parser.add_option('-r', '--report', help='scanner report',
dest='report', action='store_true', default=False)
parser.add_option('-i', '--install', help='install the scanners',
dest='install', action='store_true', default=False)
(opts, args) = parser.parse_args()
# validate
if not opts.port:
parser.print_help()
sys.exit(0)
if ( opts.report or opts.install ) and not opts.config_file:
parser.print_help()
sys.exit(0)
# create the scanner interface
si = ScannerInterface(opts.port)
if opts.install:
scanners_list = read_config_file(opts.config_file)
for scanner in scanners_list:
msg = set_prefix_cmd(scanner['name'], scanner['prefix'])
si.send(msg)
response = si.recv()
if not response:
log.error("Scanner %s could not be installed." % scanner['name'])
elif opts.report:
si.send(image_scanner_report)
rep = si.recv()
log.info(rep)
scanners_list = read_config_file(opts.config_file)
for scanner in scanners_list:
msg = get_prefix_cmd(scanner['name'])
si.send(msg)
response = si.recv()
if response:
log.info('PREFIX for scanner %s: %s' % (scanner['name'], chr(int(response[8:12][:2], 16))+chr(int(response[8:12][2:], 16)) ))
si.close()
if __name__ == "__main__":
main()
@@ -0,0 +1,76 @@
import sys, os
import serial
import array
import time
import optparse
import ConfigParser
import logging
logging.basicConfig(level=logging.INFO)
log = logging.getLogger( 'ScannerInterface' )
class ScannerInterface( object ):
cmdprefix = [22, 77, 13]
response = { 6: 'ACK', 5: 'ENQ', 21: 'NAK' }
def __init__( self, port ):
if os.name in ['posix', 'mac']:
self.port = port
elif os.name == 'nt':
self.port = int(port)
if self.port:
self.open()
def open(self):
try:
self.serial_conn = serial.Serial(self.port)
except serial.SerialException:
log.exception('Unable to open port: %s' % str(self.port))
sys.exit(1)
log.debug('Port %s is open: %s' %( str(self.port), self.serial_conn.isOpen() ) )
def is_open(self):
return self.serial_conn.isOpen()
def close(self):
self.serial_conn.close()
log.debug('Port %s is open: %s' %( str(self.port), self.serial_conn.isOpen() ) )
def send(self, msg):
message = self.cmdprefix + map(ord, msg)
byte_array = array.array('B', message)
log.debug('Sending message to %s: %s' % ( str(self.port), message) )
bytes = self.serial_conn.write( byte_array.tostring() )
log.debug('%i bytes out of %i bytes sent to the scanner' % ( bytes, len(message) ) )
def recv(self):
time.sleep(1)
self.serial_conn.flush()
nbytes = self.serial_conn.inWaiting()
log.debug('%i bytes received' % nbytes)
if nbytes:
msg = self.serial_conn.read(nbytes)
byte_array = map(ord, msg)
log.debug('Message received [%s]: %s' % (self.response.get(byte_array[len(byte_array)-2], byte_array[len(byte_array)-2]),
msg))
return msg
else:
log.error('Error!')
return None
def setup_recv(self, callback):
self.recv_callback = callback
def wait(self):
nbytes = self.serial_conn.inWaiting()
if nbytes:
msg = self.serial_conn.read(nbytes)
byte_array = map(ord, msg)
log.debug('Message received [%s]: %s' % (self.response.get(byte_array[len(byte_array)-2], byte_array[len(byte_array)-2],
msg)))
if self.recv_callback:
self.recv_callback(msg)
return
@@ -1,6 +1,6 @@
#/usr/bin/python
from datetime import datetime, timedelta
from datetime import datetime
import sys
import optparse
import os
@@ -99,10 +99,9 @@ class GalaxyDbInterface(object):
if new_state_id == -1:
return
log.debug('Updating sample_id %i state to %s' % (sample_id, new_state))
d = timedelta(hours=4)
i = self.event_table.insert()
i.execute(update_time=datetime.now()+d,
create_time=datetime.now()+d,
i.execute(update_time=datetime.utcnow(),
create_time=datetime.utcnow(),
sample_id=sample_id,
sample_state_id=int(new_state_id),
comment='bar code scanner')
@@ -122,7 +121,6 @@ class GalaxyDbInterface(object):
else:
request_state = 'Submitted'
log.debug('Updating request_id %i state to "%s"' % (self.request_id, request_state))
d = timedelta(hours=4)
i = self.request_table.update(whereclause=self.request_table.c.id==self.request_id,
values={self.request_table.c.state: request_state})
i.execute()
@@ -132,20 +130,20 @@ class GalaxyDbInterface(object):
if __name__ == '__main__':
print '''This file should not be run directly. To start the Galaxy AMQP Listener:
%sh run_galaxy_listener.sh'''
# dbstr = 'postgres://postgres:postgres@localhost/galaxy_ft'
#
# parser = optparse.OptionParser()
# parser.add_option('-n', '--name', help='name of the sample field', dest='name', \
# action='store', default='bar_code')
# parser.add_option('-v', '--value', help='value of the sample field', dest='value', \
# action='store')
# parser.add_option('-s', '--state', help='new state of the sample', dest='state', \
# action='store')
# (opts, args) = parser.parse_args()
#
# gs = GalaxyDbInterface(dbstr)
# sample_id = gs.get_sample_id(field_name=opts.name, value=opts.value)
# gs.change_state(sample_id, opts.state)
dbstr = 'postgres://postgres:postgres@localhost/galaxy_uft'
parser = optparse.OptionParser()
parser.add_option('-n', '--name', help='name of the sample field', dest='name', \
action='store', default='bar_code')
parser.add_option('-v', '--value', help='value of the sample field', dest='value', \
action='store')
parser.add_option('-s', '--state', help='new state of the sample', dest='state', \
action='store')
(opts, args) = parser.parse_args()
gs = GalaxyDbInterface(dbstr)
sample_id = gs.get_sample_id(field_name=opts.name, value=opts.value)
gs.change_state(sample_id, opts.state)