diff --git a/brewpi.py b/brewpi.py index e7eb820..b4348ed 100644 --- a/brewpi.py +++ b/brewpi.py @@ -200,7 +200,7 @@ if __name__ == '__main__': h_dispatcher = taskmgr.start() h_server = server.listen("localhost", 8765) asyncio.gather(h_dispatcher, h_server) - server.loop.run_forever() + server.run_forever() while True: print("Hallo") diff --git a/brewpi_gui.py b/brewpi_gui.py index 72c74c5..9f45d32 100644 --- a/brewpi_gui.py +++ b/brewpi_gui.py @@ -11,15 +11,16 @@ from ws.connection import IConnection class Window(QtWidgets.QMainWindow, Ui_MainWindow): def __init__(self): QtWidgets.QMainWindow.__init__(self) - self.msg_dispatch = MessageDispatcher() - self.ws_client = WsClient(listener=self.msg_dispatch) + loop = asyncio.new_event_loop() + self.msg_dispatch = MessageDispatcher(auto_subscribe=True, loop=loop) + 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_pot.set_recv_handler(self.on_pot_changed) - self.msg_sensor.set_recv_handler(self.on_sensor_changed) - self.msg_heater.set_recv_handler(self.on_heater_changed) +# self.msg_pot.set_recv_handler(self.on_pot_changed) +# self.msg_sensor.set_recv_handler(self.on_sensor_changed) +# self.msg_heater.set_recv_handler(self.on_heater_changed) self.setupUi(self) self.actionStart.triggered.connect(self.connect) @@ -41,17 +42,6 @@ class Window(QtWidgets.QMainWindow, Ui_MainWindow): def on_heater_changed(self, msg): print("on_heater_changed") - def on_connect(self): - self.send_data.put_nowait({'+': 'Pot'}) - self.send_data.put_nowait({'+': 'Sensor'}) - self.send_data.put_nowait({'+': 'Heater'}) - - def on_disconnect(self): - pass - - async def on_recv(self, data): - pass - async def on_send(self): try: data = await asyncio.get_event_loop().run_in_executor(None, self.send_data.get, True, 1) diff --git a/ws/client/ws_client.py b/ws/client/ws_client.py index 72052c3..40876fa 100644 --- a/ws/client/ws_client.py +++ b/ws/client/ws_client.py @@ -8,9 +8,9 @@ from ws.connection import IConnection class WsClient: - def __init__(self, listener: IConnection): + def __init__(self, listener: IConnection, loop=None): self.listener = listener - self.loop = asyncio.new_event_loop() + self.loop = loop self.stop = None self.bg_thread = None @@ -44,7 +44,7 @@ class WsClient: return print("handler: got connection to {}".format(websocket.remote_address)) - self.listener.on_connect(self.loop) + await self.listener.on_connect() try: path = "/" consumer_task = asyncio.ensure_future(self.handler_recv(websocket, path), loop=self.loop) @@ -55,7 +55,7 @@ class WsClient: finally: print("handler: lost connection from {}".format(websocket.remote_address)) await websocket.close() - self.listener.on_disconnect() + await self.listener.on_disconnect() async def handler_recv(self, websocket, path): while True: diff --git a/ws/message.py b/ws/message.py index 49a4ea8..60ca013 100644 --- a/ws/message.py +++ b/ws/message.py @@ -45,9 +45,11 @@ class Value: class MessageDispatcher(IConnection): - def __init__(self): + 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 def msgio_get(self, key): obj = MsgIo(key, self.send) @@ -56,10 +58,14 @@ class MessageDispatcher(IConnection): return obj return None - def on_connect(self, loop): - self.state = asyncio.Queue(loop=loop) + 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}) - def on_disconnect(self): + async def on_disconnect(self): pass async def on_recv(self, data): diff --git a/ws/server/ws_server.py b/ws/server/ws_server.py index bdf7cbd..5db8600 100644 --- a/ws/server/ws_server.py +++ b/ws/server/ws_server.py @@ -5,10 +5,13 @@ import abc class WsServer: def __init__(self, loop=None): - if loop is None: - self.loop = asyncio.get_event_loop() + self.loop = loop + + def run_forever(self): + if self.loop is None: + asyncio.get_event_loop().run_forever() else: - self.loop = loop + self.loop.run_forever() async def listen(self, host, port): return await websockets.serve(self.handler, host, port, loop=self.loop) diff --git a/ws/server/ws_server_multi_user.py b/ws/server/ws_server_multi_user.py index 51626f7..c7d5f20 100644 --- a/ws/server/ws_server_multi_user.py +++ b/ws/server/ws_server_multi_user.py @@ -21,7 +21,6 @@ class WsServerMultiUser(WsServer): self.listener = listener self.USERS = UserSet() self.global_state = {} - self.listener.on_connect(self.loop) async def notify_state(self, data): if self.USERS: # asyncio.wait doesn't accept an empty list