diff --git a/brewpi.py b/brewpi.py index 6d3ebbe..e7eb820 100644 --- a/brewpi.py +++ b/brewpi.py @@ -1,117 +1,11 @@ import asyncio -from ws_server import WsServer -from user import User, UserSet, update from time import sleep import abc -import json -from connection import IConnection - +from ws.server.ws_server_multi_user import WsServerMultiUser from components.sensor.tempSensorSim import TempSensorSim, ATemperatureSensor from components.plant.pot import Pot, APlant from components.actor.heater import Heater, AHeater - - -class IWsServer: - @abc.abstractmethod - def on_recv(self, data, loop): - pass - - @abc.abstractmethod - def on_send(self, loop): - pass - - -class WsServerMultiUser(WsServer): - def __init__(self, loop=None, listener: IConnection = None): - WsServer.__init__(self, loop) - self.listener = listener - self.USERS = UserSet() - self.global_state = {} - - async def notify_state(self, data): - if self.USERS: # asyncio.wait doesn't accept an empty list - await asyncio.wait([user.send(data) for user in self.USERS]) - - async def notify_users(self): - if self.USERS: # asyncio.wait doesn't accept an empty list - message = {"info": {"type": "users", "count": len(self.USERS)}} - await asyncio.wait([user.send(message) for user in self.USERS]) - - async def register(self, websocket): - usr = User(websocket) - usr.path_add("info") - self.USERS.add(websocket, usr) - await self.notify_users() - - async def unregister(self, websocket): - self.USERS.remove(websocket) - await self.notify_users() - - async def handler_recv(self, websocket, path): - while True: - try: - data = await websocket.recv() - usr = self.USERS.get(websocket) - processed = await usr.process(data) - if processed: - await usr.send(self.global_state) - else: - await self.listener.on_recv(data) - except Exception as e: - print(e) - break - - async def handler_send(self, websocket, path): - while True: - data = await self.listener.on_send() - try: - await self.notify_state(data) - # Update global state - self.global_state = update(self.global_state, data) - except Exception as e: - print(e) - break - - -class MsgIo: - def __init__(self, key, send): - self.key = key - self.sender = send - 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 can_key(self, key): - return self.key == key - - def get_key(self): - return self.key - - -class Value: - def __init__(self, initial=None): - self.value = initial - self.has_changed = True - - def set(self, value): - self.has_changed = self.value != value - self.value = value - - def is_changed(self): - return self.has_changed - - def get(self): - self.has_changed = False - return self.value +from ws.message import MsgIo, MessageDispatcher, Value class ATask: @@ -248,38 +142,6 @@ class ColorHandler: await self.msg_handler.send(data) -class MessageDispatcher(IConnection): - def __init__(self): - self.msg_handlers: MsgIo = [] - self.state = asyncio.Queue() - - def msgio_get(self, key): - obj = MsgIo(key, self.send) - if key not in self.msg_handlers: - self.msg_handlers.append(obj) - return obj - return None - - def on_connect(self): - pass - - 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) - - class TaskManager: def __init__(self): self.tasks: ATask = [] diff --git a/brewpi_gui.py b/brewpi_gui.py index 0cb7cd4..689e448 100644 --- a/brewpi_gui.py +++ b/brewpi_gui.py @@ -1,11 +1,8 @@ from brewpi_win import Ui_MainWindow -from PyQt5 import QtGui, QtWidgets, QtCore +from PyQt5 import QtWidgets import sys -from ws_client import WsClient -from connection import IConnection -import json -from queue import Queue -from message import MsgIo, Value, MessageDispatcher +from ws.client.ws_client import WsClient +from ws.message import MessageDispatcher class Window(QtWidgets.QMainWindow):