From dc647ee0fd3c476da7321d7cfc14b583b9b5f698 Mon Sep 17 00:00:00 2001 From: Simon Date: Sun, 19 Jul 2026 00:39:33 +0000 Subject: [PATCH] Add reporting relay and call origin metadata --- README.md | 3 + bridge_master.py | 55 ++++++---- docs/reporting-relay.md | 32 ++++++ docs/v1x-codex-changelog.md | 10 ++ hblink.py | 30 +++++- report_sql.py | 162 ++++++++++++++++++++++------ requirements.txt | 1 + tests/harness/deterministic.py | 10 +- tests/test_auxiliary_tools.py | 104 ++++++++++++++++++ tests/test_deterministic_harness.py | 8 +- 10 files changed, 354 insertions(+), 61 deletions(-) create mode 100644 docs/reporting-relay.md diff --git a/README.md b/README.md index 83de5e4..4c2b028 100755 --- a/README.md +++ b/README.md @@ -5,3 +5,6 @@ FreeDMR Peer Server - software to assist in building a peer mesh network Please see the wiki for documentation. For packet harness tests and run commands, see [docs/testing.md](docs/testing.md). + +For the single-client SQL importer and dashboard reporting relay, see +[docs/reporting-relay.md](docs/reporting-relay.md). diff --git a/bridge_master.py b/bridge_master.py index 47ea7bb..d692ea1 100644 --- a/bridge_master.py +++ b/bridge_master.py @@ -1033,7 +1033,7 @@ def stream_trimmer_loop(): logger.info('(%s) *TIME OUT* RX STREAM ID: %s SUB: %s TGID %s, TS %s, Duration: %.2f', \ system, int_id(_slot['RX_STREAM_ID']), int_id(_slot['RX_RFS']), int_id(_slot['RX_TGID']), slot, _slot['RX_TIME'] - _slot['RX_START']) if CONFIG['REPORTS']['REPORT'] and _slot.get('RX_GROUP_VOICE_STREAM') and not _slot.get('RX_DATA_STREAM'): - systems[system]._report.send_bridgeEvent('GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}'.format(system, int_id(_slot['RX_STREAM_ID']), int_id(_slot['RX_PEER']), int_id(_slot['RX_RFS']), slot, int_id(_slot['RX_TGID']), _slot['RX_TIME'] - _slot['RX_START']).encode(encoding='utf-8', errors='ignore')) + systems[system]._report.send_bridgeEvent('GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}'.format(system, int_id(_slot['RX_STREAM_ID']), int_id(_slot['RX_PEER']), int_id(_slot['RX_RFS']), slot, int_id(_slot['RX_TGID']), _slot['RX_TIME'] - _slot['RX_START']).encode(encoding='utf-8', errors='ignore'), _slot.get('RX_SOURCE_SERVER'), _slot.get('RX_SOURCE_RPTR')) _slot['RX_DATA_STREAM'] = False _slot['RX_GROUP_VOICE_STREAM'] = False #Null stream_id - for loop control @@ -1048,7 +1048,7 @@ def stream_trimmer_loop(): logger.debug('(%s) *TIME OUT* TX STREAM ID: %s SUB: %s TGID %s, TS %s, Duration: %.2f', \ system, int_id(_slot['TX_STREAM_ID']), int_id(_slot['TX_RFS']), int_id(_slot['TX_TGID']), slot, _slot['TX_TIME'] - _slot['TX_START']) if CONFIG['REPORTS']['REPORT'] and not _slot.get('TX_DATA_STREAM'): - systems[system]._report.send_bridgeEvent('GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}'.format(system, int_id(_slot['TX_STREAM_ID']), int_id(_slot['TX_PEER']), int_id(_slot['TX_RFS']), slot, int_id(_slot['TX_TGID']), _slot['TX_TIME'] - _slot['TX_START']).encode(encoding='utf-8', errors='ignore')) + systems[system]._report.send_bridgeEvent('GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}'.format(system, int_id(_slot['TX_STREAM_ID']), int_id(_slot['TX_PEER']), int_id(_slot['TX_RFS']), slot, int_id(_slot['TX_TGID']), _slot['TX_TIME'] - _slot['TX_START']).encode(encoding='utf-8', errors='ignore'), _slot.get('TX_SOURCE_SERVER'), _slot.get('TX_SOURCE_RPTR')) _slot['TX_DATA_STREAM'] = False # OBP systems @@ -1088,7 +1088,7 @@ def stream_trimmer_loop(): system, int_id(stream_id), get_alias(int_id(_stream['RFS']), subscriber_ids), get_alias(int_id(_stream['RX_PEER']), peer_ids), get_alias(int_id(_stream['TGID']), talkgroup_ids), _stream['LAST'] - _stream['START']) if CONFIG['REPORTS']['REPORT'] and not _stream.get('DATA_STREAM'): - systems[system]._report.send_bridgeEvent('GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}'.format(system, int_id(stream_id), int_id(_stream['RX_PEER']), int_id(_stream['RFS']), 1, int_id(_stream['TGID']), _stream['LAST'] - _stream['START']).encode(encoding='utf-8', errors='ignore')) + systems[system]._report.send_bridgeEvent('GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}'.format(system, int_id(stream_id), int_id(_stream['RX_PEER']), int_id(_stream['RFS']), 1, int_id(_stream['TGID']), _stream['LAST'] - _stream['START']).encode(encoding='utf-8', errors='ignore'), _stream.get('SOURCE_SERVER'), _stream.get('SOURCE_RPTR')) systems[system].STATUS[stream_id]['_to'] = True continue except Exception as e: @@ -1824,7 +1824,7 @@ class routerOBP(OPENBRIDGE): logger.debug('(%s) Conference Bridge: %s, Call Bridged to OBP System: %s TS: %s, TGID: %s', self._system, _bridge, _target['SYSTEM'], _target['TS'], int_id(_target['TGID'])) if CONFIG['REPORTS']['REPORT']: if not _data_control: - systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,START,TX,{},{},{},{},{},{}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore')) + systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,START,TX,{},{},{},{},{},{}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) else: systems[_target['SYSTEM']]._report.send_bridgeEvent('{},DATA,TX,{},{},{},{},{},{}'.format(_data_event, _target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore')) @@ -1852,7 +1852,7 @@ class routerOBP(OPENBRIDGE): logger.warning('(%s) KeyError - T_LC, Skipping', self._system) if CONFIG['REPORTS']['REPORT'] and not _target_status[_stream_id].get('DATA_STREAM'): call_duration = pkt_time - _target_status[_stream_id]['START'] - systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID']), call_duration).encode(encoding='utf-8', errors='ignore')) + systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID']), call_duration).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) if not _target_status[_stream_id].get('DATA_STREAM'): _target_status[_stream_id]['_fin'] = True # Create a Burst B-E packet (Embedded LC) @@ -1907,6 +1907,8 @@ class routerOBP(OPENBRIDGE): _target_status[_target['TS']]['TX_RFS'] = _rf_src _target_status[_target['TS']]['TX_PEER'] = _peer_id _target_status[_target['TS']]['TX_DATA_STREAM'] = _data_control + _target_status[_target['TS']]['TX_SOURCE_SERVER'] = _source_server + _target_status[_target['TS']]['TX_SOURCE_RPTR'] = _source_rptr # Generate LCs (full and EMB) for the TX stream dst_lc = b''.join([self.STATUS[_stream_id]['LC'][0:3], _target['TGID'], _rf_src]) _target_status[_target['TS']]['TX_H_LC'] = dmr_codec.encode_header_lc(dst_lc) @@ -1916,7 +1918,7 @@ class routerOBP(OPENBRIDGE): logger.debug('(%s) Conference Bridge: %s, Call Bridged to HBP System: %s TS: %s, TGID: %s', self._system, _bridge, _target['SYSTEM'], _target['TS'], int_id(_target['TGID'])) if CONFIG['REPORTS']['REPORT']: if not _data_control: - systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,START,TX,{},{},{},{},{},{}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore')) + systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,START,TX,{},{},{},{},{},{}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) else: systems[_target['SYSTEM']]._report.send_bridgeEvent('{},DATA,TX,{},{},{},{},{},{}'.format(_data_event, _target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore')) @@ -1946,7 +1948,7 @@ class routerOBP(OPENBRIDGE): dmrbits = _target_status[_target['TS']]['TX_T_LC'][0:98] + dmrbits[98:166] + _target_status[_target['TS']]['TX_T_LC'][98:197] if CONFIG['REPORTS']['REPORT'] and not _target_status[_target['TS']].get('TX_DATA_STREAM'): call_duration = pkt_time - _target_status[_target['TS']]['TX_START'] - systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID']), call_duration).encode(encoding='utf-8', errors='ignore')) + systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID']), call_duration).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) # Create a Burst B-E packet (Embedded LC) elif target_requires_emb_lc_rewrite(_dst_id, _target['TGID']) and _frame_type == HBPF_VOICE and _dtype_vseq in [1,2,3,4]: dmrbits = dmrbits[0:116] + _target_status[_target['TS']]['TX_EMB_LC'][_dtype_vseq] + dmrbits[148:264] @@ -2200,7 +2202,9 @@ class routerOBP(OPENBRIDGE): 'packets': 0, 'loss': 0, 'crcs': set(), - 'DATA_STREAM': _data_control + 'DATA_STREAM': _data_control, + 'SOURCE_SERVER': _source_server, + 'SOURCE_RPTR': _source_rptr } @@ -2222,7 +2226,7 @@ class routerOBP(OPENBRIDGE): self._system, _rx_event, int_id(_stream_id),get_alias(_rf_src, subscriber_ids),int_id(_rf_src),self.get_rptr(_source_rptr), int_id(_source_rptr), get_alias(_peer_id, peer_ids), int_id(_peer_id), get_alias(_dst_id, talkgroup_ids), int_id(_dst_id), _slot,int_id(_source_server),_inthops) if CONFIG['REPORTS']['REPORT']: if not _data_control: - self._report.send_bridgeEvent('GROUP VOICE,START,RX,{},{},{},{},{},{}'.format(self._system, int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _slot, int_id(_dst_id)).encode(encoding='utf-8', errors='ignore')) + self._report.send_bridgeEvent('GROUP VOICE,START,RX,{},{},{},{},{},{}'.format(self._system, int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _slot, int_id(_dst_id)).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) else: self._report.send_bridgeEvent('{},DATA,RX,{},{},{},{},{},{}'.format(_data_event, self._system, int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _slot, int_id(_dst_id)).encode(encoding='utf-8', errors='ignore')) @@ -2368,7 +2372,7 @@ class routerOBP(OPENBRIDGE): self._system, int_id(_stream_id), get_alias(_rf_src, subscriber_ids), int_id(_rf_src), get_alias(_peer_id, peer_ids), int_id(_peer_id), get_alias(_dst_id, talkgroup_ids), int_id(_dst_id), _slot, call_duration, packet_rate,loss) if not self.STATUS[_stream_id].get('DATA_STREAM'): if CONFIG['REPORTS']['REPORT']: - self._report.send_bridgeEvent('GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}'.format(self._system, int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _slot, int_id(_dst_id), call_duration).encode(encoding='utf-8', errors='ignore')) + self._report.send_bridgeEvent('GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}'.format(self._system, int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _slot, int_id(_dst_id), call_duration).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) self.STATUS[_stream_id]['_fin'] = True self.STATUS[_stream_id]['lastSeq'] = False @@ -2527,7 +2531,7 @@ class routerHBP(HBSYSTEM): logger.debug('(%s) Conference Bridge: %s, Call Bridged to OBP System: %s TS: %s, TGID: %s', self._system, _bridge, _target['SYSTEM'], _target['TS'], int_id(_target['TGID'])) if CONFIG['REPORTS']['REPORT']: if not _data_control: - systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,START,TX,{},{},{},{},{},{}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore')) + systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,START,TX,{},{},{},{},{},{}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) else: systems[_target['SYSTEM']]._report.send_bridgeEvent('{},DATA,TX,{},{},{},{},{},{}'.format(_data_event, _target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore')) @@ -2552,7 +2556,7 @@ class routerHBP(HBSYSTEM): dmrbits = _target_status[_stream_id]['T_LC'][0:98] + dmrbits[98:166] + _target_status[_stream_id]['T_LC'][98:197] if CONFIG['REPORTS']['REPORT'] and not _target_status[_stream_id].get('DATA_STREAM'): call_duration = pkt_time - _target_status[_stream_id]['START'] - systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID']), call_duration).encode(encoding='utf-8', errors='ignore')) + systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID']), call_duration).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) if not _target_status[_stream_id].get('DATA_STREAM'): _target_status[_stream_id]['_fin'] = True # Create a Burst B-E packet (Embedded LC) @@ -2599,6 +2603,8 @@ class routerHBP(HBSYSTEM): _target_status[_target['TS']]['TX_RFS'] = _rf_src _target_status[_target['TS']]['TX_PEER'] = _peer_id _target_status[_target['TS']]['TX_DATA_STREAM'] = _data_control + _target_status[_target['TS']]['TX_SOURCE_SERVER'] = _source_server + _target_status[_target['TS']]['TX_SOURCE_RPTR'] = _source_rptr # Generate LCs (full and EMB) for the TX stream dst_lc = b''.join([self.STATUS[_slot]['RX_LC'][0:3],_target['TGID'],_rf_src]) _target_status[_target['TS']]['TX_H_LC'] = dmr_codec.encode_header_lc(dst_lc) @@ -2608,7 +2614,7 @@ class routerHBP(HBSYSTEM): logger.debug('(%s) Conference Bridge: %s, Call Bridged to HBP System: %s TS: %s, TGID: %s', self._system, _bridge, _target['SYSTEM'], _target['TS'], int_id(_target['TGID'])) if CONFIG['REPORTS']['REPORT']: if not _data_control: - systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,START,TX,{},{},{},{},{},{}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore')) + systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,START,TX,{},{},{},{},{},{}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) else: systems[_target['SYSTEM']]._report.send_bridgeEvent('{},DATA,TX,{},{},{},{},{},{}'.format(_data_event, _target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID'])).encode(encoding='utf-8', errors='ignore')) @@ -2638,7 +2644,7 @@ class routerHBP(HBSYSTEM): dmrbits = _target_status[_target['TS']]['TX_T_LC'][0:98] + dmrbits[98:166] + _target_status[_target['TS']]['TX_T_LC'][98:197] if CONFIG['REPORTS']['REPORT'] and not _target_status[_target['TS']].get('TX_DATA_STREAM'): call_duration = pkt_time - _target_status[_target['TS']]['TX_START'] - systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID']), call_duration).encode(encoding='utf-8', errors='ignore')) + systems[_target['SYSTEM']]._report.send_bridgeEvent('GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}'.format(_target['SYSTEM'], int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _target['TS'], int_id(_target['TGID']), call_duration).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) # Create a Burst B-E packet (Embedded LC) elif target_requires_emb_lc_rewrite(_dst_id, _target['TGID']) and _frame_type == HBPF_VOICE and _dtype_vseq in [1,2,3,4]: dmrbits = dmrbits[0:116] + _target_status[_target['TS']]['TX_EMB_LC'][_dtype_vseq] + dmrbits[148:264] @@ -3147,6 +3153,8 @@ class routerHBP(HBSYSTEM): self.STATUS[_slot]['RX_START'] = pkt_time self.STATUS[_slot]['RX_DATA_STREAM'] = _data_control self.STATUS[_slot]['RX_GROUP_VOICE_STREAM'] = not _data_control + self.STATUS[_slot]['RX_SOURCE_SERVER'] = _source_server + self.STATUS[_slot]['RX_SOURCE_RPTR'] = _source_rptr if _call_type == 'vcsbk' and _data_event == 'OTHER DATA': logger.info('(%s) *VCSBK* STREAM ID: %s SUB: %s (%s) PEER: %s (%s) TGID %s (%s), TS %s _dtype_vseq: %s', @@ -3162,7 +3170,7 @@ class routerHBP(HBSYSTEM): logger.info('(%s) *CALL START* STREAM ID: %s SUB: %s (%s) PEER: %s (%s) TGID %s (%s), TS %s', \ self._system, int_id(_stream_id), get_alias(_rf_src, subscriber_ids), int_id(_rf_src), get_alias(_peer_id, peer_ids), int_id(_peer_id), get_alias(_dst_id, talkgroup_ids), int_id(_dst_id), _slot) if CONFIG['REPORTS']['REPORT']: - self._report.send_bridgeEvent('GROUP VOICE,START,RX,{},{},{},{},{},{}'.format(self._system, int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _slot, int_id(_dst_id)).encode(encoding='utf-8', errors='ignore')) + self._report.send_bridgeEvent('GROUP VOICE,START,RX,{},{},{},{},{},{}'.format(self._system, int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _slot, int_id(_dst_id)).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) # If we can, use the LC from the voice header as to keep all options intact if _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD: @@ -3287,7 +3295,7 @@ class routerHBP(HBSYSTEM): logger.info('(%s) *CALL END* STREAM ID: %s SUB: %s (%s) PEER: %s (%s) TGID %s (%s), TS %s, Duration: %.2f, Packet rate: %.2f/s, LOSS: %.2f%%', \ self._system, int_id(_stream_id), get_alias(_rf_src, subscriber_ids), int_id(_rf_src), get_alias(_peer_id, peer_ids), int_id(_peer_id), get_alias(_dst_id, talkgroup_ids), int_id(_dst_id), _slot, call_duration, packet_rate, loss) if CONFIG['REPORTS']['REPORT']: - self._report.send_bridgeEvent('GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}'.format(self._system, int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _slot, int_id(_dst_id), call_duration).encode(encoding='utf-8', errors='ignore')) + self._report.send_bridgeEvent('GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}'.format(self._system, int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _slot, int_id(_dst_id), call_duration).encode(encoding='utf-8', errors='ignore'), _source_server, _source_rptr) self.STATUS[_slot]['RX_FINISHED_STREAM_ID'] = _stream_id self.STATUS[_slot]['RX_FINISHED_STREAM_LOG'] = False @@ -3396,13 +3404,20 @@ class routerHBP(HBSYSTEM): # class bridgeReportFactory(reportFactory): - def send_bridge(self): + def send_bridge(self, client=None): serialized = pickle.dumps(BRIDGES, protocol=2) #.decode("utf-8", errors='ignore') - self.send_clients(b''.join([REPORT_OPCODES['BRIDGE_SND'],serialized])) + message = b''.join([REPORT_OPCODES['BRIDGE_SND'],serialized]) + if client is None: + self.send_clients(message) + else: + self.send_client(client, message) - def send_bridgeEvent(self, _data): + def send_bridgeEvent(self, _data, _source_server=None, _source_rptr=None): if isinstance(_data, str): - _data = _data.decode('utf-8', error='ignore') + _data = _data.encode('utf-8', errors='ignore') + if _source_server is not None and _source_rptr is not None: + origin = '{},{}'.format(int_id(_source_server), int_id(_source_rptr)).encode('ascii') + _data = b','.join([_data, origin]) self.send_clients(b''.join([REPORT_OPCODES['BRDG_EVENT'],_data])) diff --git a/docs/reporting-relay.md b/docs/reporting-relay.md new file mode 100644 index 0000000..94ede4b --- /dev/null +++ b/docs/reporting-relay.md @@ -0,0 +1,32 @@ +# Reporting Relay + +`report_sql.py` can be the sole client of the FreeDMR TCP reporting socket and +provide a compatible downstream socket for the live dashboard. This avoids +placing multiple reporting consumers on the live server while retaining the +existing netstring protocol. + +The existing six positional arguments are unchanged. Enable the relay with: + +```sh +python report_sql.py freedmr 4321 dbhost dbuser dbpassword logs \ + --relay-interface 127.0.0.1 --relay-port 4322 +``` + +Configure FDMR Monitor to connect to relay port `4322`. Use an explicit +container-facing interface instead of `127.0.0.1` when the dashboard runs in a +different container, and restrict access to that interface at the network +boundary. Report configuration and bridge-state messages contain trusted +Python pickle data and must not be exposed to untrusted clients. + +The relay forwards upstream messages unchanged and caches only the latest +configuration and bridge-state messages so a newly connected dashboard can be +initialized immediately. SQL storage accepts both legacy events and extended +events with trailing origin fields: + +```text +START: type,event,trx,system,streamid,peerid,subid,slot,dstid,source_server,source_rptr +END: type,event,trx,system,streamid,peerid,subid,slot,dstid,duration,source_server,source_rptr +``` + +Legacy events remain valid and are stored with null origin columns. The relay +is optional; omitting `--relay-port` preserves the previous importer behavior. diff --git a/docs/v1x-codex-changelog.md b/docs/v1x-codex-changelog.md index 5febec9..a8cd71b 100644 --- a/docs/v1x-codex-changelog.md +++ b/docs/v1x-codex-changelog.md @@ -1,5 +1,15 @@ # FreeDMR 1.x Changelog +## Reporting + +- Added source server and source repeater metadata as trailing fields on group + voice lifecycle reports without changing the legacy field positions. +- Extended `report_sql.py` to accept legacy and origin-aware events, write the + existing Global-LH source columns, and optionally expose a transparent, + cached downstream report socket for the live dashboard. +- Isolated failed reporting clients and made configuration/bridge requests + reply only to the requesting client. + ## Test Harnesses - Added an in-process deterministic packet harness for `bridge_master.py` diff --git a/hblink.py b/hblink.py index feac741..5023d6a 100755 --- a/hblink.py +++ b/hblink.py @@ -1295,7 +1295,7 @@ class report(NetstringReceiver): def connectionLost(self, reason): logger.info('(REPORT) HBlink reporting client disconnected: %s', self.transport.getPeer()) - self._factory.clients.remove(self) + self._factory.remove_client(self) def stringReceived(self, data): self.process_message(data) @@ -1304,7 +1304,10 @@ class report(NetstringReceiver): opcode = _message[:1] if opcode == REPORT_OPCODES['CONFIG_REQ']: logger.info('(REPORT) HBlink reporting client sent \'CONFIG_REQ\': %s', self.transport.getPeer()) - self.send_config() + self._factory.send_config(self) + elif opcode == REPORT_OPCODES['BRIDGE_REQ'] and hasattr(self._factory, 'send_bridge'): + logger.info('(REPORT) HBlink reporting client sent \'BRIDGE_REQ\': %s', self.transport.getPeer()) + self._factory.send_bridge(self) else: logger.error('(REPORT) got unknown opcode') @@ -1321,13 +1324,30 @@ class reportFactory(Factory): return None def send_clients(self, _message): - for client in self.clients: + for client in tuple(self.clients): + self.send_client(client, _message) + + def send_client(self, client, _message): + try: client.sendString(_message) + except Exception as err: + logger.warning('(REPORT) Dropping reporting client after send failure: %s', err) + self.remove_client(client) - def send_config(self): + def remove_client(self, client): + try: + self.clients.remove(client) + except ValueError: + pass + + def send_config(self, client=None): serialized = pickle.dumps(self._config['SYSTEMS'], protocol=2) #.decode('utf-8', errors='ignore') #pickle.HIGHEST_PROTOCOL) logger.debug('(REPORT) Send config') - self.send_clients(b''.join([REPORT_OPCODES['CONFIG_SND'], serialized])) + message = b''.join([REPORT_OPCODES['CONFIG_SND'], serialized]) + if client is None: + self.send_clients(message) + else: + self.send_client(client, message) #Read list of listed servers from CSV (actually TSV) file diff --git a/report_sql.py b/report_sql.py index 66824d7..983e4c6 100644 --- a/report_sql.py +++ b/report_sql.py @@ -31,7 +31,7 @@ import pickle import threading from twisted.internet import reactor -from twisted.internet.protocol import ReconnectingClientFactory +from twisted.internet.protocol import Factory, ReconnectingClientFactory from twisted.protocols.basic import NetstringReceiver import mysql.connector @@ -39,14 +39,97 @@ from mysql.connector import errorcode from reporting_const import * + +def parse_bridge_event(data): + """Parse legacy events and the extended origin-aware event shape.""" + fields = data.split(',') + if len(fields) < 9: + raise ValueError('bridge event has fewer than 9 fields') + + event = { + 'type': fields[0], + 'event': fields[1], + 'trx': fields[2], + 'system': fields[3], + 'streamid': fields[4], + 'peerid': fields[5], + 'subid': fields[6], + 'slot': fields[7], + 'dstid': fields[8], + 'duration': 0, + 'source_server': None, + 'source_rptr': None, + } + metadata_offset = 9 + if event['event'] == 'END': + if len(fields) < 10: + raise ValueError('END bridge event has no duration') + event['duration'] = fields[9] + metadata_offset = 10 + + if len(fields) >= metadata_offset + 2: + event['source_server'] = fields[metadata_offset] or None + event['source_rptr'] = fields[metadata_offset + 1] or None + return event + + +class reportRelay(NetstringReceiver): + """Downstream endpoint compatible with the FreeDMR report socket.""" + def connectionMade(self): + self.factory.add_client(self) + + def connectionLost(self, reason): + self.factory.clients.discard(self) + + def stringReceived(self, data): + if data[:1] == REPORT_OPCODES['CONFIG_REQ']: + self.factory.send_cached(self, REPORT_OPCODES['CONFIG_SND']) + elif data[:1] == REPORT_OPCODES['BRIDGE_REQ']: + self.factory.send_cached(self, REPORT_OPCODES['BRIDGE_SND']) + + +class reportRelayFactory(Factory): + protocol = reportRelay + + def __init__(self): + self.clients = set() + self._cached = {} + + def add_client(self, client): + self.clients.add(client) + self.send_cached(client, REPORT_OPCODES['CONFIG_SND']) + self.send_cached(client, REPORT_OPCODES['BRIDGE_SND']) + + def send_cached(self, client, opcode): + message = self._cached.get(opcode) + if message is not None: + self._send(client, message) + + def relay(self, message): + opcode = message[:1] + if opcode in (REPORT_OPCODES['CONFIG_SND'], REPORT_OPCODES['BRIDGE_SND']): + self._cached[opcode] = message + for client in tuple(self.clients): + self._send(client, message) + + def _send(self, client, message): + try: + client.sendString(message) + except Exception as err: + self.clients.discard(client) + print('(RELAY) dropping report client after send failure: {}'.format(err)) + class reportClient(NetstringReceiver): - def __init__(self,db,reactor): + def __init__(self,db,reactor,relay=None): self.db = db self.reactor = reactor + self.relay = relay self._db_lock = threading.Lock() def stringReceived(self, data): - + if self.relay is not None: + self.relay.relay(data) + if data[:1] == REPORT_OPCODES['BRDG_EVENT']: self.bridgeEvent(data[1:].decode('UTF-8')) elif data[:1] == REPORT_OPCODES['CONFIG_SND']: @@ -59,39 +142,30 @@ class reportClient(NetstringReceiver): print('Unkown opcode - line:',data) def bridgeEvent(self,data): - datalist = data.split(',') - event = { - 'type' : datalist[0], - 'event' : datalist[1], - 'trx' : datalist[2], - 'system' : datalist[3], - 'streamid' : datalist[4], - 'peerid' : datalist[5], - 'subid' : datalist[6], - 'slot' : datalist[7], - 'dstid' : datalist[8], - 'duration' : 0 - } - - if len(datalist) > 9: - event['duration'] = datalist[9] - - + try: + event = parse_bridge_event(data) + except ValueError as err: + print('(REPORT) ignoring malformed bridge event: {}'.format(err)) + return self.reactor.callInThread(self.send_mysql,event) def send_mysql(self,event): with self._db_lock: - while not self.db.is_connected(): + if not self.db.is_connected(): try: - self.db.reconnect() + self.db.reconnect(attempts=3, delay=1) except mysql.connector.Error as err: print('(MYSQL) error on reconnect: {}'.format(err)) + return + if not self.db.is_connected(): + print('(MYSQL) reconnect did not restore the connection') + return - print("{} {} {} {} {} {} {} {} {}".format(event['type'],event['event'], event['trx'],event['system'],event['streamid'],event['peerid'],event['subid'],event['slot'],event['dstid'],event['duration'])) + print("{} {} {} {} {} {} {} {} {} {} {} {}".format(event['type'],event['event'], event['trx'],event['system'],event['streamid'],event['peerid'],event['subid'],event['slot'],event['dstid'],event['duration'],event['source_server'],event['source_rptr'])) _cursor = self.db.cursor() try: _cursor.execute( - "insert into feed values (NULL,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)", + "insert into feed (type,event,trx,system,streamid,peerid,subid,slot,dstid,duration,source_server,source_rptr) values (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)", ( event['type'], event['event'], @@ -103,6 +177,8 @@ class reportClient(NetstringReceiver): event['slot'], event['dstid'], event['duration'], + event['source_server'], + event['source_rptr'], ) ) self.db.commit() @@ -120,10 +196,11 @@ class reportClient(NetstringReceiver): class reportClientFactory(ReconnectingClientFactory): - def __init__(self,proto,db,reactor): + def __init__(self,proto,db,reactor,relay=None): self.proto = proto self.db = db self.reactor = reactor + self.relay = relay def startedConnecting(self, connector): print('Started to connect.') @@ -132,7 +209,7 @@ class reportClientFactory(ReconnectingClientFactory): print('Connected.') print('Resetting reconnection delay') self.resetDelay() - return self.proto(self.db,self.reactor) + return self.proto(self.db,self.reactor,self.relay) def clientConnectionLost(self, connector, reason): print('Lost connection. Reason:', reason) @@ -143,7 +220,7 @@ class reportClientFactory(ReconnectingClientFactory): ReconnectingClientFactory.clientConnectionFailed(self, connector,reason) if __name__ == '__main__': - + import argparse from twisted.internet import reactor from setproctitle import setproctitle import signal @@ -164,12 +241,25 @@ if __name__ == '__main__': for sig in [signal.SIGINT, signal.SIGTERM]: signal.signal(sig, sig_handler) + parser = argparse.ArgumentParser(description='Persist and optionally relay FreeDMR reports') + parser.add_argument('host', help='FreeDMR report socket host') + parser.add_argument('port', type=int, help='FreeDMR report socket port') + parser.add_argument('db_host') + parser.add_argument('db_user') + parser.add_argument('db_password') + parser.add_argument('db_name') + parser.add_argument('--relay-port', type=int, default=0, + help='downstream report port for the dashboard (disabled by default)') + parser.add_argument('--relay-interface', default='127.0.0.1', + help='interface for the downstream relay (default: 127.0.0.1)') + args = parser.parse_args() + try: db = mysql.connector.connect( - host=sys.argv[3], - user=sys.argv[4], - password=sys.argv[5], - database=sys.argv[6], + host=args.db_host, + user=args.db_user, + password=args.db_password, + database=args.db_name, #pool_name = "master", #pool_size = 5 ) @@ -182,5 +272,11 @@ if __name__ == '__main__': sys.exit('(MYSQL) error: %s',err) - reactor.connectTCP(sys.argv[1],int(sys.argv[2]), reportClientFactory(reportClient,db,reactor)) + relay = None + if args.relay_port: + relay = reportRelayFactory() + reactor.listenTCP(args.relay_port, relay, interface=args.relay_interface) + print('(RELAY) listening on {}:{}'.format(args.relay_interface, args.relay_port)) + + reactor.connectTCP(args.host, args.port, reportClientFactory(reportClient,db,reactor,relay)) reactor.run() diff --git a/requirements.txt b/requirements.txt index 69123e4..159ca21 100755 --- a/requirements.txt +++ b/requirements.txt @@ -6,3 +6,4 @@ configparser>=3.0.0 resettabletimer>=0.7.0 setproctitle Pyro5 +mysql-connector-python diff --git a/tests/harness/deterministic.py b/tests/harness/deterministic.py index 12d13e1..1783792 100644 --- a/tests/harness/deterministic.py +++ b/tests/harness/deterministic.py @@ -224,7 +224,15 @@ class ReportCapture: def __init__(self) -> None: self.events: list[bytes] = [] - def send_bridgeEvent(self, data: bytes) -> None: + def send_bridgeEvent( + self, + data: bytes, + source_server: bytes | None = None, + source_rptr: bytes | None = None, + ) -> None: + if source_server is not None and source_rptr is not None: + origin = f"{int_id(source_server)},{int_id(source_rptr)}".encode("ascii") + data = b",".join((data, origin)) self.events.append(data) diff --git a/tests/test_auxiliary_tools.py b/tests/test_auxiliary_tools.py index 6411dc1..d446529 100644 --- a/tests/test_auxiliary_tools.py +++ b/tests/test_auxiliary_tools.py @@ -65,6 +65,8 @@ class AuxiliaryToolTests(unittest.TestCase): "slot": "2", "dstid": "2350", "duration": "0", + "source_server": None, + "source_rptr": None, } with redirect_stdout(io.StringIO()): client.send_mysql(event) @@ -73,9 +75,89 @@ class AuxiliaryToolTests(unittest.TestCase): self.assertIn("%s", statement) self.assertEqual(params[0], "GROUP VOICE") self.assertEqual(params[8], "2350") + self.assertIsNone(params[10]) + self.assertIsNone(params[11]) self.assertTrue(fake_db.committed) self.assertTrue(fake_db.cursor_obj.closed) + def test_report_sql_parses_legacy_and_extended_events(self): + self._install_mysql_stub() + import report_sql + + report_sql = importlib.reload(report_sql) + legacy = report_sql.parse_bridge_event( + "GROUP VOICE,END,RX,SYSTEM,1234,5678,9012,2,2350,4.25" + ) + extended = report_sql.parse_bridge_event( + "GROUP VOICE,END,RX,SYSTEM,1234,5678,9012,2,2350,4.25,9991,1001" + ) + start = report_sql.parse_bridge_event( + "GROUP VOICE,START,RX,SYSTEM,1234,5678,9012,2,2350,9991,1001" + ) + + self.assertEqual(legacy["duration"], "4.25") + self.assertIsNone(legacy["source_server"]) + self.assertEqual(extended["source_server"], "9991") + self.assertEqual(extended["source_rptr"], "1001") + self.assertEqual(start["duration"], 0) + self.assertEqual(start["source_server"], "9991") + + def test_report_sql_relay_caches_state_and_isolates_clients(self): + self._install_mysql_stub() + import report_sql + + report_sql = importlib.reload(report_sql) + relay = report_sql.reportRelayFactory() + good = _FakeReportClient() + failed = _FakeReportClient(fail=True) + relay.clients.update((failed, good)) + + config = report_sql.REPORT_OPCODES["CONFIG_SND"] + b"config" + bridge = report_sql.REPORT_OPCODES["BRIDGE_SND"] + b"bridge" + event = report_sql.REPORT_OPCODES["BRDG_EVENT"] + b"event" + relay.relay(config) + relay.relay(bridge) + relay.relay(event) + + self.assertEqual(good.messages, [config, bridge, event]) + self.assertNotIn(failed, relay.clients) + newcomer = _FakeReportClient() + relay.add_client(newcomer) + self.assertEqual(newcomer.messages, [config, bridge]) + + def test_report_factory_isolates_failed_clients_and_targets_replies(self): + import hblink + + factory = hblink.reportFactory({ + "REPORTS": {"REPORT_CLIENTS": ["*"]}, + "SYSTEMS": {"MASTER-A": {"MODE": "MASTER"}}, + }) + failed = _FakeReportClient(fail=True) + broadcast_client = _FakeReportClient() + request_client = _FakeReportClient() + factory.clients = [failed, broadcast_client] + + factory.send_clients(b"event") + factory.send_config(request_client) + + self.assertEqual(broadcast_client.messages, [b"event"]) + self.assertNotIn(failed, factory.clients) + self.assertEqual(len(request_client.messages), 1) + self.assertEqual(request_client.messages[0][:1], hblink.REPORT_OPCODES["CONFIG_SND"]) + self.assertEqual(broadcast_client.messages, [b"event"]) + + def test_report_sql_failed_reconnect_returns_without_spinning(self): + self._install_mysql_stub() + import report_sql + + report_sql = importlib.reload(report_sql) + db = _FakeDisconnectedDB() + client = report_sql.reportClient(db, object()) + with redirect_stdout(io.StringIO()): + client.send_mysql({}) + + self.assertEqual(db.reconnects, 1) + def test_proxy_environment_bool_parser(self): hotspot_proxy_v2, saved_modules = self._import_proxy_module() try: @@ -596,6 +678,28 @@ class _FakeDB: self.committed = True +class _FakeReportClient: + def __init__(self, fail=False): + self.fail = fail + self.messages = [] + + def sendString(self, message): + if self.fail: + raise RuntimeError("closed") + self.messages.append(message) + + +class _FakeDisconnectedDB: + def __init__(self): + self.reconnects = 0 + + def is_connected(self): + return False + + def reconnect(self, attempts, delay): + self.reconnects += 1 + + class _FakeTimer: def __init__(self): self.resets = [] diff --git a/tests/test_deterministic_harness.py b/tests/test_deterministic_harness.py index c4df88d..174ec02 100644 --- a/tests/test_deterministic_harness.py +++ b/tests/test_deterministic_harness.py @@ -3380,7 +3380,7 @@ NETWORK_DIRECT_DIAL_SLOT: {value} scenario.systems["OBP-1"].STATUS[bytes_4(stream_id)], ) self.assertIn( - b"GROUP VOICE,END,TX,OBP-1,16909060,1001,3120001,1,91,0.00", + b"GROUP VOICE,END,TX,OBP-1,16909060,1001,3120001,1,91,0.00,9990,1001", scenario.reports["OBP-1"].events, ) self.assertFalse( @@ -3436,7 +3436,7 @@ NETWORK_DIRECT_DIAL_SLOT: {value} scenario.systems["OBP-2"].STATUS[bytes_4(stream_id)], ) self.assertIn( - b"GROUP VOICE,END,TX,OBP-2,16909060,1001,3120001,1,91,0.00", + b"GROUP VOICE,END,TX,OBP-2,16909060,1001,3120001,1,91,0.00,9990,0", scenario.reports["OBP-2"].events, ) self.assertFalse( @@ -3639,24 +3639,28 @@ NETWORK_DIRECT_DIAL_SLOT: {value} self.assertTrue( any( event.startswith(b"GROUP VOICE,START,RX,MASTER-A") + and event.endswith(b",9990,1001") for event in scenario.reports["MASTER-A"].events ) ) self.assertTrue( any( event.startswith(b"GROUP VOICE,END,RX,MASTER-A") + and event.endswith(b",9990,1001") for event in scenario.reports["MASTER-A"].events ) ) self.assertTrue( any( event.startswith(b"GROUP VOICE,START,TX,MASTER-B") + and event.endswith(b",9990,1001") for event in scenario.reports["MASTER-B"].events ) ) self.assertTrue( any( event.startswith(b"GROUP VOICE,END,TX,MASTER-B") + and event.endswith(b",9990,1001") for event in scenario.reports["MASTER-B"].events ) )