diff --git a/res/beer.png b/res/beer.png new file mode 100644 index 0000000..fc94e7b Binary files /dev/null and b/res/beer.png differ diff --git a/res/beer.xcf b/res/beer.xcf new file mode 100644 index 0000000..038f1b0 Binary files /dev/null and b/res/beer.xcf differ diff --git a/res/dose.png b/res/dose.png new file mode 100644 index 0000000..ec6b378 Binary files /dev/null and b/res/dose.png differ diff --git a/res/fass.png b/res/fass.png new file mode 100644 index 0000000..f1b83c6 Binary files /dev/null and b/res/fass.png differ diff --git a/res/flasche.png b/res/flasche.png new file mode 100644 index 0000000..552d049 Binary files /dev/null and b/res/flasche.png differ diff --git a/res/glass.png b/res/glass.png new file mode 100644 index 0000000..ecbadad Binary files /dev/null and b/res/glass.png differ diff --git a/res/glass2.png b/res/glass2.png new file mode 100644 index 0000000..0ba102f Binary files /dev/null and b/res/glass2.png differ diff --git a/res/glass_flasche.png b/res/glass_flasche.png new file mode 100644 index 0000000..7b6f576 Binary files /dev/null and b/res/glass_flasche.png differ diff --git a/res/krug.png b/res/krug.png new file mode 100644 index 0000000..706a69e Binary files /dev/null and b/res/krug.png differ diff --git a/res/krug2.png b/res/krug2.png new file mode 100644 index 0000000..a4301d0 Binary files /dev/null and b/res/krug2.png differ diff --git a/ws/client/ws_client.py b/ws/client/ws_client.py new file mode 100644 index 0000000..e1fe40b --- /dev/null +++ b/ws/client/ws_client.py @@ -0,0 +1,132 @@ +import asyncio +import websockets +from threading import Thread +from queue import Queue +from time import sleep +import json +from ws.connection import IConnection + + +class WsClient: + def __init__(self, listener: IConnection): + self.listener = listener + self.loop = asyncio.new_event_loop() + self.stop = None + self.bg_thread = None + + def bg_run(self, uri): + print("bg_run: started") + asyncio.set_event_loop(self.loop) + fut = asyncio.ensure_future(self.run_client(uri), loop=self.loop) + self.loop.run_until_complete(fut) + self.stop = None + print("bg_run: terminated") + + def connect(self, uri): + if self.bg_thread is not None and self.bg_thread.is_alive(): + return + + self.stop: asyncio.Future[None] = self.loop.create_future() + self.bg_thread = Thread(target=self.bg_run, args=(uri,)) + self.bg_thread.setDaemon(True) + self.bg_thread.start() + + def disconnect(self): + if self.stop is not None: + self.loop.call_soon_threadsafe(self.stop.set_result, None) + self.bg_thread.join() + + async def run_client(self, uri): + try: + websocket = await websockets.connect(uri, loop=self.loop) + except Exception as exc: + print(f"Failed to connect to {uri}: {exc}.") + return + + print("handler: got connection to {}".format(websocket.remote_address)) + self.listener.on_connect() + try: + path = "/" + consumer_task = asyncio.ensure_future(self.handler_recv(websocket, path), loop=self.loop) + producer_task = asyncio.ensure_future(self.handler_send(websocket, path), loop=self.loop) + done, pending = await asyncio.wait([consumer_task, producer_task, self.stop], return_when=asyncio.FIRST_COMPLETED,) + for task in pending: + task.cancel() + finally: + print("handler: lost connection from {}".format(websocket.remote_address)) + await websocket.close() + self.listener.on_disconnect() + + async def handler_recv(self, websocket, path): + while True: + try: + data = await websocket.recv() + await self.listener.on_recv(data) + except: + print("handler_recv: lost connection from {}".format(websocket.remote_address)) + return + + async def handler_send(self, websocket, path): + while True: + data = await self.listener.on_send() + if data is not None: + try: + await websocket.send(json.dumps(data)) + except: + print("handler_send: lost connection from {}".format(websocket.remote_address)) + break + + +class PeriodicUpdater(IConnection): + def __init__(self): + self.ws_client = WsClient(self) + + self.state = Queue() + self.thread_cancel = True + self.thread = None + + def connect(self, uri): + self.ws_client.connect(uri) + + def disconnect(self): + self.ws_client.disconnect() + + def on_connect(self): + self.thread_cancel = False + self.thread = Thread(target=self.run) + self.thread.setDaemon(True) + self.thread.start() + + def on_disconnect(self): + self.thread_cancel = True + self.thread.join(1) + + async def on_send(self): + return await asyncio.get_event_loop().run_in_executor(None, self.state.get) + + async def on_recv(self, data): + print("R: {}".format(data)) + + def send(self, data): + self.state.put_nowait(data) + + def run(self): + count = 0 + while not self.thread_cancel: + self.send({"Value": count}) + count += 1 + sleep(0.5) + + +if __name__ == '__main__': + + ws = PeriodicUpdater() + ws.connect(uri="ws://localhost:8765") + + sleep(5) + ws.disconnect() + sleep(1) + ws.connect(uri="ws://localhost:8765") + + while True: + sleep(1) \ No newline at end of file diff --git a/ws/connection.py b/ws/connection.py new file mode 100644 index 0000000..062aaa4 --- /dev/null +++ b/ws/connection.py @@ -0,0 +1,19 @@ +import abc + + +class IConnection: + @abc.abstractmethod + def on_connect(self): + pass + + @abc.abstractmethod + def on_disconnect(self): + pass + + @abc.abstractmethod + def on_recv(self, data): + pass + + @abc.abstractmethod + def on_send(self): + return None diff --git a/ws/message.py b/ws/message.py new file mode 100644 index 0000000..3e21dad --- /dev/null +++ b/ws/message.py @@ -0,0 +1,77 @@ +import asyncio +from ws.connection import IConnection +import json + + +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 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) + diff --git a/ws/server/ws_server.py b/ws/server/ws_server.py new file mode 100644 index 0000000..bdf7cbd --- /dev/null +++ b/ws/server/ws_server.py @@ -0,0 +1,39 @@ +import asyncio +import websockets +import abc + + +class WsServer: + def __init__(self, loop=None): + if loop is None: + self.loop = asyncio.get_event_loop() + else: + self.loop = loop + + async def listen(self, host, port): + return await websockets.serve(self.handler, host, port, loop=self.loop) + + def disconnect(self): + pass + + async def handler(self, websocket, path): + print("handler: got connection from {}".format(websocket.remote_address)) + await self.register(websocket) + try: + consumer_task = asyncio.ensure_future(self.handler_recv(websocket, path), loop=self.loop) + producer_task = asyncio.ensure_future(self.handler_send(websocket, path), loop=self.loop) + done, pending = await asyncio.wait([consumer_task, producer_task], return_when=asyncio.FIRST_COMPLETED,) + for task in pending: + task.cancel() + finally: + print("handler: lost connection from {}".format(websocket.remote_address)) + await self.unregister(websocket) + + @abc.abstractmethod + async def handler_recv(self, websocket, path): + pass + + @abc.abstractmethod + async def handler_send(self, websocket, path): + pass + diff --git a/ws/server/ws_server_multi_user.py b/ws/server/ws_server_multi_user.py new file mode 100644 index 0000000..c7d5f20 --- /dev/null +++ b/ws/server/ws_server_multi_user.py @@ -0,0 +1,67 @@ +import asyncio +from ws.server.ws_server import WsServer +from ws.user import User, UserSet, update +import abc +from ws.connection import IConnection + + +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 diff --git a/ws/user.py b/ws/user.py new file mode 100644 index 0000000..046a892 --- /dev/null +++ b/ws/user.py @@ -0,0 +1,132 @@ +import dpath.util as dp +import json + + +def update(d, u, only_existing=True): + for k, v in u.items(): + if isinstance(v, dict): + if only_existing: + r = update(d.get(k, {}), v) + d[k] = r + else: + d[k] = v + else: + d[k] = u[k] + return d + + +class User: + def __init__(self, socket): + self.socket = socket + self.paths = [] + + async def send(self, data): + result = {} + for path in self.paths: + msg = dp.search(data, path) + result.update(msg) + if len(result) != 0: + await self.socket.send(json.dumps(result)) + + async def process(self, data): + handled = True + try: + # Handle message subscription + d = json.loads(data) + if "+" in d: + path = d['+'] + print("Received subscribe {}".format(path)) + self.path_add(path) + elif "-" in d: + path = d['-'] + print("Received unsubscribe {}".format(path)) + self.path_remove(path) + elif "?" in d: + print("Subscriptions: {}".format(self.paths)) + await self.socket.send(json.dumps({"?": self.paths})) + else: + handled = False + except Exception as e: + print(e) + + return handled + + def path_add(self, path): + if path not in self.paths: + self.paths.append(path) + + def path_remove(self, path): + if path in self.paths: + self.paths.remove(path) + + +class UserSet: + def __init__(self): + self.users = {} + self.keys = None + + def __len__(self): + return len(self.users) + + def clear(self): + self.users = {} + self.keys = None + + def add(self, key, user): + self.users[key] = user + + def remove(self, key): + self.users.pop(key, None) + + def get(self, key): + return self.users[key] + + def __iter__(self): + self.keys = iter(self.users.keys()) + return self + + def __next__(self): + return self.users[next(self.keys)] + + +def check_fill(u): + if u: + print("Is occupied") + else: + print("Is empty") + + +if __name__ == '__main__': + userset = UserSet() + check_fill(userset) + + userset.add(1, "A") + userset.add(2, "B") + userset.add(3, "C") + check_fill(userset) + + print("Len(userset) = {}".format(userset)) + for i in userset: + print(i) + + userset.remove(3) + for i in userset: + print(i) + + userset.remove(2) + for i in userset: + print(i) + + userset.remove(1) + for i in userset: + print(i) + + a = {'Pot': {'Temp': 15.5}} + b = {'Pot': {'Power': 2400}} + + print(a) + print(b) + + a = update(a, b) + + print("a | b = {}".format(a)) diff --git a/ws_client_gui.py b/ws_client_gui.py new file mode 100644 index 0000000..13edf0c --- /dev/null +++ b/ws_client_gui.py @@ -0,0 +1,251 @@ +# This is a sample Python script. +import asyncio +import sys +from PyQt5 import QtGui, QtWidgets, QtCore +from ws.client.ws_client import WsClient +from ws.connection import IConnection +import json +from queue import Queue + + +class Window(QtWidgets.QMainWindow, IConnection): + + # Use signals and slots to avoid crashes + # during updating gui elements from websocket context + upd_progressbar_signal = QtCore.pyqtSignal(int, int) + upd_color_signal = QtCore.pyqtSignal(str) + + def __init__(self): + QtWidgets.QMainWindow.__init__(self) + self.ws_client = WsClient(self) + self.is_exit_gracefully = True + self.textEdit = None + self.progressbar = [] + self.styleChoice = None + + self.setGeometry(50, 50, 500, 300) + self.setWindowTitle("Hallo Mausi!") + self.setWindowIcon(QtGui.QIcon('res/flasche.png')) + + self.menubar_create() + self.statusbar_create() + self.toolbar_create() + self.buttons_create() + self.checkbox_create() + + self.progressbar_create() + self.style_chooser_create() + self.calendar_create() + + self.send_data = Queue() + self.show() + + self.upd_progressbar_signal.connect(self.upd_progressbar_slot) + self.upd_color_signal.connect(self.upd_color_slot) + + def calendar_create(self): + cal = QtWidgets.QCalendarWidget(self) + cal.move(500, 200) + cal.resize(250, 250) + + def menubar_create(self): + mainMenu = self.menuBar() + fileMenu = mainMenu.addMenu("&File") + + actionExit = QtWidgets.QAction("&Exit", self) + actionExit.setShortcut("Ctrl+Q") + actionExit.setStatusTip("Leave the App") + actionExit.triggered.connect(self.on_btn_exit_clicked) + + actionEditor = QtWidgets.QAction("&Editor", self) + actionEditor.setShortcut("Ctrl+E") + actionEditor.setStatusTip("Open Editor") + actionEditor.triggered.connect(self.on_action_editor) + + actionFileOpener = QtWidgets.QAction("File &Open", self) + actionFileOpener.setShortcut("Ctrl+O") + actionFileOpener.setStatusTip("Open a file") + actionFileOpener.triggered.connect(self.on_action_file_open) + + actionFileSaver = QtWidgets.QAction("File &Save", self) + actionFileSaver.setShortcut("Ctrl+S") + actionFileSaver.setStatusTip("Save a file") + actionFileSaver.triggered.connect(self.on_action_file_save) + + fileMenu.addAction(actionExit) + fileMenu.addAction(actionEditor) + fileMenu.addAction(actionFileOpener) + fileMenu.addAction(actionFileSaver) + + def statusbar_create(self): + statusBar = self.statusBar() + + def toolbar_create(self): + def action_caller(a, n): + return lambda action: self.on_action_triggered(a, n) + + action1 = QtWidgets.QAction(QtGui.QIcon('res/glass.png'), "Drink a beer!", self) + action1.triggered.connect(action_caller(action1, 1)) + + action2 = QtWidgets.QAction(QtGui.QIcon('res/glass2.png'), "Drink more beer!", self) + action2.triggered.connect(action_caller(action2, 2)) + + action3 = QtWidgets.QAction(QtGui.QIcon('res/glass_flasche.png'), "Choose font!", self) + action3.triggered.connect(self.on_action_font_triggered) + + action4 = QtWidgets.QAction(QtGui.QIcon('res/dose.png'), "Choose color!", self) + action4.triggered.connect(self.on_action_color_triggered) + + toolbar = self.addToolBar("Toolbar") + + toolbar.addAction(action1) + toolbar.addAction(action2) + toolbar.addAction(action3) + toolbar.addAction(action4) + + def buttons_create(self): + btn = QtWidgets.QPushButton("Quit", self) + btn.clicked.connect(self.on_btn_exit_clicked) + btn.resize(btn.minimumSizeHint()) + btn.move(0, 100) + + def checkbox_create(self): + checkbox = QtWidgets.QCheckBox('Exit gracefully!', self) + checkbox.stateChanged.connect(self.on_checkbox_changed) + checkbox.setCheckState(2 if self.is_exit_gracefully else 0) + checkbox.move(0, 150) + + def progressbar_create(self): + for i in range(0, 3): + progressbar = QtWidgets.QProgressBar(self) + progressbar.setGeometry(200, 80+25*i, 250, 20) + self.progressbar.append(progressbar) + + btn = QtWidgets.QPushButton("Download", self) + btn.clicked.connect(self.on_btn_download_clicked) + btn.move(0, 70) + + def style_chooser_create(self): + print(self.style().objectName()) + self.styleChoice = QtWidgets.QLabel("Linux", self) + comboxBox = QtWidgets.QComboBox(self) + comboxBox.addItem("motif") + comboxBox.addItem("Windows") + comboxBox.addItem("cde") + comboxBox.addItem("Plastique") + comboxBox.addItem("Cleanlooks") + comboxBox.addItem("windowsvista") + comboxBox.addItem("fusion") + comboxBox.move(50, 250) + self.styleChoice.move(150, 150) + comboxBox.activated[str].connect(self.on_combobox_changed) + + def on_btn_exit_clicked(self): + print('PyQt5 on_btn_click') + self.close() + + def on_action_triggered(self, action, n): + print('PyQt5 on_action_triggered {}'.format(n)) + if n == 1: + self.ws_client.connect(uri="ws://localhost:8765") + if n == 2: + self.ws_client.disconnect() + + def on_action_editor(self): + self.textEdit = QtWidgets.QTextEdit() + self.setCentralWidget(self.textEdit) + + def on_action_file_open(self): + name, _ = QtWidgets.QFileDialog.getOpenFileName(self, 'Open file') + file = open(name, 'r') + + self.on_action_editor() + with file: + text = file.read() + self.textEdit.setText(text) + + def on_action_file_save(self): + + if self.textEdit is not None: + name, _ = QtWidgets.QFileDialog.getSaveFileName(self, 'Save file') + file = open(name, 'w') + text = self.textEdit.toPlainText() + file.write(text) + file.close() + else: + print("No file to save!") + + def on_action_font_triggered(self): + font, valid = QtWidgets.QFontDialog.getFont() + if valid: + self.styleChoice.setFont(font) + + def on_action_color_triggered(self): + color = QtWidgets.QColorDialog.getColor() + self.send_data.put({'Color': color.name()}) + + def on_checkbox_changed(self, state): + print('PyQt5 on_checkbox_changed, state={}'.format(state)) + if state == QtCore.Qt.Checked: + self.is_exit_gracefully = True + else: + self.is_exit_gracefully = False + + def on_btn_download_clicked(self): + print("Nothing happens!") + + def on_combobox_changed(self, text): + self.styleChoice.setText(text) + QtWidgets.QApplication.setStyle(QtWidgets.QStyleFactory.create(text)) + + def on_connect(self): + self.send_data.put_nowait({'+':'a'}) + self.send_data.put_nowait({'+':'b'}) + self.send_data.put_nowait({'+':'c'}) + self.send_data.put_nowait({'+':'Color'}) + + def on_disconnect(self): + pass + + @QtCore.pyqtSlot(int, int) + def upd_progressbar_slot(self, n, v): + self.progressbar[n].setValue(v) + + @QtCore.pyqtSlot(str) + def upd_color_slot(self, color_str): + self.styleChoice.setStyleSheet("QWidget {background-color: %s}" % color_str) + + async def on_recv(self, data): + try: + d = json.loads(data) + keys = {'a': 0, 'b': 1, 'c': 2} + color = {'Color': 0} + for cmd in d: + if cmd in keys.keys(): + await asyncio.get_event_loop().run_in_executor(None, self.upd_progressbar_signal.emit(keys[cmd], int(d[cmd] + 0.5))) + + if cmd in color.keys(): + await asyncio.get_event_loop().run_in_executor(None, self.upd_color_signal.emit(d[cmd])) + except: + pass + + async def on_send(self): + try: + data = await asyncio.get_event_loop().run_in_executor(None, self.send_data.get, True, 1) + return data + except: + return None + + def closeEvent(self, event): + if self.is_exit_gracefully: + print("Exitting gracefully!") + self.ws_client.disconnect() + + +if __name__ == '__main__': + app = QtWidgets.QApplication(sys.argv) + + gui = Window() + sys.exit(app.exec_()) + +# See PyCharm help at https://www.jetbrains.com/help/pycharm/