Files
brewpi/ws/message.py
T

146 lines
3.3 KiB
Python

import asyncio
from ws.connection import IConnection
import json
import queue
class MsgIo:
def __init__(self, key, send, send_from_thread):
self.key = key
self.sender = send
self.send_from_thread_cb = send_from_thread
self.receiver = None
def set_recv_handler(self, handler):
self.receiver = handler
async def on_recv(self, data):
print("MessageHandler {}".format(data))
if self.receiver is not None:
await self.receiver(data[self.key])
async def send(self, data):
return await self.sender({self.key: data})
def send_from_thread(self, data):
return self.send_from_thread_cb({self.key: data})
def can_key(self, key):
return self.key == key
def get_key(self):
return self.key
class MessageDispatcher(IConnection):
def __init__(self, auto_subscribe=False, loop=None):
self.msg_handlers: MsgIo = []
self.state = None
self.state = asyncio.Queue(loop=loop)
self.auto_subscribe = auto_subscribe
self.loop = loop
def msgio_get(self, key):
obj = MsgIo(key, self.send, self.send_from_thread)
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))
await 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 h.on_recv(d)
async def on_send(self):
return await self.state.get()
async def send(self, data):
await self.state.put(data)
def send_from_thread(self, data):
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:
asyncio.get_event_loop().run_in_executor(None, self.subscribe)
def subscribe(self):
for handler in self.msg_handlers:
key = handler.get_key()
print("Would subscribe {}".format(key))
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):
result = None
try:
result = await asyncio.get_event_loop().run_in_executor(None, self.state.get, True, 1.0)
print("Send subscribe {}".format(result))
except Exception as e:
pass
return result
def send(self, data):
self.state.put(data)