BOSWatch 3
Python Script to receive and decode German BOS Information with rtl_fm and multimon-NG
 
Loading...
Searching...
No Matches
module.multicast.BoswatchModule Class Reference

Multicast module with multi-instance support and active trigger mechanism. More...

Public Member Functions

 __init__ (self, config)
 init preload some needed locals and then call onLoad() directly
 
 onLoad (self)
 Initialize module configuration and start the global cleanup thread.
 
 doWork (self, bwPacket)
 Process an incoming packet and handle multicast logic.
 
 onUnload (self)
 Unregister instance from the global cleanup process.
 

Data Fields

 instance_id
 
 name
 
- Data Fields inherited from module.moduleBase.ModuleBase
 config
 

Protected Member Functions

 _get_packet_data (self, bwPacket)
 Safely extract all fields from packet as a dictionary.
 
 _combine_results (self, *results)
 Combine multiple result sources into a single list or status.
 
 _add_tone_ric_packet (self, freq, packet_dict)
 Add a tone-RIC to the shared buffer.
 
 _get_queued_packets (self)
 Pop and return all packets currently in the static queue.
 
 _copy_packet_dict_to_packet (self, recipient_dict, packet, index=1)
 Copy dict fields to Packet with timestamp shift for DB uniqueness.
 
 _distribute_complete (self, freq, text_packet_dict)
 Create full multicast packets with message content.
 
 _create_incomplete_multicast (self, freq, recipient_dicts)
 Generate multicast packets for timeouts (no text message).
 
 _enrich_normal_alarm (self, bwPacket, packet_dict)
 Enrich a standard single alarm with multicast metadata.
 
 _handle_delimiter (self, freq, ric, bwPacket=None)
 Handle delimiter packet and clear orphaned tone-RICs.
 
 _set_mcast_metadata (self, packet, mode, role, source="", count="1", index="1")
 Helper to set standard multicast fields and register wildcards.
 
 _apply_list_tags (self, packet, recipient_dicts)
 Helper to aggregate fields from all recipients into comma-separated lists.
 
 _register_wildcard_safe (self, wildcard, field)
 Register wildcard if not already globally registered.
 
 _cleanup_worker (self)
 Per-instance background thread for timeout management.
 
 _check_all_my_frequencies (self)
 Monitor timeouts for all frequencies assigned to this instance.
 
 _check_instance_auto_clear (self, freq)
 Check if frequency has exceeded timeout (called from doWork).
 
 _cleanup_hard_timeout (self)
 Failsafe for really old packets.
 
 _send_wakeup_trigger (self, freq, fallback_ric)
 Send a loopback trigger using the standard TCPClient class.
 
- Protected Member Functions inherited from module.moduleBase.ModuleBase
 _cleanup (self)
 Cleanup routine calls onUnload() directly.
 
 _run (self, bwPacket)
 start an run of the module.
 
 _getStatistics (self)
 Returns statistical information's from last module run.
 

Protected Attributes

 _my_frequencies
 
 _auto_clear_timeout
 
 _hard_timeout
 
 _delimiter_rics
 
 _text_rics
 
 _netident_rics
 
 _trigger_ric
 
 _trigger_host
 
 _trigger_port
 
 _tone_ric_packets
 
 _last_tone_ric_time
 
 _processing_text_ric
 
 _processing_text_ric_started
 
 _wildcards_registered
 
 _packet_queue
 
 _lock
 
 _queue_lock
 
 _running
 
 _cleanup_thread
 
 _MAGIC_WAKEUP_MSG
 
- Protected Attributes inherited from module.moduleBase.ModuleBase
 _moduleName
 
 _cumTime
 
 _moduleTime
 
 _runCount
 
 _moduleErrorCount
 

Static Protected Attributes

str _TRIGGER_HOST = "127.0.0.1"
 
int _TRIGGER_PORT = 8080
 
str _MAGIC_WAKEUP_MSG = "###_MULTICAST_WAKEUP_###"
 
str _DEFAULT_TRIGGER_RIC = "9999999"
 
- Static Protected Attributes inherited from module.moduleBase.ModuleBase
list _modulesActive = []
 

Additional Inherited Members

- Static Public Member Functions inherited from module.moduleBase.ModuleBase
 registerWildcard (newWildcard, bwPacketField)
 Register a new wildcard.
 

Detailed Description

Multicast module with multi-instance support and active trigger mechanism.

This module handles multicast alarm distribution. It manages the correlation between tone-RICs (recipients) and text-RICs (message content), ensuring reliable alarm delivery even in complex multi-frequency scenarios.

Constructor & Destructor Documentation

◆ __init__()

module.multicast.BoswatchModule.__init__ (   self,
  moduleName 
)

init preload some needed locals and then call onLoad() directly

Reimplemented from module.moduleBase.ModuleBase.

48 def __init__(self, config):
49 super().__init__(__name__, config)
50

Member Function Documentation

◆ onLoad()

module.multicast.BoswatchModule.onLoad (   self)

Initialize module configuration and start the global cleanup thread.

Parameters
None
Returns
None

Reimplemented from module.moduleBase.ModuleBase.

51 def onLoad(self):
52 r"""!Initialize module configuration and start the global cleanup thread.
53
54 @param None
55 @return None"""
56 self._my_frequencies = set()
57 self.instance_id = hex(id(self))[-4:]
58 self.name = f"MCAST_{self.instance_id}"
59
60 self._auto_clear_timeout = int(self.config.get("autoClearTimeout", default=10))
61 self._hard_timeout = self._auto_clear_timeout * 3
62
63 def parse_list(key):
64 val = self.config.get(key)
65 if val:
66 return [x.strip() for x in str(val).split(",") if x.strip()]
67 return []
68
69 self._delimiter_rics = parse_list("delimiterRics")
70 self._text_rics = parse_list("textRics")
71 self._netident_rics = parse_list("netIdentRics")
72
73 trigger_ric_cfg = self.config.get("triggerRic")
74 if trigger_ric_cfg:
75 self._trigger_ric = str(trigger_ric_cfg).strip()
76 else:
77 self._trigger_ric = None
78
79 self._trigger_host = self.config.get("triggerHost", default=self._TRIGGER_HOST)
80 self._trigger_port = int(self.config.get("triggerPort", default=self._TRIGGER_PORT))
81
82 # --- Per-instance state (replaces all former class-variables) ---
83 # Key: frequency string (e.g. "85.125M")
84 self._tone_ric_packets = defaultdict(list) # buffered tone-RICs per frequency
85 self._last_tone_ric_time = defaultdict(float) # last arrival time per frequency
86 self._processing_text_ric = defaultdict(bool) # text-RIC currently being processed?
87 self._processing_text_ric_started = defaultdict(float) # when did processing start?
88 self._wildcards_registered = set() # avoid double-registering wildcards
89 self._packet_queue = [] # deferred packets waiting for trigger
90
91 # --- Locks (only needed within this instance, no cross-instance sharing) ---
92 self._lock = threading.Lock()
93 self._queue_lock = threading.Lock()
94
95 # --- Per-instance cleanup thread ---
96 self._running = True
97 self._cleanup_thread = threading.Thread(target=self._cleanup_worker, daemon=True)
98 self._cleanup_thread.start()
99
100 logging.info("[%s] Multicast module loaded", self.name)
101
102# ============================================================
103# MAIN PROCESSING
104# ============================================================
105

◆ doWork()

module.multicast.BoswatchModule.doWork (   self,
  bwPacket 
)

Process an incoming packet and handle multicast logic.

Enriches packets with multicast metadata (mode, role, source). Does NOT filter - all packets pass through, downstream modules handle filtering.

Parameters
bwPacketA BOSWatch packet instance or list of packets
Returns
bwPacket, a list of packets, or None if no processing

Reimplemented from module.moduleBase.ModuleBase.

106 def doWork(self, bwPacket):
107 r"""!Process an incoming packet and handle multicast logic.
108
109 Enriches packets with multicast metadata (mode, role, source).
110 Does NOT filter - all packets pass through, downstream modules handle filtering.
111
112 @param bwPacket: A BOSWatch packet instance or list of packets
113 @return bwPacket, a list of packets, or None if no processing"""
114 if isinstance(bwPacket, list):
115 result_packets = []
116 for single_packet in bwPacket:
117 processed = self.doWork(single_packet)
118 if processed is not None and processed is not False:
119 if isinstance(processed, list):
120 result_packets.extend(processed)
121 else:
122 result_packets.append(processed)
123 return result_packets if result_packets else None
124
125 packet_dict = self._get_packet_data(bwPacket)
126 msg = packet_dict.get("message")
127 ric = packet_dict.get("ric")
128 freq = packet_dict.get("frequency", "default")
129 mode = packet_dict.get("mode")
130
131 # Handle wakeup triggers
132 if msg == BoswatchModule._MAGIC_WAKEUP_MSG:
133 if self._trigger_ric and ric != self._trigger_ric:
134 return None
135 logging.debug("[%s] Wakeup trigger received (RIC=%s)", self.name, ric)
136 queued = self._get_queued_packets()
137 return queued if queued else None
138
139 # Only process POCSAG
140 if mode != "pocsag":
141 queued = self._get_queued_packets()
142 return queued if queued else None
143
144 self._my_frequencies.add(freq)
145
146 # Determine if this is a text-RIC
147 is_text_ric = False
148 if self._text_rics:
149 is_text_ric = ric in self._text_rics and msg and msg.strip()
150 else:
151 with self._lock:
152 is_text_ric = msg and msg.strip() and len(self._tone_ric_packets[freq]) > 0
153
154 if is_text_ric:
155 with self._lock:
156 self._processing_text_ric[freq] = True
157 self._processing_text_ric_started[freq] = time.time()
158
159 queued_packets = self._get_queued_packets()
160 incomplete_packets = None if is_text_ric else self._check_instance_auto_clear(freq)
161
162 # === CONTROL PACKETS (netident, delimiter) ===
163 # Mark and pass through - no filtering!
164
165 if self._netident_rics and ric in self._netident_rics:
166 self._set_mcast_metadata(bwPacket, "control", "netident", ric)
167 return self._combine_results(incomplete_packets, queued_packets, [bwPacket])
168
169 if self._delimiter_rics and ric in self._delimiter_rics:
170 delimiter_incomplete = self._handle_delimiter(freq, ric, bwPacket)
171 return self._combine_results(delimiter_incomplete, incomplete_packets, queued_packets)
172
173 # === TONE-RICs (no message) ===
174 if not msg or not msg.strip():
175 self._add_tone_ric_packet(freq, packet_dict)
176 return self._combine_results(incomplete_packets, queued_packets, False)
177
178 # === TEXT-RICs (with message) ===
179 if is_text_ric and msg:
180 logging.info("[%s] Text-RIC received: RIC=%s", self.name, ric)
181 alarm_packets = self._distribute_complete(freq, packet_dict)
182 with self._lock:
183 self._processing_text_ric[freq] = False
184 self._processing_text_ric_started.pop(freq, None)
185
186 if not alarm_packets:
187 logging.warning("[%s] No tone-RICs for text-RIC=%s", self.name, ric)
188 normal = self._enrich_normal_alarm(bwPacket, packet_dict)
189 return self._combine_results(normal, incomplete_packets, queued_packets)
190 else:
191 return self._combine_results(alarm_packets, incomplete_packets, queued_packets)
192
193 # === SINGLE ALARM (message but no text-RICs configured) ===
194 if msg:
195 normal = self._enrich_normal_alarm(bwPacket, packet_dict)
196 return self._combine_results(normal, incomplete_packets, queued_packets)
197
198 return self._combine_results(incomplete_packets, queued_packets)
199
200# ============================================================
201# PACKET PROCESSING HELPERS (called by doWork)
202# ============================================================
203

◆ _get_packet_data()

module.multicast.BoswatchModule._get_packet_data (   self,
  bwPacket 
)
protected

Safely extract all fields from packet as a dictionary.

Handles both dict objects and Packet instances. Dynamically extracts all fields including those added by other modules.

Parameters
bwPacketPacket instance or dict
Returns
dict: Complete dictionary of all packet fields
204 def _get_packet_data(self, bwPacket):
205 r"""!Safely extract all fields from packet as a dictionary.
206
207 Handles both dict objects and Packet instances.
208 Dynamically extracts all fields including those added by other modules.
209
210 @param bwPacket: Packet instance or dict
211 @return dict: Complete dictionary of all packet fields"""
212 # 1. Fall: Es ist bereits ein Dictionary
213 if isinstance(bwPacket, dict):
214 return bwPacket.copy()
215
216 # 2. Fall: Es ist ein Packet-Objekt (Daten liegen in _packet)
217 if hasattr(bwPacket, '_packet'):
218 return bwPacket._packet.copy()
219
220 # 3. Fallback: Falls es ein anderes Objekt ist, versuche __dict__ ohne '_' Filter für 'packet'
221 try:
222 return {k: v for k, v in bwPacket.__dict__.items() if not k.startswith('_')}
223 except Exception as e:
224 logging.warning("[%s] Error: %s", self.name, e)
225 return {}
226

◆ _combine_results()

module.multicast.BoswatchModule._combine_results (   self,
*  results 
)
protected

Combine multiple result sources into a single list or status.

Parameters
resultsMultiple packet objects, lists, or booleans
Returns
combined list, False or None
227 def _combine_results(self, *results):
228 r"""!Combine multiple result sources into a single list or status.
229
230 @param results: Multiple packet objects, lists, or booleans
231 @return combined list, False or None"""
232 combined = []
233 has_false = False
234 for result in results:
235 if result is False:
236 has_false = True
237 continue
238 if result is None:
239 continue
240 if isinstance(result, list):
241 combined.extend(result)
242 else:
243 combined.append(result)
244 if combined:
245 return combined
246 return False if has_false else None
247
248# ============================================================
249# TONE-RIC BUFFER MANAGEMENT
250# ============================================================
251

◆ _add_tone_ric_packet()

module.multicast.BoswatchModule._add_tone_ric_packet (   self,
  freq,
  packet_dict 
)
protected

Add a tone-RIC to the shared buffer.

Parameters
freqFrequency identifier
packet_dictDictionary containing packet data
Returns
None
252 def _add_tone_ric_packet(self, freq, packet_dict):
253 r"""!Add a tone-RIC to the shared buffer.
254
255 @param freq: Frequency identifier
256 @param packet_dict: Dictionary containing packet data
257 @return None"""
258 with self._lock:
259 stored_packet = packet_dict.copy()
260 stored_packet['_multicast_timestamp'] = time.time()
261 self._tone_ric_packets[freq].append(stored_packet)
262 self._last_tone_ric_time[freq] = stored_packet['_multicast_timestamp']
263 logging.debug("[%s] Tone-RIC added: RIC=%s (total: %d on %s)", self.name, stored_packet.get('ric'), len(self._tone_ric_packets[freq]), freq)
264

◆ _get_queued_packets()

module.multicast.BoswatchModule._get_queued_packets (   self)
protected

Pop and return all packets currently in the static queue.

Parameters
None
Returns
list: List of packets or None
265 def _get_queued_packets(self):
266 r"""!Pop and return all packets currently in the static queue.
267
268 @param None
269 @return list: List of packets or None"""
270 with self._queue_lock:
271 if self._packet_queue:
272 packets = self._packet_queue[:]
273 self._packet_queue.clear()
274 return packets
275 return None
276
277# ============================================================
278# MULTICAST PACKET CREATION
279# ============================================================
280

◆ _copy_packet_dict_to_packet()

module.multicast.BoswatchModule._copy_packet_dict_to_packet (   self,
  recipient_dict,
  packet,
  index = 1 
)
protected

Copy dict fields to Packet with timestamp shift for DB uniqueness.

Parameters
recipient_dictSource dictionary
packetTarget Packet object
indexPacket index (1-based) - shifts timestamp by milliseconds
Returns
None
281 def _copy_packet_dict_to_packet(self, recipient_dict, packet, index=1):
282 r"""!Copy dict fields to Packet with timestamp shift for DB uniqueness.
283
284 @param recipient_dict: Source dictionary
285 @param packet: Target Packet object
286 @param index: Packet index (1-based) - shifts timestamp by milliseconds
287 @return None"""
288 for k, v in recipient_dict.items():
289 if k.startswith('_'):
290 continue
291
292 if k == 'timestamp':
293 try:
294 # use native float (UNIX-timestamp)
295 ts_float = float(v)
296
297 # shifting for database-unique-keys (only for index >=2)
298 if index > 1:
299 ts_float += 0.001 * (index - 1)
300
301 # back to paket - float (UNIX-timestamp)
302 packet.set(k, ts_float)
303 except (ValueError, TypeError):
304 # failsafe, if no float
305 packet.set(k, v)
306 else:
307 packet.set(k, str(v))
308

◆ _distribute_complete()

module.multicast.BoswatchModule._distribute_complete (   self,
  freq,
  text_packet_dict 
)
protected

Create full multicast packets with message content.

Parameters
freqFrequency identifier
text_packet_dictData of the message-carrying packet
Returns
list: List of fully populated Packet instances
309 def _distribute_complete(self, freq, text_packet_dict):
310 r"""!Create full multicast packets with message content.
311
312 @param freq: Frequency identifier
313 @param text_packet_dict: Data of the message-carrying packet
314 @return list: List of fully populated Packet instances"""
315 with self._lock:
316 recipient_dicts = self._tone_ric_packets[freq].copy()
317 logging.debug("Text RIC found. Matching against %d stored RICs", len(recipient_dicts))
318 self._tone_ric_packets[freq].clear()
319 self._last_tone_ric_time.pop(freq, None)
320
321 if not recipient_dicts:
322 return []
323 text_ric = text_packet_dict.get("ric")
324 message_text = text_packet_dict.get("message")
325 alarm_packets = []
326
327 for idx, recipient_dict in enumerate(recipient_dicts, 1):
328 p = Packet()
329 self._copy_packet_dict_to_packet(recipient_dict, p, idx)
330 p.set("message", message_text)
331 self._apply_list_tags(p, recipient_dicts)
332 self._set_mcast_metadata(p, "complete", "recipient", text_ric, len(recipient_dicts), idx)
333 alarm_packets.append(p)
334
335 logging.info("[%s] Generated %d complete multicast packets for text-RIC %s", self.name, len(alarm_packets), text_ric)
336 return alarm_packets
337

◆ _create_incomplete_multicast()

module.multicast.BoswatchModule._create_incomplete_multicast (   self,
  freq,
  recipient_dicts 
)
protected

Generate multicast packets for timeouts (no text message).

Parameters
freqFrequency identifier
recipient_dictsList of recipient data dictionaries
Returns
list: List of incomplete Packet instances
338 def _create_incomplete_multicast(self, freq, recipient_dicts):
339 r"""!Generate multicast packets for timeouts (no text message).
340
341 @param freq: Frequency identifier
342 @param recipient_dicts: List of recipient data dictionaries
343 @return list: List of incomplete Packet instances"""
344 if not recipient_dicts:
345 return []
346 first_ric = recipient_dicts[0].get("ric", "unknown")
347 incomplete_packets = []
348 for idx, recipient_dict in enumerate(recipient_dicts, 1):
349 p = Packet()
350 self._copy_packet_dict_to_packet(recipient_dict, p, idx)
351 p.set("message", "")
352 self._apply_list_tags(p, recipient_dicts)
353 self._set_mcast_metadata(p, "incomplete", "recipient", first_ric, len(recipient_dicts), idx)
354 incomplete_packets.append(p)
355 return incomplete_packets
356

◆ _enrich_normal_alarm()

module.multicast.BoswatchModule._enrich_normal_alarm (   self,
  bwPacket,
  packet_dict 
)
protected

Enrich a standard single alarm with multicast metadata.

Parameters
bwPacketTarget Packet object
packet_dictSource data dictionary
Returns
list: List containing the enriched packet
357 def _enrich_normal_alarm(self, bwPacket, packet_dict):
358 r"""!Enrich a standard single alarm with multicast metadata.
359
360 @param bwPacket: Target Packet object
361 @param packet_dict: Source data dictionary
362 @return list: List containing the enriched packet"""
363 self._copy_packet_dict_to_packet(packet_dict, bwPacket, index=1)
364 self._apply_list_tags(bwPacket, [packet_dict])
365 self._set_mcast_metadata(bwPacket, "single", "single", packet_dict.get("ric", ""), "1", "1")
366 logging.debug("Creating single-alarm for RIC %s", packet_dict.get('ric'))
367 return [bwPacket]
368

◆ _handle_delimiter()

module.multicast.BoswatchModule._handle_delimiter (   self,
  freq,
  ric,
  bwPacket = None 
)
protected

Handle delimiter packet and clear orphaned tone-RICs.

Parameters
freqFrequency identifier
ricDelimiter RIC
bwPacketOptional delimiter packet instance
Returns
list: Incomplete packets or delimiter control packet
369 def _handle_delimiter(self, freq, ric, bwPacket=None):
370 r"""!Handle delimiter packet and clear orphaned tone-RICs.
371
372 @param freq: Frequency identifier
373 @param ric: Delimiter RIC
374 @param bwPacket: Optional delimiter packet instance
375 @return list: Incomplete packets or delimiter control packet"""
376 with self._lock:
377 orphaned = self._tone_ric_packets[freq].copy()
378 self._tone_ric_packets[freq].clear()
379 self._last_tone_ric_time.pop(freq, None)
380 self._processing_text_ric[freq] = False
381
382 if orphaned:
383 age_seconds = time.time() - orphaned[0].get('_multicast_timestamp', time.time())
384
385 logging.debug("[%s] Delimiter RIC=%s cleared %d orphaned tone-RICs on freq %s: RICs=[%s], age=%.1fs → Creating INCOMPLETE multicast packets for forwarding", self.name, ric, len(orphaned), freq, ', '.join([packet.get('ric', 'unknown') for packet in orphaned]), age_seconds)
386 return self._create_incomplete_multicast(freq, orphaned)
387 if bwPacket is not None:
388 self._set_mcast_metadata(bwPacket, "control", "delimiter", ric, "0", "0")
389 return [bwPacket]
390 return None
391
392# ============================================================
393# PACKET METADATA HELPERS
394# ============================================================
395

◆ _set_mcast_metadata()

module.multicast.BoswatchModule._set_mcast_metadata (   self,
  packet,
  mode,
  role,
  source = "",
  count = "1",
  index = "1" 
)
protected

Helper to set standard multicast fields and register wildcards.

Parameters
packetThe Packet instance to modify
modemulticastMode (complete, incomplete, single, control)
rolemulticastRole (recipient, single, delimiter, netident)
sourceThe originating RIC
countTotal number of recipients
indexCurrent recipient index
Returns
None
396 def _set_mcast_metadata(self, packet, mode, role, source="", count="1", index="1"):
397 r"""!Helper to set standard multicast fields and register wildcards.
398
399 @param packet: The Packet instance to modify
400 @param mode: multicastMode (complete, incomplete, single, control)
401 @param role: multicastRole (recipient, single, delimiter, netident)
402 @param source: The originating RIC
403 @param count: Total number of recipients
404 @param index: Current recipient index
405 @return None"""
406 logging.debug("setting Metadata - Mode: %s, Role: %s, Index: %s of %s for RIC: %s", mode, role, index, count, source)
407 mapping = {
408 "multicastMode": (mode, "{MCAST_MODE}"),
409 "multicastRole": (role, "{MCAST_ROLE}"),
410 "multicastSourceRic": (source, "{MCAST_SOURCE}"),
411 "multicastRecipientCount": (str(count), "{MCAST_COUNT}"),
412 "multicastRecipientIndex": (str(index), "{MCAST_INDEX}")
413 }
414 for key, (val, wildcard) in mapping.items():
415 packet.set(key, val)
416 self._register_wildcard_safe(wildcard, key)
417

◆ _apply_list_tags()

module.multicast.BoswatchModule._apply_list_tags (   self,
  packet,
  recipient_dicts 
)
protected

Helper to aggregate fields from all recipients into comma-separated lists.

Parameters
packetThe target Packet instance
recipient_dictsList of dictionaries of all recipients in this group
Returns
None
418 def _apply_list_tags(self, packet, recipient_dicts):
419 r"""!Helper to aggregate fields from all recipients into comma-separated lists.
420
421 @param packet: The target Packet instance
422 @param recipient_dicts: List of dictionaries of all recipients in this group
423 @return None"""
424 all_fields = set()
425 for r in recipient_dicts:
426 all_fields.update(k for k in r.keys() if not k.startswith('_'))
427
428 for f in sorted(all_fields):
429 list_val = ", ".join([str(r.get(f, "")) for r in recipient_dicts])
430 list_key = f"{f}_list"
431 packet.set(list_key, list_val)
432 self._register_wildcard_safe("{" + f.upper() + "_LIST}", list_key)
433

◆ _register_wildcard_safe()

module.multicast.BoswatchModule._register_wildcard_safe (   self,
  wildcard,
  field 
)
protected

Register wildcard if not already globally registered.

Parameters
wildcardThe wildcard string (e.g. {MCAST_MODE})
fieldThe packet field name
Returns
None
434 def _register_wildcard_safe(self, wildcard, field):
435 r"""!Register wildcard if not already globally registered.
436
437 @param wildcard: The wildcard string (e.g. {MCAST_MODE})
438 @param field: The packet field name
439 @return None"""
440 if wildcard not in self._wildcards_registered:
441 self.registerWildcard(wildcard, field)
442 self._wildcards_registered.add(wildcard)
443
444# ============================================================
445# CLEANUP & TIMEOUT MANAGEMENT
446# ============================================================
447

◆ _cleanup_worker()

module.multicast.BoswatchModule._cleanup_worker (   self)
protected

Per-instance background thread for timeout management.

448 def _cleanup_worker(self):
449 r"""!Per-instance background thread for timeout management."""
450 logging.info("[%s] Cleanup thread started", self.name)
451 while self._running:
452 time.sleep(1)
453 try:
454 self._check_all_my_frequencies()
455 except Exception as e:
456 logging.error("[%s] Error in cleanup thread: %s", self.name, e)
457 if int(time.time()) % 60 == 0:
458 self._cleanup_hard_timeout()
459

◆ _check_all_my_frequencies()

module.multicast.BoswatchModule._check_all_my_frequencies (   self)
protected

Monitor timeouts for all frequencies assigned to this instance.

Parameters
None
Returns
None
460 def _check_all_my_frequencies(self):
461 r"""!Monitor timeouts for all frequencies assigned to this instance.
462
463 @param None
464 @return None"""
465 incomplete_packets = []
466 trigger_data = []
467
468 with self._lock:
469 current_time = time.time()
470 for freq in list(self._my_frequencies):
471 if freq not in self._tone_ric_packets or not self._tone_ric_packets[freq]:
472 continue
473
474 if self._processing_text_ric.get(freq, False):
475 flag_age = current_time - self._processing_text_ric_started.get(freq, current_time)
476 if flag_age > 2:
477 self._processing_text_ric[freq] = False
478 self._processing_text_ric_started.pop(freq, None)
479 else:
480 continue
481
482 last_time = self._last_tone_ric_time.get(freq, 0)
483 if current_time - last_time > self._auto_clear_timeout:
484 recipient_dicts = self._tone_ric_packets[freq].copy()
485 safe_ric = recipient_dicts[0].get('ric', self._DEFAULT_TRIGGER_RIC)
486 trigger_data.append((freq, safe_ric))
487 self._tone_ric_packets[freq].clear()
488 self._last_tone_ric_time.pop(freq, None)
489
490 logging.info("[%s] Auto-clear: %d tone-RICs on %s (Timeout %ds)", self.name, len(recipient_dicts), freq, self._auto_clear_timeout)
491 packets = self._create_incomplete_multicast(freq, recipient_dicts)
492 if packets:
493 incomplete_packets.extend(packets)
494
495 if incomplete_packets:
496 with self._queue_lock:
497 self._packet_queue.extend(incomplete_packets)
498 for freq, safe_ric in trigger_data:
499 self._send_wakeup_trigger(freq, safe_ric)
500

◆ _check_instance_auto_clear()

module.multicast.BoswatchModule._check_instance_auto_clear (   self,
  freq 
)
protected

Check if frequency has exceeded timeout (called from doWork).

Parameters
freqFrequency identifier
Returns
list: Incomplete packets if timeout exceeded, else None
501 def _check_instance_auto_clear(self, freq):
502 r"""!Check if frequency has exceeded timeout (called from doWork).
503
504 @param freq: Frequency identifier
505 @return list: Incomplete packets if timeout exceeded, else None"""
506 with self._lock:
507 if freq not in self._tone_ric_packets or not self._tone_ric_packets[freq]:
508 return None
509 last_time = self._last_tone_ric_time.get(freq, 0)
510 if time.time() - last_time > self._auto_clear_timeout:
511 recipient_dicts = self._tone_ric_packets[freq].copy()
512 self._tone_ric_packets[freq].clear()
513 self._last_tone_ric_time.pop(freq, None)
514 logging.warning("[%s] Auto-clear (doWork): %d packets", self.name, len(recipient_dicts))
515 return self._create_incomplete_multicast(freq, recipient_dicts)
516 return None
517

◆ _cleanup_hard_timeout()

module.multicast.BoswatchModule._cleanup_hard_timeout (   self)
protected

Failsafe for really old packets.

518 def _cleanup_hard_timeout(self):
519 r"""!Failsafe for really old packets."""
520 with self._lock:
521 current_time = time.time()
522 for freq in list(self._tone_ric_packets.keys()):
523 self._tone_ric_packets[freq] = [
524 p for p in self._tone_ric_packets[freq]
525 if current_time - p.get('_multicast_timestamp', 0) < self._hard_timeout
526 ]
527 # cleaning empty frequencies
528 if not self._tone_ric_packets[freq]:
529 del self._tone_ric_packets[freq]
530
531# ============================================================
532# TRIGGER SYSTEM
533# ============================================================
534

◆ _send_wakeup_trigger()

module.multicast.BoswatchModule._send_wakeup_trigger (   self,
  freq,
  fallback_ric 
)
protected

Send a loopback trigger using the standard TCPClient class.

535 def _send_wakeup_trigger(self, freq, fallback_ric):
536 r"""!Send a loopback trigger using the standard TCPClient class."""
537 try:
538 trigger_ric = self._trigger_ric if self._trigger_ric else fallback_ric
539 payload = {
540 "timestamp": time.time(),
541 "mode": "pocsag",
542 "bitrate": "1200",
543 "ric": trigger_ric,
544 "subric": "1",
545 "subricText": "a",
546 "message": self._MAGIC_WAKEUP_MSG,
547 "clientName": "MulticastTrigger",
548 "inputSource": "loopback",
549 "frequency": freq
550 }
551 json_str = json.dumps(payload)
552
553 # using BOSWatch-Architecture
554 client = TCPClient(timeout=2)
555 if client.connect(self._trigger_host, self._trigger_port):
556 # 1. Send
557 client.transmit(json_str)
558
559 # 2. Recieve (getting [ack] and prevents connection reset)
560 client.receive(timeout=1)
561
562 client.disconnect()
563 logging.debug("[%s] Wakeup trigger sent and acknowledged (RIC=%s)", self.name, trigger_ric)
564 else:
565 logging.error("[%s] Could not connect to local server for wakeup", self.name)
566
567 except Exception as e:
568 logging.error("[%s] Failed to send wakeup trigger: %s", self.name, e)
569
570# ============================================================
571# LIFECYCLE (End)
572# ============================================================
573

◆ onUnload()

module.multicast.BoswatchModule.onUnload (   self)

Unregister instance from the global cleanup process.

Parameters
None
Returns
None

Reimplemented from module.moduleBase.ModuleBase.

574 def onUnload(self):
575 r"""!Unregister instance from the global cleanup process.
576
577 @param None
578 @return None"""
579 self._running = False
580 logging.debug("[%s] Multicast instance unloaded", self.name)

Field Documentation

◆ _TRIGGER_HOST

str module.multicast.BoswatchModule._TRIGGER_HOST = "127.0.0.1"
staticprotected

◆ _TRIGGER_PORT

int module.multicast.BoswatchModule._TRIGGER_PORT = 8080
staticprotected

◆ _MAGIC_WAKEUP_MSG [1/2]

str module.multicast.BoswatchModule._MAGIC_WAKEUP_MSG = "###_MULTICAST_WAKEUP_###"
staticprotected

◆ _DEFAULT_TRIGGER_RIC

str module.multicast.BoswatchModule._DEFAULT_TRIGGER_RIC = "9999999"
staticprotected

◆ _my_frequencies

module.multicast.BoswatchModule._my_frequencies
protected

◆ instance_id

module.multicast.BoswatchModule.instance_id

◆ name

module.multicast.BoswatchModule.name

◆ _auto_clear_timeout

module.multicast.BoswatchModule._auto_clear_timeout
protected

◆ _hard_timeout

module.multicast.BoswatchModule._hard_timeout
protected

◆ _delimiter_rics

module.multicast.BoswatchModule._delimiter_rics
protected

◆ _text_rics

module.multicast.BoswatchModule._text_rics
protected

◆ _netident_rics

module.multicast.BoswatchModule._netident_rics
protected

◆ _trigger_ric

module.multicast.BoswatchModule._trigger_ric
protected

◆ _trigger_host

module.multicast.BoswatchModule._trigger_host
protected

◆ _trigger_port

module.multicast.BoswatchModule._trigger_port
protected

◆ _tone_ric_packets

module.multicast.BoswatchModule._tone_ric_packets
protected

◆ _last_tone_ric_time

module.multicast.BoswatchModule._last_tone_ric_time
protected

◆ _processing_text_ric

module.multicast.BoswatchModule._processing_text_ric
protected

◆ _processing_text_ric_started

module.multicast.BoswatchModule._processing_text_ric_started
protected

◆ _wildcards_registered

module.multicast.BoswatchModule._wildcards_registered
protected

◆ _packet_queue

module.multicast.BoswatchModule._packet_queue
protected

◆ _lock

module.multicast.BoswatchModule._lock
protected

◆ _queue_lock

module.multicast.BoswatchModule._queue_lock
protected

◆ _running

module.multicast.BoswatchModule._running
protected

◆ _cleanup_thread

module.multicast.BoswatchModule._cleanup_thread
protected

◆ _MAGIC_WAKEUP_MSG [2/2]

module.multicast.BoswatchModule._MAGIC_WAKEUP_MSG
protected