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

Publishes BOSWatch alarms as JSON payload via MQTT. More...

Public Member Functions

 __init__ (self, config)
 Do not change anything here!
 
 onLoad (self)
 Called by import of the plugin.
 
 setup (self)
 Called before alarm - recreates the sender if needed (self-healing)
 
 fms (self, bwPacket)
 Called on FMS alarm.
 
 pocsag (self, bwPacket)
 Called on POCSAG alarm.
 
 zvei (self, bwPacket)
 Called on ZVEI alarm.
 
 msg (self, bwPacket)
 Called on MSG packet.
 
 onUnload (self)
 Called by destruction of the plugin.
 
- Public Member Functions inherited from plugin.pluginBase.PluginBase
 teardown (self)
 Called after alarm can be inherited.
 
 parseWildcards (self, msg)
 Return the message with parsed wildcards.
 

Data Fields

 brokerAddress
 
 brokerPort
 
 topicTemplate
 
 clientId
 
 username
 
 password
 
 qos
 
 retain
 
 keepalive
 
 statusTopic
 
 queueSize
 
 maxRetries
 
 initialDelay
 
 maxDelay
 
 sender
 
- Data Fields inherited from plugin.pluginBase.PluginBase
 config
 

Protected Member Functions

 _ensure_sender (self)
 Ensures that a MqttSender instance exists - checking-with-hasattr pattern analogous to telegram.py's _ensure_sender()
 
 _publish (self, bwPacket)
 Serializes the complete bwPacket to JSON and enqueues it with the sender.
 
- Protected Member Functions inherited from plugin.pluginBase.PluginBase
 _cleanup (self)
 Cleanup routine calls onUnload() directly.
 
 _run (self, bwPacket)
 start an complete running turn of an plugin.
 
 _getStatistics (self)
 Returns statistical information's from last plugin run.
 

Static Protected Member Functions

 _json_serializer (obj)
 Fail-safe JSON serializer for non-primitive data types.
 
 _packet_to_dict (bwPacket)
 Extracts the complete content of the bwPacket object as a dict.
 

Additional Inherited Members

- Protected Attributes inherited from plugin.pluginBase.PluginBase
 _pluginName
 
 _bwPacket
 
 _sumTime
 
 _cumTime
 
 _setupTime
 
 _alarmTime
 
 _teardownTime
 
 _runCount
 
 _setupErrorCount
 
 _alarmErrorCount
 
 _teardownErrorCount
 
- Static Protected Attributes inherited from plugin.pluginBase.PluginBase
list _pluginsActive = []
 

Detailed Description

Publishes BOSWatch alarms as JSON payload via MQTT.

Constructor & Destructor Documentation

◆ __init__()

plugin.mqtt.BoswatchPlugin.__init__ (   self,
  config 
)

Do not change anything here!

Reimplemented from plugin.pluginBase.PluginBase.

215 def __init__(self, config):
216 r"""!Do not change anything here!"""
217 super().__init__(__name__, config) # you can access the config class on 'self.config'
218

Member Function Documentation

◆ onLoad()

plugin.mqtt.BoswatchPlugin.onLoad (   self)

Called by import of the plugin.

Reimplemented from plugin.pluginBase.PluginBase.

219 def onLoad(self):
220 r"""!Called by import of the plugin"""
221 self.brokerAddress = self.config.get("brokerAddress")
222 self.brokerPort = int(self.config.get("brokerPort", default=1883))
223 self.topicTemplate = self.config.get("topic", default="boswatch/alarm")
224 self.clientId = self.config.get("clientId", default="boswatch")
225 self.username = self.config.get("username", default=None)
226 self.password = self.config.get("password", default=None)
227 self.qos = int(self.config.get("qos", default=0))
228 self.retain = bool(self.config.get("retain", default=False))
229 self.keepalive = int(self.config.get("keepalive", default=60))
230 self.statusTopic = self.config.get("statusTopic", default=None)
231 self.queueSize = int(self.config.get("queueSize", default=200))
232 self.maxRetries = int(self.config.get("maxRetries", default=5))
233 self.initialDelay = int(self.config.get("initialDelay", default=2))
234 self.maxDelay = int(self.config.get("maxDelay", default=60))
235
236 if self.qos not in (0, 1, 2):
237 logging.warning("MQTT Plugin: invalid qos value '%s' - using 0", self.qos)
238 self.qos = 0
239
240 self.sender = None
241 self._ensure_sender()
242

◆ setup()

plugin.mqtt.BoswatchPlugin.setup (   self)

Called before alarm - recreates the sender if needed (self-healing)

Reimplemented from plugin.pluginBase.PluginBase.

243 def setup(self):
244 r"""!Called before alarm - recreates the sender if needed (self-healing)"""
245 self._ensure_sender()
246

◆ fms()

plugin.mqtt.BoswatchPlugin.fms (   self,
  bwPacket 
)

Called on FMS alarm.

Parameters
bwPacketbwPacket instance Remove if not implemented

Reimplemented from plugin.pluginBase.PluginBase.

247 def fms(self, bwPacket):
248 r"""!Called on FMS alarm
249
250 @param bwPacket: bwPacket instance
251 Remove if not implemented"""
252 self._publish(bwPacket)
253

◆ pocsag()

plugin.mqtt.BoswatchPlugin.pocsag (   self,
  bwPacket 
)

Called on POCSAG alarm.

Parameters
bwPacketbwPacket instance Remove if not implemented

Reimplemented from plugin.pluginBase.PluginBase.

254 def pocsag(self, bwPacket):
255 r"""!Called on POCSAG alarm
256
257 @param bwPacket: bwPacket instance
258 Remove if not implemented"""
259 self._publish(bwPacket)
260

◆ zvei()

plugin.mqtt.BoswatchPlugin.zvei (   self,
  bwPacket 
)

Called on ZVEI alarm.

Parameters
bwPacketbwPacket instance Remove if not implemented

Reimplemented from plugin.pluginBase.PluginBase.

261 def zvei(self, bwPacket):
262 r"""!Called on ZVEI alarm
263
264 @param bwPacket: bwPacket instance
265 Remove if not implemented"""
266 self._publish(bwPacket)
267

◆ msg()

plugin.mqtt.BoswatchPlugin.msg (   self,
  bwPacket 
)

Called on MSG packet.

Parameters
bwPacketbwPacket instance Remove if not implemented

Reimplemented from plugin.pluginBase.PluginBase.

268 def msg(self, bwPacket):
269 r"""!Called on MSG packet
270
271 @param bwPacket: bwPacket instance
272 Remove if not implemented"""
273 self._publish(bwPacket)
274

◆ onUnload()

plugin.mqtt.BoswatchPlugin.onUnload (   self)

Called by destruction of the plugin.

Reimplemented from plugin.pluginBase.PluginBase.

275 def onUnload(self):
276 r"""!Called by destruction of the plugin"""
277 if self.sender is not None:
278 self.sender.shutdown()
279

◆ _ensure_sender()

plugin.mqtt.BoswatchPlugin._ensure_sender (   self)
protected

Ensures that a MqttSender instance exists - checking-with-hasattr pattern analogous to telegram.py's _ensure_sender()

284 def _ensure_sender(self):
285 r"""!Ensures that a MqttSender instance exists - checking-with-hasattr pattern
286 analogous to telegram.py's _ensure_sender()"""
287 if self.sender is not None and self.sender.is_alive():
288 return
289
290 if not self.brokerAddress:
291 logging.error("MQTT Plugin: 'brokerAddress' not configured - plugin will not send messages.")
292 return
293
294 try:
295 self.sender = MqttSender(
296 broker_address=self.brokerAddress,
297 broker_port=self.brokerPort,
298 client_id=self.clientId,
299 username=self.username,
300 password=self.password,
301 keepalive=self.keepalive,
302 qos=self.qos,
303 retain=self.retain,
304 status_topic=self.statusTopic,
305 queue_size=self.queueSize,
306 max_retries=self.maxRetries,
307 initial_delay=self.initialDelay,
308 max_delay=self.maxDelay
309 )
310 logging.debug("MQTT Plugin: MqttSender initialized")
311 except Exception as e:
312 logging.error("MQTT Plugin: sender could not be initialized: %s", e)
313 self.sender = None
314

◆ _publish()

plugin.mqtt.BoswatchPlugin._publish (   self,
  bwPacket 
)
protected

Serializes the complete bwPacket to JSON and enqueues it with the sender.

Parameters
bwPacketbwPacket instance
315 def _publish(self, bwPacket):
316 r"""!Serializes the complete bwPacket to JSON and enqueues it with the sender
317
318 @param bwPacket: bwPacket instance"""
319 if self.sender is None:
320 logging.warning("MQTT Plugin not configured - packet dropped")
321 return
322
323 payload = self._packet_to_dict(bwPacket)
324 topic = self.parseWildcards(self.topicTemplate)
325
326 try:
327 json_payload = json.dumps(payload, default=self._json_serializer, ensure_ascii=False)
328 except Exception as e:
329 logging.error("MQTT Plugin: error during JSON serialization: %s", e)
330 return
331
332 self.sender.enqueue(topic, json_payload)
333

◆ _json_serializer()

plugin.mqtt.BoswatchPlugin._json_serializer (   obj)
staticprotected

Fail-safe JSON serializer for non-primitive data types.

Packet.set() currently only ever stores strings, but this fallback defensively covers the case that a future module (e.g. a custom extension) writes an exotic object into a field - so a single unexpected data type cannot crash the whole plugin.

Parameters
objthe object to serialize
Returns
: JSON-compatible representation
335 def _json_serializer(obj):
336 r"""!Fail-safe JSON serializer for non-primitive data types.
337
338 Packet.set() currently only ever stores strings, but this fallback
339 defensively covers the case that a future module (e.g. a custom
340 extension) writes an exotic object into a field - so a single
341 unexpected data type cannot crash the whole plugin.
342
343 @param obj: the object to serialize
344 @return: JSON-compatible representation"""
345 if isinstance(obj, (datetime.date, datetime.datetime)):
346 return obj.isoformat()
347 if isinstance(obj, bytes):
348 return obj.decode("utf-8", errors="replace")
349 if isinstance(obj, set):
350 return list(obj)
351 try:
352 return str(obj)
353 except Exception:
354 return repr(obj)
355

◆ _packet_to_dict()

plugin.mqtt.BoswatchPlugin._packet_to_dict (   bwPacket)
staticprotected

Extracts the complete content of the bwPacket object as a dict.

Analogous to the approach in module/packetDump.py: instead of picking individual fields, the entire internal state of the packet is taken over, so fields added by other modules (Descriptor, Multicast, Geocoding, ...) automatically end up in the MQTT payload.

Parameters
bwPacketbwPacket instance
Returns
: dict with all fields of the packet
357 def _packet_to_dict(bwPacket):
358 r"""!Extracts the complete content of the bwPacket object as a dict.
359
360 Analogous to the approach in module/packetDump.py: instead of
361 picking individual fields, the entire internal state of the packet
362 is taken over, so fields added by other modules (Descriptor,
363 Multicast, Geocoding, ...) automatically end up in the MQTT payload.
364
365 @param bwPacket: bwPacket instance
366 @return: dict with all fields of the packet"""
367 if hasattr(bwPacket, '_packet'):
368 return bwPacket._packet.copy()
369 try:
370 return {k: v for k, v in bwPacket.__dict__.items() if not k.startswith('_')}
371 except Exception as e:
372 logging.warning("MQTT Plugin: could not read packet fields: %s", e)
373 return {}

Field Documentation

◆ brokerAddress

plugin.mqtt.BoswatchPlugin.brokerAddress

◆ brokerPort

plugin.mqtt.BoswatchPlugin.brokerPort

◆ topicTemplate

plugin.mqtt.BoswatchPlugin.topicTemplate

◆ clientId

plugin.mqtt.BoswatchPlugin.clientId

◆ username

plugin.mqtt.BoswatchPlugin.username

◆ password

plugin.mqtt.BoswatchPlugin.password

◆ qos

plugin.mqtt.BoswatchPlugin.qos

◆ retain

plugin.mqtt.BoswatchPlugin.retain

◆ keepalive

plugin.mqtt.BoswatchPlugin.keepalive

◆ statusTopic

plugin.mqtt.BoswatchPlugin.statusTopic

◆ queueSize

plugin.mqtt.BoswatchPlugin.queueSize

◆ maxRetries

plugin.mqtt.BoswatchPlugin.maxRetries

◆ initialDelay

plugin.mqtt.BoswatchPlugin.initialDelay

◆ maxDelay

plugin.mqtt.BoswatchPlugin.maxDelay

◆ sender

plugin.mqtt.BoswatchPlugin.sender