diff --git a/Source/EventCenterPlugin/EventCenterPlugin.Build.cs b/Source/EventCenterPlugin/EventCenterPlugin.Build.cs index a5e1b52..bdd6c4a 100644 --- a/Source/EventCenterPlugin/EventCenterPlugin.Build.cs +++ b/Source/EventCenterPlugin/EventCenterPlugin.Build.cs @@ -18,7 +18,13 @@ public class EventCenterPlugin : ModuleRules // Add this in the Unreal Engine Plugins menu. // If not present, use the Epic Games Launcher's Unreal Library Marketplace. "RedisPlugin/Public", - "RedisPlugin/Classes" + "RedisPlugin/Classes", + + // NB(bt): Using the commercial BlueprintWebSocket plugin described in README.md + // Add this in the Unreal Engine Plugins menu. + // If not present, use the Epic Games Launcher's Unreal Library Marketplace. + "BlueprintWebSocket/Public" + //"BlueprintWebSocket/Private", } ); @@ -39,7 +45,12 @@ public class EventCenterPlugin : ModuleRules // NB(bt): Using the commercial Redis Plugin by GameSeed (AKA "SDRedis") // Add this in the Unreal Engine Plugins menu. // If not present, use the Epic Games Launcher's Unreal Library Marketplace. - "RedisPlugin" + "RedisPlugin", + + // NB(bt): Using the commercial BlueprintWebSocket plugin described in README.md + // Add this in the Unreal Engine Plugins menu. + // If not present, use the Epic Games Launcher's Unreal Library Marketplace. + "BlueprintWebSocket" } ); diff --git a/Source/EventCenterPlugin/Public/EventSourceActor.h b/Source/EventCenterPlugin/Public/EventSourceActor.h index eb79d3b..9ff9454 100644 --- a/Source/EventCenterPlugin/Public/EventSourceActor.h +++ b/Source/EventCenterPlugin/Public/EventSourceActor.h @@ -7,6 +7,9 @@ // RedisPlugin #include "RedisObject.h" +// BlueprintWebSocketPlugin +#include "BlueprintWebSocketWrapper.h" + // EventCenterPlugin #include "EventSourceActor.generated.h" diff --git a/server.py b/server.py index ad85cd2..da6f74b 100644 --- a/server.py +++ b/server.py @@ -1,5 +1,6 @@ #!/usr/bin/env python3 +import aioredis import asyncio import redis import sys @@ -7,21 +8,23 @@ 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) -#redis_connection = redis_pool.redis_connect() -#pubsub = redis_connection.pubsub() 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() @@ -38,6 +41,8 @@ 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" @@ -48,5 +53,47 @@ 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) + + 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)) + + asyncio.ensure_future(reader(mpsc)) + + redis_connection = await aioredis.create_connection(('localhost', 6379)) + + await redis_connection.subscribe(mpsc.channel('channel:1'), mpsc.channel('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) + + 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() +