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 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 class ATask: def __init__(self, interval): self.interval = interval @abc.abstractmethod def on_process(self): pass class TempSensorTask(ATask): def __init__(self, sensor: ATemperatureSensor, interval, msg_handler: MsgIo): ATask.__init__(self, interval) self.msg_handler = msg_handler msg_handler.set_recv_handler(self.recv) self.sensor = sensor self.offset = 0 async def recv(self, data): print(data) self.offset = data['Temp'] async def send(self, data): await self.msg_handler.send(data) async def on_process(self): print("{}: Started with interval {} s".format(self.msg_handler.get_key(), self.interval)) while True: temp = self.sensor.temperature() + self.offset await self.send(temp) await asyncio.sleep(self.interval) print(temp) class HeaterTask(ATask): def __init__(self, heater: AHeater, interval, msg_handler: MsgIo): ATask.__init__(self, interval) self.msg_handler = msg_handler msg_handler.set_recv_handler(self.recv) self.heater = heater async def recv(self, data): print(data) for pair in data.items(): if 'Power' in pair[0]: self.heater.setPower(pair[1]) if 'Activate' in pair[0]: self.heater.activate(bool(pair[1])) async def send(self, data): await self.msg_handler.send(data) async def on_process(self): print("{}: Started with interval {} s".format(self.msg_handler.get_key(), self.interval)) while True: self.heater.process() power = self.heater.getPower() await self.send({'Power': power}) await asyncio.sleep(self.interval) class PotTask(ATask): def __init__(self, pot: APlant, interval, msg_handler: MsgIo): ATask.__init__(self, interval) self.msg_handler = msg_handler msg_handler.set_recv_handler(self.recv) self.plant = pot async def recv(self, data): print(data) for key in data: print(key) async def send(self, data): await self.msg_handler.send(data) async def on_process(self): print("{}: Started with interval {} s".format(self.msg_handler.get_key(), self.interval)) plant_temp = Value(-1) plant_power = Value() while True: self.plant.process() plant_power.set(round(self.plant.getPower(), 1)) plant_temp.set(round(self.plant.getTemperature(), 1)) if plant_power.is_changed(): print("Plant Power {}".format(plant_power.get())) await self.send({'Power': plant_power.get()}) if plant_temp.is_changed(): print("Plant Temp {}".format(plant_temp.get())) await self.send({'Temp': plant_temp.get()}) await asyncio.sleep(self.interval) class CounterTask(ATask): def __init__(self, interval, msg_handler: MsgIo): ATask.__init__(self, interval) self.msg_handler = msg_handler msg_handler.set_recv_handler(self.recv) async def recv(self, data): print(data) async def send(self, data): await self.msg_handler.send(data) async def on_process(self): print("{}: Started with interval {} s".format(self.msg_handler.get_key(), self.interval)) while True: for count in range(0, 101): await self.send(count) await asyncio.sleep(self.interval) class ColorHandler: def __init__(self, msg_handler: MsgIo): self.msg_handler = msg_handler msg_handler.set_recv_handler(self.recv) print("{}: Constructed".format(self.msg_handler.get_key())) async def recv(self, data): print(data) for value in data.values(): await self.send(value) async def send(self, 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: def __init__(self): self.tasks: ATask = [] def add(self, task: ATask): self.tasks.append(task) def start(self): funcs = [] for task in self.tasks: funcs.append(task.on_process()) return asyncio.gather(*funcs) if __name__ == '__main__': if 0: taskmgr = PeriodicUpdaterThreaded() taskmgr.start("localhost", 8765) bg_thread = Thread(target=taskmgr.loop.run_forever) bg_thread.setDaemon(True) bg_thread.start() else: dispatcher = MessageDispatcher() server = WsServerMultiUser(listener=dispatcher) taskmgr = TaskManager() taskmgr.add(CounterTask(1.0, dispatcher.msgio_get("a"))) taskmgr.add(CounterTask(0.5, dispatcher.msgio_get("b"))) taskmgr.add(CounterTask(0.1, dispatcher.msgio_get("c"))) # Sensor sensor = TempSensorSim() taskmgr.add(TempSensorTask(sensor, 1.0, dispatcher.msgio_get("Sensor"))) pot_params = { "dt" : 1.0, "theta_amb" : 15, "C" : 4190, "M" : 20, "L" : 0.1, "Td" : 12, "kn" : 0.2 } # {"Heater": {"Activate": 1, "Power": 1000}} # Plant pot = Pot(pot_params) # pot.connect_theta_changed(sensor.set_fake_temp) taskmgr.add(PotTask(pot, 1.0, dispatcher.msgio_get("Pot"))) # Heater heater = Heater() heater.connect_power_changed(pot.setPower) taskmgr.add(HeaterTask(heater, 1.0, dispatcher.msgio_get("Heater"))) ColorHandler(dispatcher.msgio_get("Color")) h_dispatcher = taskmgr.start() h_server = server.listen("localhost", 8765) asyncio.gather(h_dispatcher, h_server) server.loop.run_forever() while True: print("Hallo") sleep(1)