reworked messagereceiver module
* use bumblebee's internal threading capabilities * various small code improvements (pylint)
This commit is contained in:
parent
01cde70e14
commit
a4a622252b
1 changed files with 29 additions and 40 deletions
|
@ -18,29 +18,36 @@ Example:
|
||||||
|
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
import socket
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
import json
|
||||||
|
|
||||||
import core.module
|
import core.module
|
||||||
import core.widget
|
import core.widget
|
||||||
import core.input
|
import core.input
|
||||||
|
|
||||||
import socket
|
|
||||||
import threading
|
|
||||||
import logging
|
|
||||||
import os
|
|
||||||
import queue
|
|
||||||
import json
|
|
||||||
|
|
||||||
|
class Module(core.module.Module):
|
||||||
|
@core.decorators.never
|
||||||
|
def __init__(self, config, theme):
|
||||||
|
super().__init__(config, theme, core.widget.Widget(self.message))
|
||||||
|
|
||||||
class Worker(threading.Thread):
|
self.background = True
|
||||||
def __init__(self, unix_socket_address, queue):
|
|
||||||
threading.Thread.__init__(self)
|
|
||||||
self.__unix_socket_address = unix_socket_address
|
|
||||||
self.__queue = queue
|
|
||||||
|
|
||||||
def run(self):
|
self.__unix_socket_address = self.parameter("address", "")
|
||||||
|
|
||||||
|
self.__message = ""
|
||||||
|
self.__state = []
|
||||||
|
|
||||||
|
def message(self, widget):
|
||||||
|
return self.__message
|
||||||
|
|
||||||
|
def __read_data_from_socket(self):
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
os.unlink(self.__unix_socket_address)
|
os.unlink(self.__unix_socket_address)
|
||||||
except OSError as e:
|
except OSError:
|
||||||
if os.path.exists(self.__unix_socket_address):
|
if os.path.exists(self.__unix_socket_address):
|
||||||
logging.exception(
|
logging.exception(
|
||||||
"Couldn't bind to unix socket %s", self.__unix_socket_address
|
"Couldn't bind to unix socket %s", self.__unix_socket_address
|
||||||
|
@ -57,37 +64,19 @@ class Worker(threading.Thread):
|
||||||
data = conn.recv(1024)
|
data = conn.recv(1024)
|
||||||
if not data:
|
if not data:
|
||||||
break
|
break
|
||||||
self.__queue.put(data.decode("utf-8"))
|
yield data.decode("utf-8")
|
||||||
|
|
||||||
|
|
||||||
class Module(core.module.Module):
|
|
||||||
@core.decorators.every(seconds=1)
|
|
||||||
def __init__(self, config, theme):
|
|
||||||
super().__init__(config, theme, core.widget.Widget(self.message))
|
|
||||||
|
|
||||||
self.__unix_socket_address = self.parameter("address", "")
|
|
||||||
|
|
||||||
self.__message = ""
|
|
||||||
self.__state = []
|
|
||||||
|
|
||||||
self.__queue = queue.Queue()
|
|
||||||
self.__worker = Worker(self.__unix_socket_address, self.__queue)
|
|
||||||
self.__worker.daemon = True
|
|
||||||
self.__worker.start()
|
|
||||||
|
|
||||||
def message(self, widget):
|
|
||||||
return self.__message
|
|
||||||
|
|
||||||
def update(self):
|
def update(self):
|
||||||
try:
|
try:
|
||||||
received_data = self.__queue.get(block=False)
|
for received_data in self.__read_data_from_socket():
|
||||||
parsed_data = json.loads(received_data)
|
parsed_data = json.loads(received_data)
|
||||||
self.__message = parsed_data["message"]
|
self.__message = parsed_data["message"]
|
||||||
self.__state = parsed_data["state"]
|
self.__state = parsed_data["state"]
|
||||||
except json.JSONDecodeError as e:
|
core.event.trigger("update", [self.id], redraw_only=True)
|
||||||
|
except json.JSONDecodeError:
|
||||||
logging.exception("Couldn't parse message")
|
logging.exception("Couldn't parse message")
|
||||||
except queue.Empty as e:
|
except Exception:
|
||||||
pass
|
logging.exception("Unexpected exception while reading from socket")
|
||||||
|
|
||||||
def state(self, widget):
|
def state(self, widget):
|
||||||
return self.__state
|
return self.__state
|
||||||
|
|
Loading…
Reference in a new issue