Try to stop properly the threads with SIGTERM and SIGINT
This commit is contained in:
+26
-10
@@ -3,17 +3,18 @@ import time
|
|||||||
import threading
|
import threading
|
||||||
import os
|
import os
|
||||||
import json
|
import json
|
||||||
|
import signal
|
||||||
|
|
||||||
#mine
|
#mine
|
||||||
import mqtt
|
import mqtt
|
||||||
import yamlparser
|
import yamlparser
|
||||||
from xiaomihub import XiaomiHub
|
from xiaomihub import XiaomiHub
|
||||||
|
|
||||||
logging.basicConfig(level=logging.INFO)
|
logging.basicConfig(level=logging.DEBUG)
|
||||||
_LOGGER = logging.getLogger(__name__)
|
_LOGGER = logging.getLogger(__name__)
|
||||||
|
|
||||||
def process_gateway_messages(gateway, client):
|
def process_gateway_messages(gateway, client, stop_event):
|
||||||
while True:
|
while not stop_event.is_set():
|
||||||
try:
|
try:
|
||||||
packet = gateway._queue.get()
|
packet = gateway._queue.get()
|
||||||
_LOGGER.debug("data from queuee: " + format(packet))
|
_LOGGER.debug("data from queuee: " + format(packet))
|
||||||
@@ -28,10 +29,11 @@ def process_gateway_messages(gateway, client):
|
|||||||
gateway._queue.task_done()
|
gateway._queue.task_done()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
_LOGGER.error('Error while sending from gateway to mqtt: ', str(e))
|
_LOGGER.error('Error while sending from gateway to mqtt: ', str(e))
|
||||||
|
_LOGGER.info("Stopping Gateway Thread ...")
|
||||||
|
|
||||||
def read_motion_data(gateway, client, polling_interval, polling_models):
|
def read_motion_data(gateway, client, polling_interval, polling_models, stop_event):
|
||||||
first = True
|
first = True
|
||||||
while True:
|
while not stop_event.is_set():
|
||||||
try:
|
try:
|
||||||
for device_type in gateway.XIAOMI_DEVICES:
|
for device_type in gateway.XIAOMI_DEVICES:
|
||||||
devices = gateway.XIAOMI_DEVICES[device_type]
|
devices = gateway.XIAOMI_DEVICES[device_type]
|
||||||
@@ -59,9 +61,10 @@ def read_motion_data(gateway, client, polling_interval, polling_models):
|
|||||||
except Exception as e:
|
except Exception as e:
|
||||||
_LOGGER.error('Error while sending from mqtt to gateway: ', str(e))
|
_LOGGER.error('Error while sending from mqtt to gateway: ', str(e))
|
||||||
time.sleep(polling_interval)
|
time.sleep(polling_interval)
|
||||||
|
_LOGGER.info("Stopping Polling Thread ...")
|
||||||
|
|
||||||
def process_mqtt_messages(gateway, client):
|
def process_mqtt_messages(gateway, client, stop_event):
|
||||||
while True:
|
while not stop_event.is_set():
|
||||||
try:
|
try:
|
||||||
data = client._queue.get()
|
data = client._queue.get()
|
||||||
_LOGGER.debug("data from mqtt: " + format(data))
|
_LOGGER.debug("data from mqtt: " + format(data))
|
||||||
@@ -73,6 +76,13 @@ def process_mqtt_messages(gateway, client):
|
|||||||
client._queue.task_done()
|
client._queue.task_done()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
_LOGGER.error('Error while sending from mqtt to gateway: ', str(e))
|
_LOGGER.error('Error while sending from mqtt to gateway: ', str(e))
|
||||||
|
_LOGGER.info("Stopping MQTT Thread ...")
|
||||||
|
|
||||||
|
def exit_handler(signal, frame):
|
||||||
|
print('Exiting')
|
||||||
|
stop_event.set()
|
||||||
|
gateway.stop()
|
||||||
|
client.disconnect()
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
_LOGGER.info("Loading config file...")
|
_LOGGER.info("Loading config file...")
|
||||||
@@ -81,6 +91,9 @@ if __name__ == "__main__":
|
|||||||
polling_interval = config['gateway'].get("polling_interval", 2)
|
polling_interval = config['gateway'].get("polling_interval", 2)
|
||||||
polling_models = config['gateway'].get("polling_models", ['motion'])
|
polling_models = config['gateway'].get("polling_models", ['motion'])
|
||||||
|
|
||||||
|
signal.signal(signal.SIGINT, exit_handler)
|
||||||
|
signal.signal(signal.SIGTERM, exit_handler)
|
||||||
|
|
||||||
_LOGGER.info("Init mqtt client.")
|
_LOGGER.info("Init mqtt client.")
|
||||||
client = mqtt.Mqtt(config)
|
client = mqtt.Mqtt(config)
|
||||||
client.connect()
|
client.connect()
|
||||||
@@ -90,17 +103,20 @@ if __name__ == "__main__":
|
|||||||
client.subscribe("plug", "+", "status", "set")
|
client.subscribe("plug", "+", "status", "set")
|
||||||
|
|
||||||
gateway = XiaomiHub(gateway_pass)
|
gateway = XiaomiHub(gateway_pass)
|
||||||
t1 = threading.Thread(target=process_gateway_messages, args=[gateway, client])
|
stop_event= threading.Event()
|
||||||
|
t1 = threading.Thread(target=process_gateway_messages, args=[gateway, client, stop_event])
|
||||||
t1.daemon = True
|
t1.daemon = True
|
||||||
t1.start()
|
t1.start()
|
||||||
|
|
||||||
t2 = threading.Thread(target=process_mqtt_messages, args=[gateway, client])
|
t2 = threading.Thread(target=process_mqtt_messages, args=[gateway, client, stop_event])
|
||||||
t2.daemon = True
|
t2.daemon = True
|
||||||
t2.start()
|
t2.start()
|
||||||
|
|
||||||
t3 = threading.Thread(target=read_motion_data, args=[gateway, client, polling_interval, polling_models])
|
t3 = threading.Thread(target=read_motion_data, args=[gateway, client, polling_interval, polling_models, stop_event])
|
||||||
t3.daemon = True
|
t3.daemon = True
|
||||||
t3.start()
|
t3.start()
|
||||||
|
|
||||||
while True:
|
while True:
|
||||||
|
if stop_event.is_set():
|
||||||
|
break
|
||||||
time.sleep(10)
|
time.sleep(10)
|
||||||
|
|||||||
@@ -55,6 +55,9 @@ class Mqtt:
|
|||||||
t1.start()
|
t1.start()
|
||||||
self._threads.append(t1)
|
self._threads.append(t1)
|
||||||
|
|
||||||
|
def disconnect(self):
|
||||||
|
self._client.disconnect()
|
||||||
|
|
||||||
def subscribe(self, model="+", name="+", prop="+", command=None):
|
def subscribe(self, model="+", name="+", prop="+", command=None):
|
||||||
topic = self.prefix + "/" + model + "/" + name + "/" + prop
|
topic = self.prefix + "/" + model + "/" + name + "/" + prop
|
||||||
if command is not None:
|
if command is not None:
|
||||||
|
|||||||
Reference in New Issue
Block a user