Add reporting relay and call origin metadata

master
Simon 2 months ago
parent 91a24e91b2
commit 40a70abea9

@ -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).

@ -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]))

@ -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.

@ -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

@ -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()

@ -6,3 +6,4 @@ configparser>=3.0.0
resettabletimer>=0.7.0
setproctitle
Pyro5
mysql-connector-python

@ -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)

@ -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 = []

@ -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
)
)

Loading…
Cancel
Save

Powered by TurnKey Linux.