mirror of
https://github.com/storytold/EventCenterPlugin.git
synced 2026-10-09 00:09:42 +00:00
This works! WOO
This commit is contained in:
@@ -63,9 +63,9 @@ void AEventSourceActor::Tick(float DeltaTime)
|
|||||||
|
|
||||||
if (time - timerLastTime > 15.0) {
|
if (time - timerLastTime > 15.0) {
|
||||||
timerLastTime = time;
|
timerLastTime = time;
|
||||||
GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "seconds elapsed");
|
//GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "seconds elapsed");
|
||||||
redis->Unsubscribe(TEXT("second"));
|
//redis->Unsubscribe(TEXT("second"));
|
||||||
redis->Subscribe(TEXT("second"));
|
//redis->Subscribe(TEXT("second"));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -143,7 +143,7 @@ void AEventSourceActor::OnConnected()
|
|||||||
GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "OnConnected");
|
GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "OnConnected");
|
||||||
UE_LOG(EventSourceLog, Log, TEXT("OnConnected"));
|
UE_LOG(EventSourceLog, Log, TEXT("OnConnected"));
|
||||||
|
|
||||||
WebSocket->SendMessage(TEXT("Hello Server!"));
|
//WebSocket->SendMessage(TEXT("Hello Server!"));
|
||||||
}
|
}
|
||||||
|
|
||||||
void AEventSourceActor::OnConnectionError(const FString& Error)
|
void AEventSourceActor::OnConnectionError(const FString& Error)
|
||||||
@@ -189,10 +189,20 @@ void AEventSourceActor::OnMessageSent(const FString& Message)
|
|||||||
|
|
||||||
void AEventSourceActor::SetupWebSockets()
|
void AEventSourceActor::SetupWebSockets()
|
||||||
{
|
{
|
||||||
|
GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "SetupWebSockets");
|
||||||
|
UE_LOG(EventSourceLog, Log, TEXT("SetupWebSockets"));
|
||||||
|
|
||||||
if (WebSocket == nullptr) {
|
if (WebSocket == nullptr) {
|
||||||
WebSocket = UBlueprintWebSocket::CreateWebSocket();
|
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"));
|
WebSocket->Connect(TEXT("ws://localhost:8765/"), TEXT("ws"));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -64,28 +64,46 @@ print('defining server', flush=True)
|
|||||||
async def server(websocket, path):
|
async def server(websocket, path):
|
||||||
print('subscribing to redis pubsub', flush=True)
|
print('subscribing to redis pubsub', flush=True)
|
||||||
|
|
||||||
|
queue = asyncio.Queue()
|
||||||
|
|
||||||
mpsc = Receiver(loop=loop)
|
mpsc = Receiver(loop=loop)
|
||||||
async def reader(mpsc):
|
async def reader(mpsc):
|
||||||
async for channel, msg in mpsc.iter():
|
async for channel, msg in mpsc.iter():
|
||||||
assert isinstance(channel, AbcChannel)
|
assert isinstance(channel, AbcChannel)
|
||||||
print("Got {!r} in channel {!r}".format(msg, channel))
|
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))
|
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)
|
print('starting websocket server', flush=True)
|
||||||
try:
|
try:
|
||||||
while True:
|
while True:
|
||||||
# Receive data from "the outside world"
|
# Receive data from "the outside world"
|
||||||
message = await websocket.recv()
|
#message = await websocket.recv()
|
||||||
|
|
||||||
# Feed this data to the PUBLISH co-routine
|
# Feed this data to the PUBLISH co-routine
|
||||||
#await publish_to_redis(message, path)
|
#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:
|
except websockets.exceptions.ConnectionClosed:
|
||||||
print('Connection Closed!')
|
print('Connection Closed!')
|
||||||
|
|||||||
Reference in New Issue
Block a user