- removed MsgIo::send_from_thread
- created sync variant of MessageDispatcher and MsgIo - fixed delay in sending client messages
This commit is contained in:
+65
-2
@@ -1,7 +1,7 @@
|
||||
import asyncio
|
||||
from ws.connection import IConnection
|
||||
import json
|
||||
|
||||
import queue
|
||||
|
||||
class MsgIo:
|
||||
def __init__(self, key, send, send_from_thread):
|
||||
@@ -70,4 +70,67 @@ class MessageDispatcher(IConnection):
|
||||
await self.state.put(data)
|
||||
|
||||
def send_from_thread(self, data):
|
||||
asyncio.ensure_future(self.send(data), loop=self.loop)
|
||||
self.loop.create_task(self.send(data))
|
||||
|
||||
|
||||
class MsgIoSync:
|
||||
def __init__(self, key, send):
|
||||
self.key = key
|
||||
self.sender = send
|
||||
self.receiver = None
|
||||
|
||||
def set_recv_handler(self, handler):
|
||||
self.receiver = handler
|
||||
|
||||
def on_recv(self, data):
|
||||
print("MessageHandler {}".format(data))
|
||||
if self.receiver is not None:
|
||||
self.receiver(data[self.key])
|
||||
|
||||
def send(self, data):
|
||||
return self.sender({self.key: data})
|
||||
|
||||
def can_key(self, key):
|
||||
return self.key == key
|
||||
|
||||
def get_key(self):
|
||||
return self.key
|
||||
|
||||
|
||||
class MessageDispatcherSync(IConnection):
|
||||
def __init__(self, auto_subscribe=False):
|
||||
self.msg_handlers: MsgIoSync = []
|
||||
self.state = None
|
||||
self.state = queue.Queue()
|
||||
self.auto_subscribe = auto_subscribe
|
||||
|
||||
def msgio_get(self, key):
|
||||
obj = MsgIoSync(key, self.send)
|
||||
if key not in self.msg_handlers:
|
||||
self.msg_handlers.append(obj)
|
||||
return obj
|
||||
return None
|
||||
|
||||
async def on_connect(self):
|
||||
if self.auto_subscribe:
|
||||
for handler in self.msg_handlers:
|
||||
key = handler.get_key()
|
||||
print("Would subscribe {}".format(key))
|
||||
asyncio.get_event_loop().run_in_executor(None, self.send, {'+': key})
|
||||
|
||||
async def on_disconnect(self):
|
||||
pass
|
||||
|
||||
async def on_recv(self, data):
|
||||
d = json.loads(data)
|
||||
for key_req in d.keys():
|
||||
for h in self.msg_handlers:
|
||||
if h.can_key(key_req):
|
||||
await asyncio.get_event_loop().run_in_executor(None, h.on_recv, d)
|
||||
|
||||
async def on_send(self):
|
||||
return await asyncio.get_event_loop().run_in_executor(None, self.state.get)
|
||||
|
||||
def send(self, data):
|
||||
self.state.put(data)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user