-
Notifications
You must be signed in to change notification settings - Fork 1
/
Copy pathvoevent_handler.py
executable file
·309 lines (255 loc) · 12.6 KB
/
voevent_handler.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
#!/usr/bin/env python
"""Handler daemon, runs continuously and accepts VO-Events in XML format via Pyro from the push_voevent.py
script (called by the COMET event broker, possibly multiple times in parallel). XML packets are serialised
in a queue and processed one by one by passing them to one or more plugins, defined in the 'handlers' module.
Use the '-p' argument to run in 'pretend' mode, where the MWA schedule is not actually changed when a trigger is
sent.
"""
import logging
import os
import pwd
import socket
import sys
import threading
import time
import traceback
if sys.version_info.major == 2:
from ConfigParser import SafeConfigParser as conparser
import Queue
else:
from configparser import ConfigParser as conparser
import queue as Queue
from astropy.time import Time
import voeventparse
IVORN_LIST = [] # Maintain list of ivorn names that have already been processed by the queue
EXCEPTION_NOTIFY_LIST = ["[email protected]"]
EXCEPTION_EMAIL_TEMPLATE = """
The VOEvent handler threw an exception:
%s
"""
############### set up the logging before importing Pyro4
class MWALogFormatter(logging.Formatter):
"""
Add a time string to the start of any log messages sent to the log file.
"""
def format(self, record):
now = Time.now()
return "%s=(%d): %s" % (now.iso, int(now.gps), record.getMessage())
LOGLEVEL_LOGFILE = logging.DEBUG # Logging level for logfile
# Make the log file name include the username, to avoid permission errors
LOGFILE = "/var/log/mwa/voevents-%s.log" % pwd.getpwuid(os.getuid()).pw_name
formatter = MWALogFormatter()
filehandler = logging.FileHandler(LOGFILE)
filehandler.setLevel(LOGLEVEL_LOGFILE)
filehandler.setFormatter(formatter)
DEFAULTLOGGER = logging.getLogger('voevent')
DEFAULTLOGGER.addHandler(filehandler)
##############
import Pyro4
import Pyro4.errors
import Pyro4.socketutil
sys.excepthook = Pyro4.util.excepthook
Pyro4.config.DETAILED_TRACEBACK = True
from mwa_trigger import handlers
from mwa_trigger import GRB_fermi_swift, FlareStar_swift_maxi, GW_LIGO, Neutrino
PRETEND = False # Set to true to trigger event in 'pretend' mode, not actually schedule observations.
# One or more handler functions - all will be called in turn on each XML event.
EVENTHANDLERS = [
GRB_fermi_swift.processevent,
# FlareStar_swift_maxi.processevent,
# GW_LIGO.processevent,
Neutrino.processevent,
]
Pyro4.config.COMMTIMEOUT = 10.0
Pyro4.config.THREADPOOL_SIZE_MIN = 8
Pyro4.config.SERIALIZERS_ACCEPTED.add('pickle')
REFERENCEIP = '8.8.8.8' # A host guaranteed to be visible on the network interface that we want the Pyro server to bind to
EXITING = None
PYRO_DAEMON = None
CPPATH = ['/usr/local/etc/trigger.conf', './trigger.conf'] # Path list to look for configuration file
############## Point to a running Pyro nameserver #####################
# If not on site, start one before running this code, using pyro_nameserver.py
CP = conparser()
CP.read(CPPATH)
if CP.has_option(section='pyro', option='ns_host'):
Pyro4.config.NS_HOST = CP.get(section='pyro', option='ns_host')
else:
Pyro4.config.NS_HOST = 'localhost'
if CP.has_option(section='pyro', option='ns_port'):
Pyro4.config.NS_PORT = int(CP.get(section='pyro', option='ns_port'))
else:
Pyro4.config.NS_PORT = 9090
if Pyro4.config.NS_HOST in ['helios', 'mwa-db']:
Pyro4.config.SERIALIZER = 'pickle' # We must be on site, where we have an ancient Pyro4 install and nameserver running
############### Main event handler - receives VOEvent objects by RPC and queues them for processing #################
class VOEventHandler(object):
"""
Implements Pyro4 methods that can be called remotely via RPC, to process VO event XML strings.
"""
def __init__(self, logger=DEFAULTLOGGER):
self.logger = logger
@Pyro4.expose
def ping(self):
"""
Empty function, call to see if RPC connection is live.
"""
pass
@Pyro4.expose
def putEvent(self, event=None):
"""
Called by the remote client to send a VOEvent XML packet to this server. It just
pushes the complete XML packet onto the queue for background processing, and returns
immediately.
:param event: string containing XML format VOEvent.
"""
EventQueue.put(event)
self.logger.info("Queued VOEvent XML, current queue size is %d" % EventQueue.qsize())
def servePyroRequests(self):
"""
When called, start serving Pyro requests. Only exits if the global EXITING is set to True
externally, to trigger a clean shutdown.
"""
global EXITING, PYRO_DAEMON
iface = None
self.logger.info('Getting interface address for VOEventHandler Pyro server')
while iface is None:
try:
iface = Pyro4.socketutil.getInterfaceAddress(REFERENCEIP) # What is the network IP of this receiver?
except socket.error:
self.logger.info("Network down, can't start pycontroller Pyro server, sleeping for 10 seconds")
time.sleep(10)
ns = None
while not EXITING:
self.logger.info("Starting VOEventHandler Pyro4 server")
if ns is not None:
try:
ns._pyroRelease()
except Pyro4.errors.PyroError:
pass
try:
ns = Pyro4.locateNS()
except Pyro4.errors.PyroError:
self.logger.error("Can't locate Pyro nameserver - waiting 10 sec to retry")
if not EXITING:
time.sleep(10)
continue
try:
existing = ns.lookup("VOEventHandler")
self.logger.info("VOEventHandler still exists in Pyro nameserver with id: %s" % existing.object)
self.logger.info("Previous Pyro daemon socket port: %d" % existing.port)
# start the daemon on the previous port
PYRO_DAEMON = Pyro4.Daemon(host=iface, port=existing.port)
# register the object in the daemon with the old objectId
PYRO_DAEMON.register(self, objectId=existing.object)
except (Pyro4.errors.PyroError, socket.error):
try:
# just start a new daemon on a random port
PYRO_DAEMON = Pyro4.Daemon(host=iface)
# register the object in the daemon and let it get a new objectId
# also need to register in name server because it's not there yet.
uri = PYRO_DAEMON.register(self)
ns.register("VOEventHandler", uri)
except (Pyro4.errors.PyroError, socket.error):
if not EXITING:
self.logger.error("Exception in VOEventHandler Pyro4 startup. Retrying in 10 sec: %s" % (traceback.format_exc(),))
time.sleep(10)
else:
self.logger.error("Exception in VOEventHandler, EXITING is true, shutting down: %s" % (traceback.format_exc(),))
continue
except Exception:
if not EXITING:
self.logger.error("Exception in VOEventHandler Pyro4 start. Retrying in 10 sec: %s" % (traceback.format_exc(),))
handlers.send_email(from_address='[email protected]',
to_addresses=EXCEPTION_NOTIFY_LIST,
subject='Exception in Pyro4 daemon request loop',
msg_text=EXCEPTION_EMAIL_TEMPLATE % traceback.format_exc())
time.sleep(10)
else:
self.logger.error("Exception in VOEventHandler, EXITING true, shutting down: %s" % (traceback.format_exc(),))
continue
if not EXITING:
try:
ns._pyroRelease()
PYRO_DAEMON.requestLoop()
except Exception:
if not EXITING:
self.logger.error("Exception in VOEventHandler Pyro4 server. Restarting in 10 sec: %s" % (traceback.format_exc(),))
handlers.send_email(from_address='[email protected]',
to_addresses=EXCEPTION_NOTIFY_LIST,
subject='Exception in Pyro4 daemon request loop',
msg_text=EXCEPTION_EMAIL_TEMPLATE % traceback.format_exc())
time.sleep(10)
else:
self.logger.error("Exception in VOEventHandler Pyro4 server, EXITING true, shutting down: %s" % (traceback.format_exc(),))
PYRO_DAEMON.close()
def QueueWorker():
"""
Worker thread to process incoming message packets in the EventQueue. It is spawned on startup, and run continuously,
blocking on EventQueue.get() if there's nothing to process. When an item is 'put' on the queue, the EventQueue.get()
returns and the event is processed.
Only exits if the global EXITING is set to True externally, to trigger a clean shutdown.
"""
global EXITING
global IVORN_LIST
try:
while not EXITING:
eventxml = EventQueue.get()
if sys.version_info.major == 2:
# event arrives as a unicode string but loads requires a non-unicode string.
v = voeventparse.loads(str(eventxml))
else:
v = voeventparse.loads(eventxml.encode('latin-1'))
if v.attrib['ivorn'] in IVORN_LIST:
DEFAULTLOGGER.info("Already seen event %s, discarding. Current queue size is %d" % (v.attrib['ivorn'],
EventQueue.qsize()))
else:
DEFAULTLOGGER.info("Processing event %s. Current queue size is %d" % (v.attrib['ivorn'],
EventQueue.qsize()))
IVORN_LIST.append(v.attrib['ivorn'])
for hfunc in EVENTHANDLERS:
handled = hfunc(event=eventxml, pretend=PRETEND)
if handled: # One of the handlers accepted this event
break # Don't try any more event handlers.
EventQueue.task_done()
except Exception:
DEFAULTLOGGER.error("Exception in QueueWorker. Restarting in 10 sec: %s" % (traceback.format_exc(),))
handlers.send_email(from_address='[email protected]',
to_addresses=EXCEPTION_NOTIFY_LIST,
subject='Exception in QueueWorker loop',
msg_text=EXCEPTION_EMAIL_TEMPLATE % traceback.format_exc())
if __name__ == '__main__':
if (len(sys.argv) > 1) and '-p' in sys.argv:
PRETEND = True
if PRETEND:
DEFAULTLOGGER.info('Working in PRETEND mode, not actually scheduling observations.')
EventQueue = Queue.Queue(maxsize=10)
while True:
# Start a background thread accepting network connections that add events to the queue.
rpcHandler = VOEventHandler(logger=DEFAULTLOGGER)
pyro_thread = threading.Thread(target=rpcHandler.servePyroRequests, name='PyroDaemon')
pyro_thread.daemon = True
DEFAULTLOGGER.info('Starting Pyro4 request handler.')
pyro_thread.start()
# Start a background thread to process incoming events from the queue, one by one.
queue_thread = threading.Thread(target=QueueWorker, name='QueueDaemon')
queue_thread.daemon = True
DEFAULTLOGGER.info('Starting Queue handler.')
queue_thread.start()
try:
while True:
time.sleep(5)
if not pyro_thread.is_alive():
DEFAULTLOGGER.error('Pyro request handler thread has died - restarting.')
handlers.send_email(from_address='[email protected]',
to_addresses=EXCEPTION_NOTIFY_LIST,
subject='Exception in Pyro4 daemon request loop',
msg_text=EXCEPTION_EMAIL_TEMPLATE % 'Pyro request handler thread has died - restarting.')
break
if not queue_thread.is_alive():
DEFAULTLOGGER.error('Queue handler thread has died - restarting.')
break
finally:
EXITING = True
PYRO_DAEMON.shutdown()
time.sleep(15)
EXITING = False