#!/usr/bin/env python3 import aioredis import asyncio import redis import sys import websockets #redis_client = redis.Redis(host='localhost', port=6379, db=0) # Redis pubsub (blocking) """ pool = redis.ConnectionPool(host='localhost', port=6379, password='', db=0) redis_pool = redis.StrictRedis(connection_pool=pool) pubsub = redis_pool.pubsub() pubsub.subscribe('testing') pubsub.subscribe('foo') #for message in pubsub.listen(): # if message['data'] == 'exit': # sys.exit() # print(message) """ # Single request/response websocket """ async def hello(websocket, path): name = await websocket.recv() print(f"< {name}") greeting = f"Hello {name}!" await websocket.send(greeting) print(f"> {greeting}") start_server = websockets.serve(hello, "localhost", 8765) asyncio.get_event_loop().run_until_complete(start_server) asyncio.get_event_loop().run_forever() """ # Infinitely publishing websocket server: """ async def time(websocket, path): while True: now = datetime.datetime.utcnow().isoformat() + "Z" await websocket.send(now) await asyncio.sleep(random.random() * 3) start_server = websockets.serve(time, "127.0.0.1", 5678) asyncio.get_event_loop().run_until_complete(start_server) asyncio.get_event_loop().run_forever() """ from aioredis.pubsub import Receiver from aioredis.abc import AbcChannel 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('redis://localhost') await redis.subscribe(mpsc.channel('channel:1'), mpsc.channel('channel:testing'), mpsc.channel('testing')) print('starting websocket server', flush=True) try: while True: # Receive data from "the outside world" #message = await websocket.recv() # Feed this data to the PUBLISH co-routine #await publish_to_redis(message, path) #print('sending ping', flush=True) #await websocket.send("Websocket ping") 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()