BOSWatch 3
Python Script to receive and decode German BOS Information with rtl_fm and multimon-NG
 
Loading...
Searching...
No Matches
plugin.mqtt.MqttSender Class Reference

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
 

Detailed Description

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.

Constructor & Destructor Documentation

◆ __init__()

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 
)
46 queue_size=200, max_retries=5, initial_delay=2, max_delay=60):
47 self._stop_event = threading.Event()
48
49 self.broker_address = broker_address
50 self.broker_port = broker_port
51 self.keepalive = keepalive
52 self.qos = qos
53 self.retain = retain
54 self.status_topic = status_topic
55 self.max_retries = max_retries
56 self.initial_delay = initial_delay
57 self.max_delay = max_delay
58
59 self._connected = False
60 self._msg_queue = queue.Queue(maxsize=queue_size)
61
62 self.client = mqtt.Client(
63 client_id=client_id,
64 callback_api_version=mqtt.CallbackAPIVersion.VERSION2
65 )
66 if username:
67 self.client.username_pw_set(username, password)
68
69 # Last Will: published by the broker itself if the connection drops
70 # ungracefully (e.g. host crash) - must be set before connect().
71 if self.status_topic:
72 self.client.will_set(self.status_topic, payload="offline", qos=1, retain=True)
73
74 self.client.on_connect = self._on_connect
75 self.client.on_disconnect = self._on_disconnect
76
77 # built-in reconnect mechanism of paho with exponential backoff
78 self.client.reconnect_delay_set(min_delay=1, max_delay=30)
79
80 self._worker = threading.Thread(target=self._worker_loop, daemon=True)
81 self._worker.start()
82
83 self._connect()
84

Member Function Documentation

◆ enqueue()

plugin.mqtt.MqttSender.enqueue (   self,
  topic,
  json_payload 
)

Buffers a message for delivery.

Parameters
topicMQTT topic
json_payloadalready serialized JSON string
85 def enqueue(self, topic, json_payload):
86 r"""!Buffers a message for delivery.
87
88 @param topic: MQTT topic
89 @param json_payload: already serialized JSON string"""
90 try:
91 self._msg_queue.put_nowait((topic, json_payload, 0))
92 except queue.Full:
93 try:
94 dropped_topic, _, _ = self._msg_queue.get_nowait()
95 self._msg_queue.put_nowait((topic, json_payload, 0))
96 logging.warning("MQTT Plugin: buffer full (max %d) - dropping oldest buffered message (topic '%s') in favor of the new one",
97 self._msg_queue.maxsize, dropped_topic)
98 except queue.Empty: # pragma: no cover - race, practically unreachable
99 logging.error("MQTT Plugin: message for topic '%s' discarded (buffer full)", topic)
100

◆ is_alive()

plugin.mqtt.MqttSender.is_alive (   self)

Checks whether the worker thread is still running.

Returns
True or False
101 def is_alive(self):
102 r"""!Checks whether the worker thread is still running.
103
104 @return True or False"""
105 return self._worker.is_alive()
106

◆ shutdown()

plugin.mqtt.MqttSender.shutdown (   self)

Graceful shutdown with queue drain (analogous to TelegramSender.shutdown)

107 def shutdown(self):
108 r"""!Graceful shutdown with queue drain (analogous to TelegramSender.shutdown)"""
109 logging.info("MQTT Plugin: shutting down connection...")
110 self._stop_event.set()
111
112 timeout = time.time() + 5
113 while not self._msg_queue.empty() and time.time() < timeout:
114 time.sleep(0.1)
115
116 remaining = self._msg_queue.qsize()
117 if remaining > 0:
118 logging.warning("MQTT Plugin: %d unsent messages discarded during shutdown", remaining)
119
120 if self.status_topic and self._connected:
121 try:
122 info = self.client.publish(self.status_topic, payload="offline", qos=1, retain=True)
123 info.wait_for_publish(timeout=2)
124 except Exception:
125 pass
126
127 self._worker.join(timeout=5)
128 try:
129 self.client.loop_stop()
130 self.client.disconnect()
131 except Exception:
132 pass
133 logging.debug("MQTT Plugin: connection closed")
134

◆ _connect()

plugin.mqtt.MqttSender._connect (   self)
protected

Establishes the connection to the broker asynchronously and starts the network thread.

139 def _connect(self):
140 r"""!Establishes the connection to the broker asynchronously and starts the network thread"""
141 try:
142 self.client.connect_async(self.broker_address, self.broker_port, keepalive=self.keepalive)
143 self.client.loop_start()
144 except Exception as e:
145 logging.error("MQTT Plugin: failed to connect to %s:%d: %s", self.broker_address, self.broker_port, e)
146

◆ _on_connect()

plugin.mqtt.MqttSender._on_connect (   self,
  client,
  userdata,
  flags,
  reason_code,
  properties = None 
)
protected

paho-mqtt callback on (re-)connect

147 def _on_connect(self, client, userdata, flags, reason_code, properties=None):
148 r"""!paho-mqtt callback on (re-)connect"""
149 rc = getattr(reason_code, "value", reason_code)
150 if rc == 0:
151 logging.info("MQTT Plugin: connected to %s:%d", self.broker_address, self.broker_port)
152 self._connected = True
153 if self.status_topic:
154 client.publish(self.status_topic, payload="online", qos=1, retain=True)
155 else:
156 logging.error("MQTT Plugin: connection refused (reason_code=%s)", reason_code)
157 self._connected = False
158

◆ _on_disconnect()

plugin.mqtt.MqttSender._on_disconnect (   self,
  client,
  userdata,
  disconnect_flags,
  reason_code,
  properties = None 
)
protected

paho-mqtt callback on connection loss

159 def _on_disconnect(self, client, userdata, disconnect_flags, reason_code, properties=None):
160 r"""!paho-mqtt callback on connection loss"""
161 logging.warning("MQTT Plugin: connection lost (reason_code=%s) - buffered messages are waiting for reconnect", reason_code)
162 self._connected = False
163

◆ _requeue()

plugin.mqtt.MqttSender._requeue (   self,
  topic,
  json_payload,
  retry_count 
)
protected

Puts a message back into the queue after a failed delivery attempt.

164 def _requeue(self, topic, json_payload, retry_count):
165 r"""!Puts a message back into the queue after a failed delivery attempt"""
166 try:
167 self._msg_queue.put_nowait((topic, json_payload, retry_count))
168 except queue.Full:
169 logging.error("MQTT Plugin: message for topic '%s' discarded while requeuing (buffer full)", topic)
170

◆ _worker_loop()

plugin.mqtt.MqttSender._worker_loop (   self)
protected

Processes the queue sequentially: connected -> send, otherwise buffer & wait.

171 def _worker_loop(self):
172 r"""!Processes the queue sequentially: connected -> send, otherwise buffer & wait"""
173 delay = self.initial_delay
174 while not self._stop_event.is_set():
175 try:
176 topic, json_payload, retry_count = self._msg_queue.get(timeout=1)
177 except queue.Empty:
178 continue
179
180 if not self._connected:
181 logging.debug("MQTT Plugin: not connected - message for topic '%s' stays buffered", topic)
182 self._requeue(topic, json_payload, retry_count)
183 time.sleep(min(delay, self.max_delay))
184 continue
185
186 try:
187 result = self.client.publish(topic, json_payload, qos=self.qos, retain=self.retain)
188 if result.rc != mqtt.MQTT_ERR_SUCCESS:
189 raise RuntimeError(f"publish() returned error code {result.rc}")
190
191 result.wait_for_publish(timeout=5)
192 if not result.is_published():
193 raise RuntimeError("no delivery confirmation within 5s")
194
195 logging.debug("MQTT Plugin: message published on topic '%s' (MID=%s, QoS=%d)", topic, result.mid, self.qos)
196 delay = self.initial_delay
197
198 except Exception as e:
199 if retry_count >= self.max_retries:
200 logging.error("MQTT Plugin: message for topic '%s' permanently discarded after %d attempts: %s", topic, self.max_retries, e)
201 else:
202 logging.warning("MQTT Plugin: delivery attempt for topic '%s' failed (%d/%d): %s", topic, retry_count + 1, self.max_retries, e)
203 self._requeue(topic, json_payload, retry_count + 1)
204 time.sleep(min(delay, self.max_delay))
205 delay = min(delay * 2, self.max_delay)
206
207
208# ===========================
209# BoswatchPlugin-Class
210# ===========================
211

Field Documentation

◆ _stop_event

plugin.mqtt.MqttSender._stop_event
protected

◆ broker_address

plugin.mqtt.MqttSender.broker_address

◆ broker_port

plugin.mqtt.MqttSender.broker_port

◆ keepalive

plugin.mqtt.MqttSender.keepalive

◆ qos

plugin.mqtt.MqttSender.qos

◆ retain

plugin.mqtt.MqttSender.retain

◆ status_topic

plugin.mqtt.MqttSender.status_topic

◆ max_retries

plugin.mqtt.MqttSender.max_retries

◆ initial_delay

plugin.mqtt.MqttSender.initial_delay

◆ max_delay

plugin.mqtt.MqttSender.max_delay

◆ _connected

plugin.mqtt.MqttSender._connected
protected

◆ _msg_queue

plugin.mqtt.MqttSender._msg_queue
protected

◆ client

plugin.mqtt.MqttSender.client

◆ _worker

plugin.mqtt.MqttSender._worker
protected