You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
283 lines
10 KiB
283 lines
10 KiB
###############################################################################
|
|
# Copyright (C) 2022 Simon Adlem, G7RZU <g7rzu@gb7fr.org.uk>
|
|
#
|
|
# This program is free software; you can redistribute it and/or modify
|
|
# it under the terms of the GNU General Public License as published by
|
|
# the Free Software Foundation; either version 3 of the License, or
|
|
# (at your option) any later version.
|
|
#
|
|
# This program is distributed in the hope that it will be useful,
|
|
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
|
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
|
# GNU General Public License for more details.
|
|
#
|
|
# You should have received a copy of the GNU General Public License
|
|
# along with this program; if not, write to the Free Software Foundation,
|
|
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
|
###############################################################################
|
|
|
|
# Does anybody read this stuff? There's a PEP somewhere that says I should do this.
|
|
__author__ = 'Simon Adlem - G7RZU'
|
|
__copyright__ = 'Copyright (c) Simon Adlem, G7RZU 2022'
|
|
__credits__ = ''
|
|
__license__ = 'GNU GPLv3'
|
|
__maintainer__ = 'Simon Adlem G7RZU'
|
|
__email__ = 'simon@gb7fr.org.uk'
|
|
|
|
#This is example code to connect to the report service in FreeDMR / HBLink3
|
|
#It can be used as a skeleton to build logging and monitoring tools.
|
|
|
|
import pickle
|
|
import threading
|
|
|
|
from twisted.internet import reactor
|
|
from twisted.internet.protocol import Factory, ReconnectingClientFactory
|
|
from twisted.protocols.basic import NetstringReceiver
|
|
|
|
import mysql.connector
|
|
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,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']:
|
|
self.configSend(data[1:])
|
|
elif data[:1] == REPORT_OPCODES['BRIDGE_SND']:
|
|
self.bridgeSend(data[1:])
|
|
elif data == b'bridge updated':
|
|
pass
|
|
else:
|
|
print('Unkown opcode - line:',data)
|
|
|
|
def bridgeEvent(self,data):
|
|
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:
|
|
if not self.db.is_connected():
|
|
try:
|
|
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'],event['source_server'],event['source_rptr']))
|
|
_cursor = self.db.cursor()
|
|
try:
|
|
_cursor.execute(
|
|
"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'],
|
|
event['trx'],
|
|
event['system'],
|
|
event['streamid'],
|
|
event['peerid'],
|
|
event['subid'],
|
|
event['slot'],
|
|
event['dstid'],
|
|
event['duration'],
|
|
event['source_server'],
|
|
event['source_rptr'],
|
|
)
|
|
)
|
|
self.db.commit()
|
|
except mysql.connector.Error as err:
|
|
print('(MYSQL) error, problem with cursor execute: {}'.format(err))
|
|
finally:
|
|
_cursor.close()
|
|
|
|
|
|
def bridgeSend(self,data):
|
|
self.BRIDGES = pickle.loads(data)
|
|
|
|
def configSend(self,data):
|
|
self.CONFIG = pickle.loads(data)
|
|
|
|
|
|
class reportClientFactory(ReconnectingClientFactory):
|
|
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.')
|
|
|
|
def buildProtocol(self, addr):
|
|
print('Connected.')
|
|
print('Resetting reconnection delay')
|
|
self.resetDelay()
|
|
return self.proto(self.db,self.reactor,self.relay)
|
|
|
|
def clientConnectionLost(self, connector, reason):
|
|
print('Lost connection. Reason:', reason)
|
|
ReconnectingClientFactory.clientConnectionLost(self, connector, reason)
|
|
|
|
def clientConnectionFailed(self, connector, reason):
|
|
print('Connection failed. Reason:', reason)
|
|
ReconnectingClientFactory.clientConnectionFailed(self, connector,reason)
|
|
|
|
if __name__ == '__main__':
|
|
import argparse
|
|
from twisted.internet import reactor
|
|
from setproctitle import setproctitle
|
|
import signal
|
|
import sys
|
|
import os
|
|
|
|
#Set process title early
|
|
setproctitle(__file__)
|
|
|
|
# Change the current directory to the location of the application
|
|
os.chdir(os.path.dirname(os.path.realpath(sys.argv[0])))
|
|
|
|
def sig_handler(_signal, _frame):
|
|
print('SHUTDOWN: TERMINATING WITH SIGNAL {}'.format(str(_signal)))
|
|
reactor.stop()
|
|
|
|
# Set signal handers so that we can gracefully exit if need be
|
|
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=args.db_host,
|
|
user=args.db_user,
|
|
password=args.db_password,
|
|
database=args.db_name,
|
|
#pool_name = "master",
|
|
#pool_size = 5
|
|
)
|
|
except mysql.connector.Error as err:
|
|
if err.errno == errorcode.ER_ACCESS_DENIED_ERROR:
|
|
sys.exit('(MYSQL) username or password error')
|
|
elif err.errno == errorcode.ER_BAD_DB_ERROR:
|
|
sys.exit('(MYSQL) DB Error')
|
|
else:
|
|
sys.exit('(MYSQL) error: %s',err)
|
|
|
|
|
|
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()
|