import shlex
import sys
import time
+import zmq.asyncio
from db import DBConnection
from plugins import Plugins
import bot
import cat
+import eddn
db = DBConnection()
}
}
+asyncio.set_event_loop(zmq.asyncio.ZMQEventLoop())
client = discord.Client()
bot.client = client
bot.db = db
plugins.on_ready()
asyncio.ensure_future(maybe_sleep())
asyncio.ensure_future(request_offline_members())
+ asyncio.ensure_future(eddn.listen(plugins))
threads_ready.value = True
def mentioned_in(message, explicit = True):
--- /dev/null
+import asyncio
+import json
+import logging
+import zlib
+import zmq
+import zmq.asyncio
+
+log = logging.getLogger('eddn')
+
+url = 'tcp://eddn-relay.elite-markets.net:9500'
+topic = b''
+timeout = 600000
+
+async def listen(plugins):
+ waittime = 5
+ context = zmq.asyncio.Context()
+ while True:
+ try:
+ subscriber = context.socket(zmq.SUB)
+ subscriber.connect(url)
+ subscriber.setsockopt(zmq.SUBSCRIBE, topic)
+ subscriber.setsockopt(zmq.RCVTIMEO, timeout)
+ poller = zmq.asyncio.Poller()
+ poller.register(subscriber, zmq.POLLOUT)
+ while True:
+ try:
+ messages = await subscriber.recv_multipart()
+ except zmq.error.Again:
+ log.exception('recv_multipart')
+ messages = False
+ if messages == False:
+ break
+ for message in messages:
+ data = await decode_message(message)
+ if data is not None:
+ await plugins.eddn_message(data['$schemaRef'], data['header'], data['message'])
+ except zmq.ZMQError:
+ log.exception('ZMQSocketException')
+ subscriber.disconnect(url)
+ asyncio.sleep(waittime)
+ context.destroy()
+
+async def decode_message(compressed):
+ try:
+ uncompressed = zlib.decompress(compressed)
+ except zlib.error:
+ log.exception('Malformed gzipped message')
+ return None
+ try:
+ return json.loads(uncompressed.decode('utf-8'))
+ except ValueError:
+ log.exception('Malformed JSON in message')
+ return None
'on_member_join',
'on_member_update',
'handle_command',
- 'handle_help'
+ 'handle_help',
+ 'eddn_message'
]
self.plugins_supporting_methods = {}
self.loaded = {}
continue
return True
return False
+
+ async def eddn_message(self, schema, data, message):
+ for name in self.plugins():
+ if self.has_method(name, 'eddn_message'):
+ asyncio.ensure_future(self.call_coroutine(name, 'eddn_message', schema, data, message))