- refactored
This commit is contained in:
@@ -1,117 +1,11 @@
|
|||||||
import asyncio
|
import asyncio
|
||||||
from ws_server import WsServer
|
|
||||||
from user import User, UserSet, update
|
|
||||||
from time import sleep
|
from time import sleep
|
||||||
import abc
|
import abc
|
||||||
import json
|
from ws.server.ws_server_multi_user import WsServerMultiUser
|
||||||
from connection import IConnection
|
|
||||||
|
|
||||||
from components.sensor.tempSensorSim import TempSensorSim, ATemperatureSensor
|
from components.sensor.tempSensorSim import TempSensorSim, ATemperatureSensor
|
||||||
from components.plant.pot import Pot, APlant
|
from components.plant.pot import Pot, APlant
|
||||||
from components.actor.heater import Heater, AHeater
|
from components.actor.heater import Heater, AHeater
|
||||||
|
from ws.message import MsgIo, MessageDispatcher, Value
|
||||||
|
|
||||||
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
|
|
||||||
|
|
||||||
|
|
||||||
class ATask:
|
class ATask:
|
||||||
@@ -248,38 +142,6 @@ class ColorHandler:
|
|||||||
await self.msg_handler.send(data)
|
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:
|
class TaskManager:
|
||||||
def __init__(self):
|
def __init__(self):
|
||||||
self.tasks: ATask = []
|
self.tasks: ATask = []
|
||||||
|
|||||||
+3
-6
@@ -1,11 +1,8 @@
|
|||||||
from brewpi_win import Ui_MainWindow
|
from brewpi_win import Ui_MainWindow
|
||||||
from PyQt5 import QtGui, QtWidgets, QtCore
|
from PyQt5 import QtWidgets
|
||||||
import sys
|
import sys
|
||||||
from ws_client import WsClient
|
from ws.client.ws_client import WsClient
|
||||||
from connection import IConnection
|
from ws.message import MessageDispatcher
|
||||||
import json
|
|
||||||
from queue import Queue
|
|
||||||
from message import MsgIo, Value, MessageDispatcher
|
|
||||||
|
|
||||||
|
|
||||||
class Window(QtWidgets.QMainWindow):
|
class Window(QtWidgets.QMainWindow):
|
||||||
|
|||||||
Reference in New Issue
Block a user