From 5c97ac4d5c438e445732cbeb89e850bccf251404 Mon Sep 17 00:00:00 2001 From: jens Date: Tue, 1 Dec 2020 12:41:01 +0100 Subject: [PATCH] - removed MsgIo::send_from_thread - created sync variant of MessageDispatcher and MsgIo - fixed delay in sending client messages --- brewpi.py | 12 ++++----- brewpi_gui.py | 40 +++++++++++++++++------------- ws/message.py | 67 +++++++++++++++++++++++++++++++++++++++++++++++++-- 3 files changed, 94 insertions(+), 25 deletions(-) diff --git a/brewpi.py b/brewpi.py index 9de1af0..2208312 100644 --- a/brewpi.py +++ b/brewpi.py @@ -65,11 +65,9 @@ class HeaterTask(ATask): self.msg_handler = msg_handler msg_handler.set_recv_handler(self.recv) self.heater = heater - self.y = 0 - self.power_soll = 0 def actor(self, y): - self.y = y + self.heater.setPower(max(0, 200 + self.heater.get_power_max() * y)) async def recv(self, data): print(data) @@ -88,8 +86,6 @@ class HeaterTask(ATask): power = Value() while True: - self.power_soll = max(0, 200 + self.heater.get_power_max()*self.y) - self.heater.setPower(self.power_soll) self.heater.process() power.set(round(self.heater.getPower())) if power.is_changed(): @@ -243,7 +239,7 @@ if __name__ == '__main__': } tc = temp_controller.TempController(tc_params) tc.set_on_changed("y_hold", heater_task.actor) - tc_task = TcTask(tc, TC_DT, dispatcher.msgio_get("Tc")) + tc_task = TcTask(tc, TC_DT, dispatcher.msgio_get("TempCtrl")) taskmgr.add(tc_task) sensor.set_on_changed("temp", ValueChanged(tc.set_theta_ist, prec=2).set) @@ -256,3 +252,7 @@ if __name__ == '__main__': while True: print("Hallo") sleep(1) + +{"+":"Pot"} +{"+":"Heater"} +{"Heater": {"Power": 2000, "Activate": 1}} diff --git a/brewpi_gui.py b/brewpi_gui.py index ae94cf7..d6b545e 100644 --- a/brewpi_gui.py +++ b/brewpi_gui.py @@ -2,7 +2,7 @@ from brewpi_win import Ui_MainWindow from PyQt5 import QtWidgets, QtGui import sys from ws.client.ws_client import WsClient -from ws.message import MessageDispatcher +from ws.message import MessageDispatcherSync from queue import Queue import asyncio from ws.connection import IConnection @@ -15,12 +15,12 @@ class Window(QtWidgets.QMainWindow, Ui_MainWindow): loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) - self.msg_dispatch = MessageDispatcher(auto_subscribe=True, loop=loop) + self.msg_dispatch = MessageDispatcherSync(auto_subscribe=True) self.ws_client = WsClient(listener=self.msg_dispatch, loop=loop) self.msg_pot = self.msg_dispatch.msgio_get('Pot') self.msg_sensor = self.msg_dispatch.msgio_get('Sensor') self.msg_heater = self.msg_dispatch.msgio_get('Heater') - self.msg_tempctrl = self.msg_dispatch.msgio_get('Tc') + self.msg_tempctrl = self.msg_dispatch.msgio_get('TempCtrl') self.msg_pot.set_recv_handler(self.on_pot_changed) self.msg_sensor.set_recv_handler(self.on_sensor_changed) @@ -35,10 +35,14 @@ class Window(QtWidgets.QMainWindow, Ui_MainWindow): self.checkbox_heater_activate.stateChanged.connect(self.on_checkbox_changed) self.send_data = Queue() - self.slider_initial_update = True + self.slider_pwr_initial_update = True + self.slider_temp_soll_initial_update = True + self.checkbox_heater_activate_initial_update = True def connect(self): - self.slider_initial_update = True + self.slider_pwr_initial_update = True + self.slider_temp_soll_initial_update = True + self.checkbox_heater_activate_initial_update = True self.ws_client.connect(uri="ws://localhost:8765") def disconnect(self): @@ -46,41 +50,43 @@ class Window(QtWidgets.QMainWindow, Ui_MainWindow): def on_slider_pwr_soll_changed(self, value): print("on_slider_pwr_soll_changed {}".format(value)) - self.msg_heater.send_from_thread({'Power': value}) + self.msg_heater.send({'Power': value}) def on_slider_temp_soll_changed(self, value): print("on_slider_temp_soll_changed {}".format(value)) - self.msg_tempctrl.send_from_thread({'Soll': value}) + self.msg_tempctrl.send({'Soll': value}) def on_checkbox_changed(self, value): print("on_checkbox_changed {}".format(value)) - self.msg_heater.send_from_thread({'Activate': int(value == 2)}) + self.msg_heater.send({'Activate': int(value == 2)}) - async def on_pot_changed(self, msg): + def on_pot_changed(self, msg): print("on_pot_changed {}".format(msg)) if "Power" in msg: self.lcdNumber_2.display(msg['Power']) - async def on_sensor_changed(self, msg): + def on_sensor_changed(self, msg): print("on_sensor_changed {}".format(msg)) if "Temp" in msg: self.lcdNumber_3.display(str(msg['Temp'])) - async def on_tempctrl_changed(self, msg): + def on_tempctrl_changed(self, msg): print("on_tempctrl_changed {}".format(msg)) -# if "Tc" in msg: -# self.lcdNumber_3.display(str(msg['Temp'])) + # if "Tempctrl" in msg: + # self.lcdNumber_3.display(str(msg['Temp'])) - async def on_heater_changed(self, msg): + def on_heater_changed(self, msg): print("on_heater_changed {}".format(msg)) for key in msg: if "Activate" in key: - self.checkbox_heater_activate.setCheckState(2 if msg['Activate'] == 1 else 0) + if self.checkbox_heater_activate_initial_update: + self.checkbox_heater_activate.setCheckState(2 if msg['Activate'] == 1 else 0) + self.checkbox_heater_activate_initial_update = False elif "Power" in key: self.lcdNumber.display(msg['Power']) - if self.slider_initial_update: + if self.slider_pwr_initial_update: self.Slider_pwr_soll.setValue(msg['Power']) - self.slider_initial_update = False + self.slider_pwr_initial_update = False elif "Capabilities" in key: submsg = msg['Capabilities'] if "Power" in submsg: diff --git a/ws/message.py b/ws/message.py index 899d1f9..db77d3d 100644 --- a/ws/message.py +++ b/ws/message.py @@ -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) +