mirror of
https://github.com/storytold/EventCenterPlugin.git
synced 2026-10-09 00:09:42 +00:00
protocol fixes (still not working)
This commit is contained in:
@@ -190,13 +190,19 @@ void AEventSourceActor::OnMessageSent(const FString& Message)
|
||||
bool result = Message.Split(TEXT("|"), &channel, &payload);
|
||||
|
||||
if (!result) {
|
||||
UE_LOG(EventSourceLog, Log, TEXT("Could not split"));
|
||||
GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "Could not split");
|
||||
return;
|
||||
}
|
||||
|
||||
if (!delegateSubscriptions.Contains(channel)) {
|
||||
UE_LOG(EventSourceLog, Log, TEXT("No delegate"));
|
||||
GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "No delegate");
|
||||
return;
|
||||
}
|
||||
|
||||
UE_LOG(EventSourceLog, Log, TEXT("Executing delegate"));
|
||||
GEngine->AddOnScreenDebugMessage(-1, 15.0f, FColor::Red, "Executing delegate");
|
||||
delegateSubscriptions[channel].Execute(channel, payload);
|
||||
|
||||
}
|
||||
|
||||
@@ -75,14 +75,19 @@ async def server(websocket, path):
|
||||
#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)
|
||||
queue.put_nowait("{}|{}".format(channel, msg))
|
||||
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'))
|
||||
await redis.subscribe(mpsc.channel('channel:1'), mpsc.channel('channel:testing'), mpsc.channel('testing'))
|
||||
|
||||
print('starting websocket server', flush=True)
|
||||
try:
|
||||
@@ -94,8 +99,8 @@ async def server(websocket, path):
|
||||
#await publish_to_redis(message, path)
|
||||
|
||||
|
||||
print('sending ping', flush=True)
|
||||
await websocket.send("Websocket ping")
|
||||
#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?)
|
||||
@@ -103,7 +108,7 @@ async def server(websocket, path):
|
||||
print('sending pubsub msg over websocket: {}'.format(msg), flush=True)
|
||||
await websocket.send(msg)
|
||||
else:
|
||||
await asyncio.sleep(10)
|
||||
await asyncio.sleep(1)
|
||||
|
||||
except websockets.exceptions.ConnectionClosed:
|
||||
print('Connection Closed!')
|
||||
|
||||
Reference in New Issue
Block a user