Encapsulates the connection and delivery to an MQTT broker. More...
Public Member Functions | |
| __init__ (self, broker_address, broker_port, client_id, username=None, password=None, keepalive=60, qos=0, retain=False, status_topic=None, queue_size=200, max_retries=5, initial_delay=2, max_delay=60) | |
| enqueue (self, topic, json_payload) | |
| Buffers a message for delivery. | |
| is_alive (self) | |
| Checks whether the worker thread is still running. | |
| shutdown (self) | |
| Graceful shutdown with queue drain (analogous to TelegramSender.shutdown) | |
Data Fields | |
| broker_address | |
| broker_port | |
| keepalive | |
| qos | |
| retain | |
| status_topic | |
| max_retries | |
| initial_delay | |
| max_delay | |
| client | |
Protected Member Functions | |
| _connect (self) | |
| Establishes the connection to the broker asynchronously and starts the network thread. | |
| _on_connect (self, client, userdata, flags, reason_code, properties=None) | |
| paho-mqtt callback on (re-)connect | |
| _on_disconnect (self, client, userdata, disconnect_flags, reason_code, properties=None) | |
| paho-mqtt callback on connection loss | |
| _requeue (self, topic, json_payload, retry_count) | |
| Puts a message back into the queue after a failed delivery attempt. | |
| _worker_loop (self) | |
| Processes the queue sequentially: connected -> send, otherwise buffer & wait. | |
Protected Attributes | |
| _stop_event | |
| _connected | |
| _msg_queue | |
| _worker | |
Encapsulates the connection and delivery to an MQTT broker.
Every message is handed off to an internal queue and processed by a dedicated worker thread (analogous to TelegramSender in telegram.py). If the broker is currently unreachable, messages simply stay in the queue and are delivered once the connection is restored (zero-packet-loss within the bounds of the configured queue size). If the queue is full, the OLDEST message is dropped so that new alarms are not blocked - this is logged explicitly.
| plugin.mqtt.MqttSender.__init__ | ( | self, | |
| broker_address, | |||
| broker_port, | |||
| client_id, | |||
username = None, |
|||
password = None, |
|||
keepalive = 60, |
|||
qos = 0, |
|||
retain = False, |
|||
status_topic = None, |
|||
queue_size = 200, |
|||
max_retries = 5, |
|||
initial_delay = 2, |
|||
max_delay = 60 |
|||
| ) |
| plugin.mqtt.MqttSender.enqueue | ( | self, | |
| topic, | |||
| json_payload | |||
| ) |
Buffers a message for delivery.
| topic | MQTT topic |
| json_payload | already serialized JSON string |
| plugin.mqtt.MqttSender.is_alive | ( | self | ) |
Checks whether the worker thread is still running.
| plugin.mqtt.MqttSender.shutdown | ( | self | ) |
Graceful shutdown with queue drain (analogous to TelegramSender.shutdown)
|
protected |
Establishes the connection to the broker asynchronously and starts the network thread.
|
protected |
paho-mqtt callback on (re-)connect
|
protected |
paho-mqtt callback on connection loss
|
protected |
Puts a message back into the queue after a failed delivery attempt.
|
protected |
Processes the queue sequentially: connected -> send, otherwise buffer & wait.
|
protected |
| plugin.mqtt.MqttSender.broker_address |
| plugin.mqtt.MqttSender.broker_port |
| plugin.mqtt.MqttSender.keepalive |
| plugin.mqtt.MqttSender.qos |
| plugin.mqtt.MqttSender.retain |
| plugin.mqtt.MqttSender.status_topic |
| plugin.mqtt.MqttSender.max_retries |
| plugin.mqtt.MqttSender.initial_delay |
| plugin.mqtt.MqttSender.max_delay |
|
protected |
|
protected |
| plugin.mqtt.MqttSender.client |
|
protected |