From 33a22303237f9d6a946932acfc8a4ebd15da79a7 Mon Sep 17 00:00:00 2001 From: Brandon Thomas Date: Thu, 18 Feb 2021 07:18:39 -0500 Subject: [PATCH] This works! WOO --- .../Private/EventSourceActor.cpp | 18 ++++++++++--- server.py | 26 ++++++++++++++++--- 2 files changed, 36 insertions(+), 8 deletions(-) diff --git a/Source/EventCenterPlugin/Private/EventSourceActor.cpp b/Source/EventCenterPlugin/Private/EventSourceActor.cpp index 4f40151..4598514 100644 --- a/Source/EventCenterPlugin/Private/EventSourceActor.cpp +++ b/Source/EventCenterPlugin/Private/EventSourceActor.cpp @@ -63,9 +63,9 @@ void AEventSourceActor::Tick(float DeltaTime) if (time - timerLastTime > 15.0) { timerLastTime = time; - GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "seconds elapsed"); - redis->Unsubscribe(TEXT("second")); - redis->Subscribe(TEXT("second")); + //GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "seconds elapsed"); + //redis->Unsubscribe(TEXT("second")); + //redis->Subscribe(TEXT("second")); } @@ -143,7 +143,7 @@ void AEventSourceActor::OnConnected() GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "OnConnected"); UE_LOG(EventSourceLog, Log, TEXT("OnConnected")); - WebSocket->SendMessage(TEXT("Hello Server!")); + //WebSocket->SendMessage(TEXT("Hello Server!")); } void AEventSourceActor::OnConnectionError(const FString& Error) @@ -189,10 +189,20 @@ void AEventSourceActor::OnMessageSent(const FString& Message) void AEventSourceActor::SetupWebSockets() { + GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "SetupWebSockets"); + UE_LOG(EventSourceLog, Log, TEXT("SetupWebSockets")); + if (WebSocket == nullptr) { WebSocket = UBlueprintWebSocket::CreateWebSocket(); } + WebSocket->OnConnectedEvent.AddDynamic(this, &AEventSourceActor::OnConnected); + WebSocket->OnConnectionErrorEvent.AddDynamic(this, &AEventSourceActor::OnConnectionError); + WebSocket->OnCloseEvent.AddDynamic(this, &AEventSourceActor::OnClosed); + WebSocket->OnMessageEvent.AddDynamic(this, &AEventSourceActor::OnMessage); + WebSocket->OnRawMessageEvent.AddDynamic(this, &AEventSourceActor::OnRawMessage); + WebSocket->OnMessageSentEvent.AddDynamic(this, &AEventSourceActor::OnMessageSent); + WebSocket->Connect(TEXT("ws://localhost:8765/"), TEXT("ws")); } diff --git a/server.py b/server.py index da6f74b..4aa41a0 100644 --- a/server.py +++ b/server.py @@ -64,28 +64,46 @@ 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)) + + print('enqueuing websocket message', flush=True) + queue.put_nowait(msg) + asyncio.ensure_future(reader(mpsc)) - redis_connection = await aioredis.create_connection(('localhost', 6379)) + redis = await aioredis.create_redis_pool('redis://localhost') - await redis_connection.subscribe(mpsc.channel('channel:1'), mpsc.channel('channel:testing')) + await redis.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() + #message = await websocket.recv() # Feed this data to the PUBLISH co-routine #await publish_to_redis(message, path) - await asyncio.sleep(1) + + print('sending ping', flush=True) + await websocket.send("Websocket ping") + + if not queue.empty(): + msg = queue.get_nowait() # TODO THROWS EXCEPTION + + print('sending pubsub msg over websocket: {}'.format(msg), flush=True) + await websocket.send(msg) + else: + await asyncio.sleep(10) except websockets.exceptions.ConnectionClosed: print('Connection Closed!')