}
}
-def exception_callback(loop, context):
- if 'exception' in context:
- log.exception('asyncio.loop')
- raise
-
-loop = zmq.asyncio.ZMQEventLoop()
+loop = eddn.init()
asyncio.set_event_loop(loop)
-loop.set_exception_handler(exception_callback)
client = discord.Client()
bot.client = client
bot.db = db
log.info('Opening WebSocket.')
loop.run_until_complete(client.connect())
except KeyboardInterrupt:
- log.info('Logging out.')
- loop.run_until_complete(client.logout())
+ log.info('Interrupted.')
except RuntimeError:
log.exception('RuntimeError')
ret = 100
log.exception('loop')
ret = 99
finally:
+ if loop.is_running():
+ log.info('Logging out.')
+ loop.run_until_complete(client.logout())
try:
- loop.close()
+ eddn.close()
except:
ret = 101
sys.exit(ret)
schema_prefix = 'http://schemas.elite-markets.net/eddn'
topics = []
timeout = 600000
+state = { 'loop': None, 'context': None, 'subscriber': None }
+
+def exception_callback(loop, context):
+ if 'exception' in context:
+ log.exception('asyncio.loop')
+ raise
+
+def init():
+ state['loop'] = zmq.asyncio.ZMQEventLoop()
+ state['loop'].set_exception_handler(exception_callback)
+ return state['loop']
+
+def close():
+ loop = state['loop']
+ context = state['context']
+ subscriber = state['subscriber']
+
+ if loop is None:
+ log.info('Event loop is not initialised.')
+ return True
+
+ if loop.is_closed():
+ log.info('Event loop is closed.')
+ return True
+
+ if context is None:
+ log.info('Context is stopped.')
+ return True
+
+ if subscriber is not None:
+ log.info('Closing subscriber.')
+ subscriber.close()
+ else:
+ log.info('Subscriber is cloed.')
+
+ log.info('Stopping context.')
+ context.term()
+
+ if loop.is_running():
+ log.info('Stopping loop.')
+ loop.stop()
+
+ return True
async def listen(plugins):
waittime = 5
context = zmq.asyncio.Context()
+ state['context'] = context
while True:
try:
subscriber = context.socket(zmq.SUB)
+ state['subscriber'] = subscriber
log.info('Connecting to {}'.format(url))
subscriber.connect(url)
if len(topics):
except zmq.ZMQError:
log.exception('ZMQSocketException')
subscriber.disconnect(url)
+ state['subscriber'] = None
asyncio.sleep(waittime)
context.destroy()
+ state['context'] = None
async def decode_message(compressed):
try: