mirror of
https://github.com/storytold/EventCenterPlugin.git
synced 2026-10-09 00:09:42 +00:00
75 lines
2.0 KiB
Python
75 lines
2.0 KiB
Python
#!/usr/bin/env python3
|
|
|
|
import aioredis
|
|
import asyncio
|
|
import redis
|
|
import sys
|
|
import toml
|
|
import websockets
|
|
|
|
from aioredis.pubsub import Receiver
|
|
from aioredis.abc import AbcChannel
|
|
|
|
CHANNELS = [
|
|
'vocode',
|
|
'spawn',
|
|
'testing',
|
|
]
|
|
|
|
print('reading secrets', flush=True)
|
|
SECRETS = toml.load("secrets.toml")
|
|
|
|
print('defining server', flush=True)
|
|
|
|
async def server(websocket, path):
|
|
print('subscribing to redis pubsub', flush=True)
|
|
|
|
queue = asyncio.Queue()
|
|
|
|
mpsc = Receiver(loop=loop)
|
|
async def reader(mpsc):
|
|
async for channel, msg in mpsc.iter():
|
|
assert isinstance(channel, AbcChannel)
|
|
print("Got {!r} in channel {!r}".format(msg, channel))
|
|
|
|
#print('sending message over websocket we got from redis pubsub', flush=True)
|
|
#await websocket.send("Got message: {}".format(msg))
|
|
|
|
channel_name = channel.name.decode('utf-8')
|
|
message = msg.decode('utf-8')
|
|
|
|
print('enqueuing websocket message', flush=True)
|
|
payload = "{}|{}".format(channel_name, message)
|
|
print('payload: {}'.format(payload, flush=True))
|
|
queue.put_nowait(payload)
|
|
|
|
asyncio.ensure_future(reader(mpsc))
|
|
|
|
redis = await aioredis.create_redis_pool(SECRETS['redis_url'])
|
|
|
|
channels = [mpsc.channel(ch) for ch in CHANNELS]
|
|
await redis.subscribe(*channels)
|
|
|
|
print('starting websocket server', flush=True)
|
|
try:
|
|
while True:
|
|
# TODO: Do we need to ping to keepalive?
|
|
if not queue.empty():
|
|
msg = queue.get_nowait() # TODO THROWS EXCEPTION IF EMPTY (data race?)
|
|
|
|
print('sending pubsub msg over websocket: {}'.format(msg), flush=True)
|
|
await websocket.send(msg)
|
|
else:
|
|
await asyncio.sleep(1)
|
|
|
|
except websockets.exceptions.ConnectionClosed:
|
|
print('Connection Closed!')
|
|
|
|
loop = asyncio.get_event_loop()
|
|
loop.set_debug(True)
|
|
|
|
ws_server = websockets.serve(server, 'localhost', 8765)
|
|
loop.run_until_complete(ws_server)
|
|
loop.run_forever()
|
|
|