bugfix pass

master
Simon 4 months ago
parent f23a9ade26
commit 426feb0db0

6
.gitignore vendored

@ -81,6 +81,8 @@ celerybeat-schedule
# virtualenv
venv/
ENV/
.venv/
pyvenv.cfg
# Spyder project settings
.spyderproject
@ -93,6 +95,10 @@ hblink.cfg
*.config
*.bak
rules.py
bridge.pkl
config.pkl
sub_map.pkl
keys.json
subscriber_ids.*
local_subscriber_ids.*
peer_ids.*

@ -1,37 +1,46 @@
import sys
from time import time
import logging
from twisted.internet import reactor,task
from twisted.internet.defer import Deferred
from twisted.internet.protocol import ClientFactory,ClientFactory,Protocol
from twisted.internet import reactor
from twisted.internet.protocol import ClientFactory
from twisted.protocols.basic import LineReceiver
logger = logging.getLogger(__name__)
class AMI():
def __init__(self,host,port,username,secret,nodenum):
self._AMIClient = self.AMIClient
self.host = host
self.port = port
self.username = username.encode('utf-8')
self.secret = secret.encode('utf-8')
self.nodenum = str(nodenum)
self.nodenum = str(nodenum).encode('utf-8')
self.CF = None
def send_command(self,command):
self._AMIClient.command = command.encode('utf-8')
self._AMIClient.username = self.username
self._AMIClient.secret = self.secret
self._AMIClient.nodenum = self.nodenum.encode('utf-8')
self.command = command
self.CF = reactor.connectTCP(self.host, self.port, self.AMIClientFactory(self._AMIClient))
factory = self.AMIClientFactory(
self.AMIClient,
self.username,
self.secret,
self.nodenum,
command.encode('utf-8')
)
self.CF = reactor.connectTCP(self.host, self.port, factory)
def closeConnection(self):
self.transport.loseConnection()
if self.CF is not None:
self.CF.disconnect()
class AMIClient(LineReceiver):
delimiter = b'\r\n'
def __init__(self, username, secret, nodenum, command):
self.username = username
self.secret = secret
self.nodenum = nodenum
self.command = command
def connectionMade(self):
self.sendLine(b'Action: login')
@ -40,31 +49,40 @@ class AMI():
self.sendLine(self.delimiter)
def lineReceived(self,line):
print(line)
logger.debug('(AMI) RX: %r', line)
if line == b'Asterisk Call Manager/1.0':
return
if line == b'Response: Success':
self.sendLine(b'Action: command')
#print(b''.join([b'Command: ',b'rpt cmd ',self.nodenum,b' ',self.command]))
self.sendLine(b''.join([b'Command: ',b'rpt cmd ',self.nodenum,b' ',self.command]))
#self.sendLine(b'Command: ' + b'rpt cmd 29177 ilink 3 2001')
self.sendLine(self.delimiter)
self.transport.loseConnection()
class AMIClientFactory(ClientFactory):
def __init__(self,AMIClient):
#self.command = command
self.done = Deferred()
def __init__(self,AMIClient,username,secret,nodenum,command):
self.protocol = AMIClient
#self.protocol.command = command
self.username = username
self.secret = secret
self.nodenum = nodenum
self.command = command
def buildProtocol(self, addr):
return self.protocol(
self.username,
self.secret,
self.nodenum,
self.command
)
def clientConnectionFailed(self, connector, reason):
ClientFactory.clientConnectionLost(self, connector, reason)
logger.warning('(AMI) Connection failed: %s', reason)
ClientFactory.clientConnectionFailed(self, connector, reason)
def clientConnectionLost(self, connector, reason):
logger.debug('(AMI) Connection lost: %s', reason)
ClientFactory.clientConnectionLost(self, connector, reason)

276
API.py

@ -17,132 +17,194 @@
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
from spyne import ServiceBase, rpc, Integer, Decimal, UnsignedInteger32, Unicode, Iterable, error
from dmr_utils3.utils import bytes_3, bytes_4
import json
import logging
from twisted.web.resource import Resource
class FD_APIUserDefinedContext(object):
def __init__(self,CONFIG,BRIDGES):
self.CONFIG = CONFIG
self.BRIDGES = BRIDGES
from dmr_utils3.utils import bytes_4
def getconfig(self):
return self.CONFIG
def getbridges(self):
return self.BRIDGES
def validateKey(self,dmrid,key):
systems = self.CONFIG['SYSTEMS']
dmrid = bytes_4(dmrid)
print(dmrid)
for system in systems:
if systems[system]['MODE'] == 'MASTER':
for peerid in systems[system]['PEERS']:
print(peerid)
if peerid == dmrid:
try:
if key == systems[system]['_opt_key']:
return(system)
else:
return(False)
except KeyError:
return(False)
return(False)
def validateSystemKey(self,systemkey):
if systemkey == self.CONFIG['GLOBAL']['SYSTEM_API_KEY']:
return True
else:
return False
def reset(self,system):
self.CONFIG['SYSTEMS'][system]['_reset'] = True
logger = logging.getLogger(__name__)
MAX_API_BODY = 8192
def options(self,system,options):
self.CONFIG['SYSTEMS'][system]['OPTIONS'] = options
def getoptions(self,system):
return self.CONFIG['SYSTEMS'][system]['OPTIONS']
class APIError(Exception):
def __init__(self, status, message):
self.status = status
self.message = message
super().__init__(message)
def killserver(self):
self.CONFIG['GLOBAL']['_KILL_SERVER'] = True
def resetAllConnections(self):
systems = self.CONFIG['SYSTEMS']
for system in systems:
self.CONFIG['SYSTEMS'][system]['_reset'] = True
class FD_APIController(object):
"""Small, non-blocking control-plane API for live FreeDMR state.
The methods intentionally perform only in-memory dictionary reads/writes.
Expensive snapshots of CONFIG/BRIDGES are not exposed here because this API
shares the Twisted reactor with packet handling.
"""
version = 1
class FD_API(ServiceBase):
_version = 0.1
def __init__(self, CONFIG, BRIDGES):
self.CONFIG = CONFIG
self.BRIDGES = BRIDGES
#return API version
@rpc(Unicode, _returns=Decimal())
def version(ctx, sessionid):
return(FD_API._version)
def validateKey(self, dmrid, key):
try:
peer_id = bytes_4(int(dmrid))
except (TypeError, ValueError, OverflowError):
return False
@rpc()
def dummy(ctx):
pass
for system, system_config in self.CONFIG['SYSTEMS'].items():
if system_config['MODE'] != 'MASTER':
continue
if peer_id not in system_config['PEERS']:
continue
if key == system_config.get('_opt_key'):
return system
return False
return False
######################
#User level API calls#
######################
@rpc(UnsignedInteger32,Unicode)
def reset(ctx,dmrid,key):
system = ctx.udc.validateKey(int(dmrid),key)
if system:
ctx.udc.reset(system)
else:
raise error.InvalidCredentialsError()
def validateSystemKey(self, systemkey):
return systemkey == self.CONFIG['GLOBAL'].get('SYSTEM_API_KEY')
@rpc(UnsignedInteger32,Unicode,Unicode)
def setoptions(ctx,dmrid,key,options):
system = ctx.udc.validateKey(int(dmrid),key)
if system:
ctx.udc.options(system,options)
else:
raise error.InvalidCredentialsError()
def reset(self, system):
self.CONFIG['SYSTEMS'][system]['_reset'] = True
@rpc(UnsignedInteger32,Unicode,_returns=Unicode())
def getoptions(ctx,dmrid,key):
system = ctx.udc.validateKey(int(dmrid),key)
if system:
return ctx.udc.getoptions(system)
else:
raise error.InvalidCredentialsError()
########################
#System level API calls#
########################
@rpc(Unicode)
def killserver(ctx,systemkey):
if ctx.udc.validateSystemKey(systemkey):
return ctx.udc.killserver()
else:
raise error.InvalidCredentialsError()
def options(self, system, options):
self.CONFIG['SYSTEMS'][system]['OPTIONS'] = options
@rpc(Unicode)
def resetall(ctx,systemkey):
if ctx.udc.validateSystemKey(systemkey):
return ctx.udc.resetAllConnections()
def getoptions(self, system):
options = self.CONFIG['SYSTEMS'][system].get('OPTIONS')
has_options = options is not None
if isinstance(options, bytes):
options = options.decode('utf-8', 'ignore')
elif options is None:
options = ''
else:
raise error.InvalidCredentialsError()
options = str(options)
return {
'connected': bool(self.CONFIG['SYSTEMS'][system]['PEERS']),
'has_options': has_options,
'options': options,
}
@rpc(Unicode,_returns=Unicode())
def getconfig(ctx,systemkey):
if ctx.udc.validateSystemKey(systemkey):
return ctx.udc.getconfig()
else:
raise error.InvalidCredentialsError()
def killserver(self):
self.CONFIG['GLOBAL']['_KILL_SERVER'] = True
@rpc(Unicode,_returns=Unicode())
def getbridges(ctx,systemkey):
if ctx.udc.validateSystemKey(systemkey):
return ctx.udc.getbridges()
else:
raise error.InvalidCredentialsError()
def resetAllConnections(self):
for system in self.CONFIG['SYSTEMS']:
self.CONFIG['SYSTEMS'][system]['_reset'] = True
class FD_APIResource(Resource):
isLeaf = True
def __init__(self, controller):
Resource.__init__(self)
self.controller = controller
def render_GET(self, request):
path = self._path(request)
if path == '/api/v1/version':
return self._json(request, 200, {'ok': True, 'version': self.controller.version})
if path == '/api/v1/health':
return self._json(request, 200, {'ok': True})
return self._json(request, 404, {'ok': False, 'error': 'not_found'})
def render_POST(self, request):
try:
payload = self._payload(request)
path = self._path(request)
if path == '/api/v1/reset':
return self._user_call(request, payload, self._reset)
if path == '/api/v1/options/get':
return self._user_call(request, payload, self._get_options)
if path == '/api/v1/options/set':
return self._user_call(request, payload, self._set_options)
if path == '/api/v1/system/kill':
return self._system_call(request, payload, self._kill_server)
if path == '/api/v1/system/resetall':
return self._system_call(request, payload, self._reset_all)
return self._json(request, 404, {'ok': False, 'error': 'not_found'})
except APIError as exc:
return self._json(request, exc.status, {'ok': False, 'error': exc.message})
except Exception:
logger.exception('(API) Unhandled API error')
return self._json(request, 500, {'ok': False, 'error': 'internal_error'})
def _path(self, request):
return (b'/' + b'/'.join(request.postpath)).decode('utf-8', 'ignore')
def _payload(self, request):
content_length = None
if hasattr(request, 'getHeader'):
content_length = request.getHeader('content-length')
if content_length is not None:
try:
if int(content_length) > MAX_API_BODY:
raise APIError(413, 'request_too_large')
except ValueError:
raise APIError(400, 'invalid_content_length')
body = request.content.read()
if len(body) > MAX_API_BODY:
raise APIError(413, 'request_too_large')
if not body:
return {}
try:
payload = json.loads(body.decode('utf-8'))
except (TypeError, ValueError):
raise APIError(400, 'invalid_json')
if not isinstance(payload, dict):
raise APIError(400, 'invalid_json')
return payload
def _user_call(self, request, payload, handler):
dmrid = payload.get('dmrid')
key = payload.get('key')
system = self.controller.validateKey(dmrid, key)
if not system:
raise APIError(401, 'invalid_credentials')
result = handler(system, payload)
return self._json(request, 200, {'ok': True, **result})
def _system_call(self, request, payload, handler):
if not self.controller.validateSystemKey(payload.get('systemkey')):
raise APIError(401, 'invalid_credentials')
result = handler(payload)
return self._json(request, 200, {'ok': True, **result})
def _reset(self, system, payload):
self.controller.reset(system)
return {'reset': True, 'system': system}
def _get_options(self, system, payload):
return self.controller.getoptions(system)
def _set_options(self, system, payload):
options = payload.get('options')
if not isinstance(options, str):
raise APIError(400, 'missing_options')
self.controller.options(system, options)
return {'updated': True, 'system': system}
def _kill_server(self, payload):
self.controller.killserver()
return {'killserver': True}
def _reset_all(self, payload):
self.controller.resetAllConnections()
return {'resetall': True}
def _json(self, request, status, payload):
request.setResponseCode(status)
request.setHeader(b'content-type', b'application/json')
return json.dumps(payload, separators=(',', ':')).encode('utf-8')
def make_api_resource(CONFIG, BRIDGES):
return FD_APIResource(FD_APIController(CONFIG, BRIDGES))

@ -3,3 +3,5 @@
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).

@ -1,23 +1,26 @@
from twisted.internet import reactor
from twisted.web.xmlrpc import Proxy
def printValue(value):
print(repr(value))
reactor.stop()
def printError(error):
print("error", error)
reactor.stop()
def capitalize(value):
print(value)
proxy = Proxy(b"http://localhost:7080/xmlrpc")
# The callRemote method accepts a method name and an argument list.
proxy.callRemote("FD_API.reset", '2', '55555').addCallbacks(capitalize, printError)
reactor.run()
import json
import sys
from urllib import request
def post(path, payload):
body = json.dumps(payload).encode("utf-8")
req = request.Request(
"http://127.0.0.1:8000" + path,
data=body,
headers={"content-type": "application/json"},
method="POST",
)
with request.urlopen(req, timeout=3) as response:
return json.loads(response.read().decode("utf-8"))
if __name__ == "__main__":
if len(sys.argv) != 4:
print("usage: api_client.py <dmrid> <key> <options>")
raise SystemExit(2)
print(post("/api/v1/options/set", {
"dmrid": int(sys.argv[1]),
"key": sys.argv[2],
"options": sys.argv[3],
}))

@ -67,6 +67,12 @@ __email__ = 'n0mjs@me.com'
# Module gobal varaibles
def dmrd_seq_delta(seq, last_seq):
if last_seq is False or last_seq is None:
return None
return (seq - last_seq) % 256
# Timed loop used for reporting HBP status
#
# REPORT BASED ON THE TYPE SELECTED IN THE MAIN CONFIG FILE
@ -323,16 +329,16 @@ class routerOBP(OPENBRIDGE):
if self.STATUS[_stream_id]['lastData'] and self.STATUS[_stream_id]['lastData'] == _data and _seq > 1:
logger.warning("(%s) *PacketControl* last packet is a complete duplicate of the previous one, disgarding. Stream ID:, %s TGID: %s",self._system,int_id(_stream_id),int_id(_dst_id))
return
#Handle inbound duplicates
if _seq and _seq == self.STATUS[_stream_id]['lastSeq']:
_seq_delta = dmrd_seq_delta(_seq, self.STATUS[_stream_id]['lastSeq'])
if _seq_delta == 0:
logger.warning("(%s) *PacketControl* Duplicate sequence number %s, disgarding. Stream ID:, %s TGID: %s",self._system,_seq,int_id(_stream_id),int_id(_dst_id))
return
#Inbound out-of-order packets
if _seq and self.STATUS[_stream_id]['lastSeq'] and (_seq != 1) and (_seq < self.STATUS[_stream_id]['lastSeq']):
if _seq_delta is not None and _seq_delta > 127:
logger.warning("%s) *PacketControl* Out of order packet - last SEQ: %s, this SEQ: %s, disgarding. Stream ID:, %s TGID: %s ",self._system,self.STATUS[_stream_id]['lastSeq'],_seq,int_id(_stream_id),int_id(_dst_id))
return
#Inbound missed packets
if _seq and self.STATUS[_stream_id]['lastSeq'] and _seq > (self.STATUS[_stream_id]['lastSeq']+1):
if _seq_delta is not None and _seq_delta > 1:
logger.warning("(%s) *PacketControl* Missed packet(s) - last SEQ: %s, this SEQ: %s. Stream ID:, %s TGID: %s ",self._system,self.STATUS[_stream_id]['lastSeq'],_seq,int_id(_stream_id),int_id(_dst_id))
#Save this sequence number
@ -495,14 +501,15 @@ 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)
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.STATUS[_stream_id]['_fin'] = True
self.STATUS[_stream_id]['_fin'] = True
#removed = self.STATUS.pop(_stream_id)
#logger.debug('(%s) OpenBridge sourced call stream end, remove terminated Stream ID: %s', self._system, int_id(_stream_id))
#if not removed:
#selflogger.error('(%s) *CALL END* STREAM ID: %s NOT IN LIST -- THIS IS A REAL PROBLEM', self._system, int_id(_stream_id))
#Reset sequence number
self._lastSeq = False
#Reset sequence tracking
self.STATUS[_stream_id]['lastSeq'] = False
self.STATUS[_stream_id]['lastData'] = False
class routerHBP(HBSYSTEM):
@ -539,7 +546,9 @@ class routerHBP(HBSYSTEM):
4: b'\x00',
},
'lastSeq': False,
'lastData': False
'lastData': False,
'RX_FINISHED_STREAM_ID': b'\x00',
'RX_FINISHED_STREAM_LOG': False
},
2: {
@ -568,7 +577,9 @@ class routerHBP(HBSYSTEM):
4: b'\x00',
},
'lastSeq': False,
'lastData': False
'lastData': False,
'RX_FINISHED_STREAM_ID': b'\x00',
'RX_FINISHED_STREAM_LOG': False
}
}
@ -584,9 +595,19 @@ class routerHBP(HBSYSTEM):
_source_rptr = _peer_id
if _call_type == 'group':
if self.STATUS[_slot].get('RX_FINISHED_STREAM_ID') == _stream_id:
if not self.STATUS[_slot].get('RX_FINISHED_STREAM_LOG'):
logger.warning("(%s) HBP *LoopControl* STREAM ID: %s ALREADY FINISHED FROM THIS SOURCE, IGNORING",self._system, int_id(_stream_id))
self.STATUS[_slot]['RX_FINISHED_STREAM_LOG'] = True
return
# Is this a new call stream?
if (_stream_id != self.STATUS[_slot]['RX_STREAM_ID']):
self.STATUS[_slot]['RX_FINISHED_STREAM_ID'] = b'\x00'
self.STATUS[_slot]['RX_FINISHED_STREAM_LOG'] = False
self.STATUS[_slot]['lastSeq'] = False
self.STATUS[_slot]['lastData'] = False
if (self.STATUS[_slot]['RX_TYPE'] != HBPF_SLT_VTERM) and (pkt_time < (self.STATUS[_slot]['RX_TIME'] + STREAM_TO)) and (_rf_src != self.STATUS[_slot]['RX_RFS']):
logger.warning('(%s) Packet received with STREAM ID: %s <FROM> SUB: %s PEER: %s <TO> TGID %s, SLOT %s collided with existing call', self._system, int_id(_stream_id), int_id(_rf_src), int_id(_peer_id), int_id(_dst_id), _slot)
return
@ -639,16 +660,16 @@ class routerHBP(HBSYSTEM):
if self.STATUS[_slot]['lastData'] and self.STATUS[_slot]['lastData'] == _data and _seq > 1:
logger.warning("(%s) *PacketControl* last packet is a complete duplicate of the previous one, disgarding. Stream ID:, %s TGID: %s",self._system,int_id(_stream_id),int_id(_dst_id))
return
#Handle inbound duplicates
if _seq and _seq == self.STATUS[_slot]['lastSeq']:
_seq_delta = dmrd_seq_delta(_seq, self.STATUS[_slot]['lastSeq'])
if _seq_delta == 0:
logger.warning("(%s) *PacketControl* Duplicate sequence number %s, disgarding. Stream ID:, %s TGID: %s",self._system,_seq,int_id(_stream_id),int_id(_dst_id))
return
#Inbound out-of-order packets
if _seq and self.STATUS[_slot]['lastSeq'] and (_seq != 1) and (_seq < self.STATUS[_slot]['lastSeq']):
if _seq_delta is not None and _seq_delta > 127:
logger.warning("%s) *PacketControl* Out of order packet - last SEQ: %s, this SEQ: %s, disgarding. Stream ID:, %s TGID: %s ",self._system,self.STATUS[_slot]['lastSeq'],_seq,int_id(_stream_id),int_id(_dst_id))
return
#Inbound missed packets
if _seq and self.STATUS[_slot]['lastSeq'] and _seq > (self.STATUS[_slot]['lastSeq']+1):
if _seq_delta is not None and _seq_delta > 1:
logger.warning("(%s) *PacketControl* Missed packet(s) - last SEQ: %s, this SEQ: %s. Stream ID:, %s TGID: %s ",self._system,self.STATUS[_slot]['lastSeq'],_seq,int_id(_stream_id),int_id(_dst_id))
#Save this sequence number
@ -812,6 +833,10 @@ class routerHBP(HBSYSTEM):
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)
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.STATUS[_slot]['RX_FINISHED_STREAM_ID'] = _stream_id
self.STATUS[_slot]['RX_FINISHED_STREAM_LOG'] = False
self.STATUS[_slot]['lastSeq'] = False
self.STATUS[_slot]['lastData'] = False
#
# Begin in-band signalling for call end. This has nothign to do with routing traffic directly.

File diff suppressed because it is too large Load Diff

@ -137,7 +137,7 @@ def build_config(_config_file):
'PATH': config.get(section, 'PATH',fallback='./'),
'PING_TIME': config.getint(section, 'PING_TIME', fallback=10),
'MAX_MISSED': config.getint(section, 'MAX_MISSED', fallback=3),
'USE_ACL': config.get(section, 'USE_ACL', fallback=True),
'USE_ACL': config.getboolean(section, 'USE_ACL', fallback=True),
'REG_ACL': config.get(section, 'REG_ACL', fallback='PERMIT:ALL'),
'SUB_ACL': config.get(section, 'SUB_ACL', fallback='DENY:1'),
'TG1_ACL': config.get(section, 'TGID_TS1_ACL', fallback='PERMIT:ALL'),
@ -149,7 +149,8 @@ def build_config(_config_file):
'DATA_GATEWAY': config.getboolean(section, 'DATA_GATEWAY', fallback=False),
'VALIDATE_SERVER_IDS': config.getboolean(section, 'VALIDATE_SERVER_IDS', fallback=True),
'DEBUG_BRIDGES' : config.getboolean(section, 'DEBUG_BRIDGES', fallback=False),
'ENABLE_API' : config.getboolean(section, 'ENABLE_API', fallback=False)
'ENABLE_API' : config.getboolean(section, 'ENABLE_API', fallback=False),
'_KILL_SERVER': False
})

@ -0,0 +1,113 @@
# FreeDMR API
FreeDMR includes an experimental HTTP/JSON API for small live control-plane
actions. It is intended for local administration and automation, not for public
internet exposure.
Enable it with:
```ini
[GLOBAL]
ENABLE_API: True
```
When enabled, the API listens on TCP port `8000`.
## Safety Notes
FreeDMR is a live voice routing process. API requests are deliberately limited
to small in-memory operations so they do not delay DMR voice packet handling.
Request bodies larger than 8192 bytes are rejected.
Bind or firewall port `8000` appropriately. Do not expose it publicly without a
trusted reverse proxy and access controls.
## Authentication
User-level endpoints require:
- `dmrid`: the connected HBP peer/repeater DMR ID
- `key`: the session options key for that peer
System-level endpoints require:
- `systemkey`: the FreeDMR system API key
## Endpoints
### Health
```bash
curl http://127.0.0.1:8000/api/v1/health
```
### Version
```bash
curl http://127.0.0.1:8000/api/v1/version
```
### Get Options
```bash
curl -X POST http://127.0.0.1:8000/api/v1/options/get \
-H 'content-type: application/json' \
-d '{"dmrid":1234567,"key":"secret"}'
```
If no live options are present, the response is:
```json
{"ok":true,"connected":true,"has_options":false,"options":""}
```
### Set Options
```bash
curl -X POST http://127.0.0.1:8000/api/v1/options/set \
-H 'content-type: application/json' \
-d '{"dmrid":1234567,"key":"secret","options":"KEY=secret;TS1=91;DIAL=2350"}'
```
The `options` value must be the complete FreeDMR `OPTIONS` string. The API does
not add or preserve `KEY=...` automatically.
### Reset Peer Session
```bash
curl -X POST http://127.0.0.1:8000/api/v1/reset \
-H 'content-type: application/json' \
-d '{"dmrid":1234567,"key":"secret"}'
```
FreeDMR expects one HBP peer per master instance, so this resets the master
instance that owns the authenticated peer session.
### Reset All Connections
```bash
curl -X POST http://127.0.0.1:8000/api/v1/system/resetall \
-H 'content-type: application/json' \
-d '{"systemkey":"system-secret"}'
```
### Stop FreeDMR
```bash
curl -X POST http://127.0.0.1:8000/api/v1/system/kill \
-H 'content-type: application/json' \
-d '{"systemkey":"system-secret"}'
```
## Responses
Successful responses include `"ok": true`. Failed responses include
`"ok": false` and an `error` string.
Common errors:
- `invalid_credentials`
- `invalid_json`
- `missing_options`
- `request_too_large`
- `not_found`

@ -377,8 +377,16 @@ class OPENBRIDGE(DatagramProtocol):
return
elif _packet[:4] == DMRE:
if len(_packet) < 56:
h,p = _sockaddr
logger.warning('(%s) FreeBridge packet too short, discarded - OPCODE: %s LENGTH: %s SRC IP: %s SRC PORT: %s', self._system, _packet[:4], len(_packet), h, p)
return
if _packet[55] > 4:
if len(_packet) < 89:
h,p = _sockaddr
logger.warning('(%s) FreeBridge v%s packet too short, discarded - OPCODE: %s LENGTH: %s SRC IP: %s SRC PORT: %s', self._system, _packet[55], _packet[:4], len(_packet), h, p)
return
_data = _packet[:53]
_ber = _packet[53:54]
_rssi = _packet[54:55]
@ -393,6 +401,10 @@ class OPENBRIDGE(DatagramProtocol):
_h = blake2b(key=self._config['PASSPHRASE'], digest_size=16)
_h.update(_packet[:73])
else:
if len(_packet) < 85:
h,p = _sockaddr
logger.warning('(%s) FreeBridge v%s packet too short, discarded - OPCODE: %s LENGTH: %s SRC IP: %s SRC PORT: %s', self._system, _packet[55], _packet[:4], len(_packet), h, p)
return
_data = _packet[:53]
_ber = _packet[53:54]
_rssi = _packet[54:55]
@ -595,10 +607,10 @@ class OPENBRIDGE(DatagramProtocol):
if _packet[:4] == BCST:
#_data = _packet[:11]
_hash = _packet[4:]
_ckhs = hmac_new(self._config['PASSPHRASE'],_packet[4:],sha1).digest()
_ckhs = hmac_new(self._config['PASSPHRASE'],_packet[:4],sha1).digest()
if compare_digest(_hash, _ckhs):
logger.trace('(%s) *BridgeControl* BCST STUN request received for TGID: %s, Stream ID: %s',self._system,int_id(_tgid), int_id(_stream_id))
self._config['_STUN'] = True
logger.trace('(%s) *BridgeControl* BCST STUN request received',self._system)
self._CONFIG['STUN'] = True
else:
h,p = _sockaddr
logger.warning('(%s) *BridgeControl* BCST invalid STUN, packet discarded - OPCODE: %s DATA: %s HMAC LENGTH: %s HMAC: %s SRC IP: %s SRC PORT: %s', self._system, _packet[:4], repr(_packet[:53]), len(_packet[53:]), repr(_packet[53:]),h,p)
@ -843,6 +855,9 @@ class HBSYSTEM(DatagramProtocol):
# Extract the command, which is various length, all but one 4 significant characters -- RPTCL
_command = _data[:4]
if _command == DMRD: # DMRData -- encapsulated DMR data frame
if len(_data) < 53:
logger.warning('(%s) DMRD packet too short, discarded - LENGTH: %s SRC IP: %s SRC PORT: %s', self._system, len(_data), _sockaddr[0], _sockaddr[1])
return
_peer_id = _data[11:15]
if _peer_id in self._peers \
and self._peers[_peer_id]['CONNECTION'] == 'YES' \
@ -1094,6 +1109,9 @@ class HBSYSTEM(DatagramProtocol):
_command = _data[:4]
if _command == DMRD: # DMRData -- encapsulated DMR data frame
if len(_data) < 53:
logger.warning('(%s) DMRD packet too short, discarded - LENGTH: %s SRC IP: %s SRC PORT: %s', self._system, len(_data), _sockaddr[0], _sockaddr[1])
return
_peer_id = _data[11:15]
if self._config['LOOSE'] or _peer_id == self._config['RADIO_ID']: # Validate the Radio_ID unless using loose validation
#_seq = _data[4:5]

@ -35,6 +35,9 @@ __license__ = 'GNU GPLv3'
__maintainer__ = 'Simon Adlem G7RZU'
__email__ = 'simon@gb7fr.org.uk'
def bool_from_env(value):
return str(value).strip().lower() in ('1', 'true', 'yes', 'on')
def IsIPv4Address(ip):
try:
ipaddress.IPv4Address(ip)
@ -362,26 +365,23 @@ if __name__ == '__main__':
print('(PROXY)(GLOBAL) SHUTDOWN: PROXY IS TERMINATING WITH SIGNAL {}'.format(str(_signal)))
reactor.stop()
def sigt(_signal,_frame):
print('oooh')
#Install signal handlers
signal.signal(signal.SIGINT, sig_handler)
signal.signal(signal.SIGTERM, sigt)
signal.signal(signal.SIGTERM, sig_handler)
#readState()
#If IPv6 is enabled by enivornment variable...
if ListenIP == '' and 'FDPROXY_IPV6' in os.environ and bool(os.environ['FDPROXY_IPV6']):
if ListenIP == '' and bool_from_env(os.environ.get('FDPROXY_IPV6')):
ListenIP = '::'
#Override static config from Environment
if 'FDPROXY_STATS' in os.environ:
Stats = bool(os.environ['FDPROXY_STATS'])
Stats = bool_from_env(os.environ['FDPROXY_STATS'])
#if 'FDPROXY_DEBUG' in os.environ:
# Debug = bool(os.environ['FDPROXY_DEBUG'])
# Debug = bool_from_env(os.environ['FDPROXY_DEBUG'])
if 'FDPROXY_CLIENTINFO' in os.environ:
ClientInfo = bool(os.environ['FDPROXY_CLIENTINFO'])
ClientInfo = bool_from_env(os.environ['FDPROXY_CLIENTINFO'])
if 'FDPROXY_LISTENPORT' in os.environ:
ListenPort = int(os.environ['FDPROXY_LISTENPORT'])

@ -1,3 +0,0 @@
home = /usr/bin
include-system-site-packages = false
version = 3.10.12

@ -29,6 +29,9 @@ from reporting_const import *
from pprint import pprint
def bool_flag(value):
return str(value).strip().lower() in ('1', 'true', 'yes', 'on')
class reportClient(NetstringReceiver):
def stringReceived(self, data):
@ -36,10 +39,10 @@ class reportClient(NetstringReceiver):
if data[:1] == REPORT_OPCODES['BRDG_EVENT']:
self.bridgeEvent(data[1:].decode('UTF-8'))
elif data[:1] == REPORT_OPCODES['CONFIG_SND']:
if cli_args.CONFIG:
if bool_flag(cli_args.CONFIG):
self.configSend(data[1:])
elif data[:1] == REPORT_OPCODES['BRIDGE_SND']:
if cli_args.BRIDGES:
if bool_flag(cli_args.BRIDGES):
self.bridgeSend(data[1:])
elif data == b'bridge updated':
pass
@ -64,12 +67,12 @@ class reportClient(NetstringReceiver):
if len(datalist) > 9:
event['duration'] = datalist[9]
if cli_args.EVENTS:
if bool_flag(cli_args.EVENTS):
pprint(event, compact=True)
def bridgeSend(self,data):
self.BRIDGES = pickle.loads(data)
if cli_args.STATS:
if bool_flag(cli_args.STATS):
print('There are currently {} active bridges in the bridge table:\n'.format(len(self.BRIDGES)))
for _bridge in self.BRIDGES.keys():
print('{},'.format({str(_bridge)}))

@ -28,6 +28,7 @@ __email__ = 'simon@gb7fr.org.uk'
#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 ReconnectingClientFactory
@ -42,6 +43,7 @@ class reportClient(NetstringReceiver):
def __init__(self,db,reactor):
self.db = db
self.reactor = reactor
self._db_lock = threading.Lock()
def stringReceived(self, data):
@ -75,25 +77,39 @@ class reportClient(NetstringReceiver):
event['duration'] = datalist[9]
#self.reactor.callInThread(self.send_mysql,event)
self.send_mysql(event)
self.reactor.callInThread(self.send_mysql,event)
def send_mysql(self,event):
with self._db_lock:
while not self.db.is_connected():
try:
self.db.reconnect()
except mysql.connector.Error as err:
print('(MYSQL) error on reconnect: {}'.format(err))
while not self.db.is_connected():
print("{} {} {} {} {} {} {} {} {}".format(event['type'],event['event'], event['trx'],event['system'],event['streamid'],event['peerid'],event['subid'],event['slot'],event['dstid'],event['duration']))
_cursor = self.db.cursor()
try:
self.db.reconnect()
_cursor.execute(
"insert into feed values (NULL,%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'],
)
)
self.db.commit()
except mysql.connector.Error as err:
print('(MYSQL) error on reconnect: {}'.format(err))
print("{} {} {} {} {} {} {} {} {}".format(event['type'],event['event'], event['trx'],event['system'],event['streamid'],event['peerid'],event['subid'],event['slot'],event['dstid'],event['duration']))
_cursor = self.db.cursor()
try:
_cursor.execute("insert into feed values (NULL,'{}','{}','{}','{}','{}','{}','{}','{}','{}','{}')".format(event['type'],event['event'], event['trx'],event['system'],event['streamid'],event['peerid'],event['subid'],event['slot'],event['dstid'],event['duration']))
self.db.commit()
except mysql.connector.Error as err:
_cursor.close()
print('(MYSQL) error, problem with cursor execute: {}'.format(err))
print('(MYSQL) error, problem with cursor execute: {}'.format(err))
finally:
_cursor.close()
def bridgeSend(self,data):
@ -116,7 +132,7 @@ class reportClientFactory(ReconnectingClientFactory):
print('Connected.')
print('Resetting reconnection delay')
self.resetDelay()
return self.proto(db,reactor)
return self.proto(self.db,self.reactor)
def clientConnectionLost(self, connector, reason):
print('Lost connection. Reason:', reason)

@ -6,4 +6,3 @@ configparser>=3.0.0
resettabletimer>=0.7.0
setproctitle
Pyro5
spyne

@ -0,0 +1 @@
"""FreeDMR test package."""

@ -0,0 +1 @@
"""Test harness helpers for FreeDMR."""

@ -0,0 +1,491 @@
"""In-process deterministic packet harness for bridge_master tests.
This module is test-only. It avoids UDP sockets and replaces production
network sends with capture functions while leaving production modules unchanged.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from types import SimpleNamespace
import copy
import importlib
import unittest
DMRD = b"DMRD"
HBPF_VOICE = 0x0
HBPF_VOICE_SYNC = 0x1
HBPF_DATA_SYNC = 0x2
HBPF_SLT_VHEAD = 0x1
HBPF_SLT_VTERM = 0x2
ID_MAX = 16776415
PEER_MAX = 4294967295
def require_bridge_master():
"""Import bridge_master or skip tests when runtime deps are unavailable."""
try:
return importlib.import_module("bridge_master")
except ModuleNotFoundError as exc:
raise unittest.SkipTest(
f"bridge_master runtime dependency is not installed: {exc.name}"
) from exc
def bytes_3(value: int | bytes) -> bytes:
if isinstance(value, bytes):
if len(value) != 3:
raise ValueError("expected exactly 3 bytes")
return value
return int(value).to_bytes(3, "big")
def bytes_4(value: int | bytes) -> bytes:
if isinstance(value, bytes):
if len(value) != 4:
raise ValueError("expected exactly 4 bytes")
return value
return int(value).to_bytes(4, "big")
def int_id(value: int | bytes) -> int:
if isinstance(value, int):
return value
return int.from_bytes(value, "big")
def acl_permit_all(max_id: int = ID_MAX) -> tuple[bool, list[tuple[int, int]]]:
return True, [(1, max_id)]
def hbp_bits(slot: int, call_type: str, frame_type: int, dtype_vseq: int) -> int:
bits = ((frame_type & 0x3) << 4) | (dtype_vseq & 0xF)
if slot == 2:
bits |= 0x80
if call_type == "unit":
bits |= 0x40
return bits
def parse_dmr_fields(packet: bytes) -> dict[str, object]:
if len(packet) < 20 or packet[:4] != DMRD:
return {"raw": packet}
bits = packet[15]
if bits & 0x40:
call_type = "unit"
elif (bits & 0x23) == 0x23:
call_type = "vcsbk"
else:
call_type = "group"
return {
"opcode": packet[:4],
"seq": packet[4],
"rf_src": packet[5:8],
"dst_id": packet[8:11],
"peer_id": packet[11:15],
"bits": bits,
"slot": 2 if bits & 0x80 else 1,
"call_type": call_type,
"frame_type": (bits & 0x30) >> 4,
"dtype_vseq": bits & 0xF,
"stream_id": packet[16:20],
"dmr_payload": packet[20:53],
"ber": packet[53:54],
"rssi": packet[54:55],
}
@dataclass(frozen=True)
class PacketSpec:
peer_id: int | bytes = 1001
rf_src: int | bytes = 3120001
dst_id: int | bytes = 91
slot: int = 2
stream_id: int | bytes = 0x01020304
seq: int = 0
call_type: str = "group"
frame_type: int = HBPF_VOICE
dtype_vseq: int = 0
payload: bytes = b"\x00" * 33
ber: bytes = b"\x00"
rssi: bytes = b"\x00"
delay: float = 0.0
def data(self) -> bytes:
if len(self.payload) != 33:
raise ValueError("DMR payload must be exactly 33 bytes")
return b"".join(
[
DMRD,
bytes([self.seq & 0xFF]),
bytes_3(self.rf_src),
bytes_3(self.dst_id),
bytes_4(self.peer_id),
bytes([hbp_bits(self.slot, self.call_type, self.frame_type, self.dtype_vseq)]),
bytes_4(self.stream_id),
self.payload,
self.ber,
self.rssi,
]
)
def decoded_args(self) -> tuple[bytes, bytes, bytes, int, int, str, int, int, bytes, bytes]:
return (
bytes_4(self.peer_id),
bytes_3(self.rf_src),
bytes_3(self.dst_id),
self.seq & 0xFF,
self.slot,
self.call_type,
self.frame_type,
self.dtype_vseq,
bytes_4(self.stream_id),
self.data(),
)
def decoded_obp_args(
self,
packet_hash: bytes = b"",
hops: bytes = b"",
source_server: int | bytes = 9990,
source_rptr: int | bytes = 0,
) -> tuple[bytes, bytes, bytes, int, int, str, int, int, bytes, bytes, bytes, bytes, bytes, bytes, bytes, bytes]:
return (
bytes_4(self.peer_id),
bytes_3(self.rf_src),
bytes_3(self.dst_id),
self.seq & 0xFF,
self.slot,
self.call_type,
self.frame_type,
self.dtype_vseq,
bytes_4(self.stream_id),
self.data(),
packet_hash,
hops,
bytes_4(source_server),
self.ber,
self.rssi,
bytes_4(source_rptr),
)
@dataclass
class CapturedPacket:
target_system: str
packet: bytes
hops: bytes | None = None
ber: bytes = b"\x00"
rssi: bytes = b"\x00"
source_server: bytes = b"\x00\x00\x00\x00"
source_rptr: bytes = b"\x00\x00\x00\x00"
fields: dict[str, object] = field(init=False)
def __post_init__(self) -> None:
self.fields = parse_dmr_fields(self.packet)
class PacketCapture:
def __init__(self) -> None:
self.packets: list[CapturedPacket] = []
def recorder(self, target_system: str):
def record(
packet: bytes,
hops: bytes | None = b"",
ber: bytes = b"\x00",
rssi: bytes = b"\x00",
source_server: bytes = b"\x00\x00\x00\x00",
source_rptr: bytes = b"\x00\x00\x00\x00",
) -> None:
self.packets.append(
CapturedPacket(
target_system=target_system,
packet=packet,
hops=hops,
ber=ber,
rssi=rssi,
source_server=source_server,
source_rptr=source_rptr,
)
)
return record
def for_system(self, system: str) -> list[CapturedPacket]:
return [packet for packet in self.packets if packet.target_system == system]
class ReportCapture:
def __init__(self) -> None:
self.events: list[bytes] = []
def send_bridgeEvent(self, data: bytes) -> None:
self.events.append(data)
class FakeClock:
def __init__(self, start: float = 1_700_000_000.0) -> None:
self.now = float(start)
def time(self) -> float:
return self.now
def advance(self, seconds: float) -> float:
self.now += seconds
return self.now
class FakeReactor:
def __init__(self) -> None:
self.later: list[tuple[float, object, tuple, dict]] = []
self.thread_calls: list[tuple[object, tuple, dict]] = []
def callLater(self, delay, func, *args, **kwargs):
self.later.append((delay, func, args, kwargs))
return SimpleNamespace(cancel=lambda: None, active=lambda: True)
def callInThread(self, func, *args, **kwargs):
self.thread_calls.append((func, args, kwargs))
def callFromThread(self, func, *args, **kwargs):
return func(*args, **kwargs)
class FakeTransport:
def __init__(self) -> None:
self.writes: list[tuple[bytes, tuple[str, int] | None]] = []
def write(self, packet: bytes, sockaddr=None) -> None:
self.writes.append((packet, sockaddr))
def minimal_config(system_names: tuple[str, ...] = ("MASTER-A", "MASTER-B")) -> dict:
config = {
"GLOBAL": {
"SERVER_ID": bytes_4(9990),
"USE_ACL": False,
"TG1_ACL": acl_permit_all(),
"TG2_ACL": acl_permit_all(),
"SUB_ACL": acl_permit_all(),
"GEN_STAT_BRIDGES": False,
"DATA_GATEWAY": False,
"VALIDATE_SERVER_IDS": False,
},
"REPORTS": {"REPORT": False},
"ALIASES": {"PATH": "./", "SUB_MAP_FILE": ""},
"ALLSTAR": {"ENABLED": False},
"SYSTEMS": {},
"_SUB_IDS": {},
"_PEER_IDS": {},
"_LOCAL_SUBSCRIBER_IDS": {},
"_SERVER_IDS": {},
"CHECKSUMS": {},
}
for name in system_names:
config["SYSTEMS"][name] = {
"MODE": "MASTER",
"ENABLED": True,
"REPEAT": True,
"MAX_PEERS": 1,
"IP": "127.0.0.1",
"PORT": 0,
"PASSPHRASE": b"",
"GROUP_HANGTIME": 0,
"USE_ACL": False,
"REG_ACL": acl_permit_all(PEER_MAX),
"SUB_ACL": acl_permit_all(),
"TG1_ACL": acl_permit_all(),
"TG2_ACL": acl_permit_all(),
"DEFAULT_UA_TIMER": 1,
"SINGLE_MODE": True,
"VOICE_IDENT": False,
"TS1_STATIC": "",
"TS2_STATIC": "",
"DEFAULT_REFLECTOR": 0,
"GENERATOR": 0,
"ANNOUNCEMENT_LANGUAGE": "en_GB",
"ALLOW_UNREG_ID": True,
"PROXY_CONTROL": False,
"OVERRIDE_IDENT_TG": False,
"PEERS": {},
}
return config
def add_openbridge_system(config: dict, name: str = "OBP-1", network_id: int = 1) -> dict:
config["SYSTEMS"][name] = {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"NETWORK_ID": bytes_4(network_id),
"IP": "127.0.0.1",
"PORT": 0,
"PASSPHRASE": b"test-passphrase\x00\x00\x00\x00\x00\x00",
"TARGET_IP": "127.0.0.1",
"TARGET_PORT": 0,
"TARGET_SOCK": ("127.0.0.1", 0),
"USE_ACL": False,
"SUB_ACL": acl_permit_all(),
"TG1_ACL": acl_permit_all(),
"TG2_ACL": acl_permit_all(),
"RELAX_CHECKS": True,
"ENHANCED_OBP": False,
"VER": 5,
}
return config
def active_bridge(
name: str,
tg_id: int,
entries: tuple[tuple[str, int], ...],
timeout_minutes: int = 1,
) -> dict[str, list[dict]]:
tg_bytes = bytes_3(tg_id)
return {
name: [
{
"SYSTEM": system,
"TS": slot,
"TGID": tg_bytes,
"ACTIVE": True,
"TIMEOUT": timeout_minutes * 60,
"TO_TYPE": "ON",
"OFF": [],
"ON": [tg_bytes],
"RESET": [],
"TIMER": 0,
}
for system, slot in entries
]
}
class DeterministicScenario:
def __init__(self, config: dict | None = None, bridges: dict | None = None) -> None:
self.config = config or minimal_config()
self.bridges = bridges or {}
self.clock = FakeClock()
self.capture = PacketCapture()
self.reports: dict[str, ReportCapture] = {}
self.transports: dict[str, FakeTransport] = {}
self.reactor = FakeReactor()
self.bm = None
self._saved_attrs: dict[str, object] = {}
self._saved_systems: dict | None = None
def __enter__(self):
self.bm = require_bridge_master()
self._saved_systems = dict(self.bm.systems)
for attr in (
"CONFIG",
"BRIDGES",
"SUB_MAP",
"peer_ids",
"subscriber_ids",
"talkgroup_ids",
"local_subscriber_ids",
"server_ids",
"checksums",
"reactor",
"time",
"words",
):
if hasattr(self.bm, attr):
self._saved_attrs[attr] = getattr(self.bm, attr)
self.bm.CONFIG = self.config
self.bm.BRIDGES = copy.deepcopy(self.bridges)
self.bm.SUB_MAP = {}
self.bm.peer_ids = {}
self.bm.subscriber_ids = {}
self.bm.talkgroup_ids = {}
self.bm.local_subscriber_ids = {}
self.bm.server_ids = {}
self.bm.checksums = {}
self.bm.words = {"en_GB": {"silence": b"", "busy": b"", "notlinked": b"", "linkedto": b"", "to": b""}}
self.bm.reactor = self.reactor
self.bm.time = self.clock.time
self.bm.systems.clear()
for system_name, system_config in self.config["SYSTEMS"].items():
report = ReportCapture()
self.reports[system_name] = report
if system_config["MODE"] == "MASTER":
system = self.bm.routerHBP(system_name, self.config, report)
elif system_config["MODE"] == "OPENBRIDGE":
system = self.bm.routerOBP(system_name, self.config, report)
else:
continue
system.send_system = self.capture.recorder(system_name)
transport = FakeTransport()
system.transport = transport
self.transports[system_name] = transport
self.bm.systems[system_name] = system
return self
def __exit__(self, exc_type, exc, tb) -> None:
if self.bm is None:
return
self.bm.systems.clear()
if self._saved_systems is not None:
self.bm.systems.update(self._saved_systems)
for attr in (
"CONFIG",
"BRIDGES",
"SUB_MAP",
"peer_ids",
"subscriber_ids",
"talkgroup_ids",
"local_subscriber_ids",
"server_ids",
"checksums",
"reactor",
"time",
"words",
):
if attr in self._saved_attrs:
setattr(self.bm, attr, self._saved_attrs[attr])
elif hasattr(self.bm, attr):
delattr(self.bm, attr)
@property
def systems(self):
return self.bm.systems
@property
def bridge_state(self):
return self.bm.BRIDGES
def inject_hbp(self, system_name: str, packet: PacketSpec) -> None:
self.systems[system_name].dmrd_received(*packet.decoded_args())
def inject_obp(self, system_name: str, packet: PacketSpec) -> None:
self.systems[system_name].dmrd_received(*packet.decoded_obp_args())
def inject_datagram(self, system_name: str, packet: bytes, sockaddr=("127.0.0.1", 50000)) -> None:
self.systems[system_name].datagramReceived(packet, sockaddr)
def register_peer(
self,
system_name: str,
peer_id: int | bytes = 1001,
sockaddr=("127.0.0.1", 50000),
callsign: bytes = b"TEST ",
) -> bytes:
peer = bytes_4(peer_id)
self.config["SYSTEMS"][system_name]["PEERS"][peer] = {
"CONNECTION": "YES",
"SOCKADDR": sockaddr,
"CALLSIGN": callsign,
"RADIO_ID": peer,
"LAST_PING": self.clock.time(),
}
return peer

File diff suppressed because it is too large Load Diff

@ -0,0 +1,162 @@
import io
import json
import sys
import types
import unittest
def install_dmr_utils_stub():
if "dmr_utils3.utils" in sys.modules:
return None
dmr_utils3 = types.ModuleType("dmr_utils3")
utils = types.ModuleType("dmr_utils3.utils")
def bytes_4(value):
return int(value).to_bytes(4, "big")
utils.bytes_4 = bytes_4
sys.modules["dmr_utils3"] = dmr_utils3
sys.modules["dmr_utils3.utils"] = utils
return ("dmr_utils3", "dmr_utils3.utils")
class FakeRequest:
def __init__(self, path, payload=None):
self.postpath = [part.encode("utf-8") for part in path.strip("/").split("/") if part]
self.content = io.BytesIO(
b"" if payload is None else json.dumps(payload).encode("utf-8")
)
self.code = None
self.headers = {}
def setResponseCode(self, code):
self.code = code
def setHeader(self, name, value):
self.headers[name] = value
def getHeader(self, name):
if name == "content-length":
return str(len(self.content.getvalue()))
return None
class APITest(unittest.TestCase):
def setUp(self):
try:
import twisted.web.resource # noqa: F401
except ModuleNotFoundError as exc:
self.skipTest(f"Twisted is not installed: {exc}")
self.stubbed_modules = install_dmr_utils_stub()
import API
self.api = API
self.peer_id = (1234567).to_bytes(4, "big")
self.config = {
"GLOBAL": {"SYSTEM_API_KEY": "system-secret", "_KILL_SERVER": False},
"SYSTEMS": {
"MASTER-A": {
"MODE": "MASTER",
"PEERS": {self.peer_id: {}},
"_opt_key": "peer-secret",
},
"OBP-A": {
"MODE": "OPENBRIDGE",
"PEERS": {},
},
},
}
self.bridges = {}
self.controller = API.FD_APIController(self.config, self.bridges)
def tearDown(self):
if self.stubbed_modules:
for module in self.stubbed_modules:
sys.modules.pop(module, None)
def test_getoptions_returns_clear_no_options_response(self):
result = self.controller.getoptions("MASTER-A")
self.assertEqual(
result,
{"connected": True, "has_options": False, "options": ""},
)
def test_getoptions_decodes_byte_options_for_json(self):
self.config["SYSTEMS"]["MASTER-A"]["OPTIONS"] = b"KEY=peer-secret;TS1=91"
result = self.controller.getoptions("MASTER-A")
self.assertEqual(result["options"], "KEY=peer-secret;TS1=91")
self.assertTrue(result["has_options"])
def test_setoptions_stores_full_options_string_unchanged(self):
options = "KEY=peer-secret;TS1=91;DIAL=2350"
self.controller.options("MASTER-A", options)
self.assertEqual(self.config["SYSTEMS"]["MASTER-A"]["OPTIONS"], options)
def test_user_reset_is_allowed_only_for_matching_peer_key(self):
system = self.controller.validateKey(1234567, "peer-secret")
self.assertEqual(system, "MASTER-A")
self.controller.reset(system)
self.assertTrue(self.config["SYSTEMS"]["MASTER-A"]["_reset"])
self.assertFalse(self.controller.validateKey(1234567, "wrong"))
def test_system_kill_sets_existing_control_flag(self):
self.assertTrue(self.controller.validateSystemKey("system-secret"))
self.controller.killserver()
self.assertTrue(self.config["GLOBAL"]["_KILL_SERVER"])
def test_options_get_endpoint_returns_json(self):
resource = self.api.make_api_resource(self.config, self.bridges)
request = FakeRequest(
"/api/v1/options/get",
{"dmrid": 1234567, "key": "peer-secret"},
)
body = resource.render_POST(request)
self.assertEqual(request.code, 200)
self.assertEqual(
json.loads(body.decode("utf-8")),
{"ok": True, "connected": True, "has_options": False, "options": ""},
)
def test_options_get_endpoint_rejects_bad_key(self):
resource = self.api.make_api_resource(self.config, self.bridges)
request = FakeRequest(
"/api/v1/options/get",
{"dmrid": 1234567, "key": "wrong"},
)
body = resource.render_POST(request)
self.assertEqual(request.code, 401)
self.assertEqual(
json.loads(body.decode("utf-8")),
{"ok": False, "error": "invalid_credentials"},
)
def test_endpoint_rejects_large_request_body(self):
resource = self.api.make_api_resource(self.config, self.bridges)
request = FakeRequest(
"/api/v1/options/set",
{"dmrid": 1234567, "key": "peer-secret", "options": "A" * 9000},
)
body = resource.render_POST(request)
self.assertEqual(request.code, 413)
self.assertEqual(
json.loads(body.decode("utf-8")),
{"ok": False, "error": "request_too_large"},
)
if __name__ == "__main__":
unittest.main()

@ -0,0 +1,165 @@
import importlib
import io
import sys
import types
import unittest
from contextlib import redirect_stdout
class AuxiliaryToolTests(unittest.TestCase):
def test_report_receiver_bool_flag(self):
import report_receiver
self.assertTrue(report_receiver.bool_flag("1"))
self.assertTrue(report_receiver.bool_flag("true"))
self.assertTrue(report_receiver.bool_flag("yes"))
self.assertFalse(report_receiver.bool_flag("0"))
self.assertFalse(report_receiver.bool_flag(""))
self.assertFalse(report_receiver.bool_flag(None))
def test_ami_factory_builds_protocol_with_instance_state(self):
try:
import AMI
except ModuleNotFoundError as exc:
self.skipTest(str(exc))
factory = AMI.AMI.AMIClientFactory(
AMI.AMI.AMIClient,
b"user",
b"secret",
b"1234",
b"ilink 3 2350",
)
protocol = factory.buildProtocol(None)
self.assertEqual(protocol.username, b"user")
self.assertEqual(protocol.secret, b"secret")
self.assertEqual(protocol.nodenum, b"1234")
self.assertEqual(protocol.command, b"ilink 3 2350")
def test_report_sql_uses_factory_db_and_parameterized_insert(self):
self._install_mysql_stub()
try:
import report_sql
report_sql = importlib.reload(report_sql)
except ModuleNotFoundError as exc:
self.skipTest(str(exc))
fake_db = _FakeDB()
fake_reactor = object()
factory = report_sql.reportClientFactory(report_sql.reportClient, fake_db, fake_reactor)
with redirect_stdout(io.StringIO()):
client = factory.buildProtocol(None)
self.assertIs(client.db, fake_db)
self.assertIs(client.reactor, fake_reactor)
event = {
"type": "GROUP VOICE",
"event": "START",
"trx": "RX",
"system": "SYSTEM",
"streamid": "1234",
"peerid": "5678",
"subid": "9012",
"slot": "2",
"dstid": "2350",
"duration": "0",
}
with redirect_stdout(io.StringIO()):
client.send_mysql(event)
statement, params = fake_db.cursor_obj.executed
self.assertIn("%s", statement)
self.assertEqual(params[0], "GROUP VOICE")
self.assertEqual(params[8], "2350")
self.assertTrue(fake_db.committed)
self.assertTrue(fake_db.cursor_obj.closed)
def test_proxy_environment_bool_parser(self):
saved_modules = self._install_proxy_stubs()
try:
import hotspot_proxy_v2
hotspot_proxy_v2 = importlib.reload(hotspot_proxy_v2)
self.assertTrue(hotspot_proxy_v2.bool_from_env("1"))
self.assertTrue(hotspot_proxy_v2.bool_from_env("true"))
self.assertTrue(hotspot_proxy_v2.bool_from_env("yes"))
self.assertFalse(hotspot_proxy_v2.bool_from_env("0"))
self.assertFalse(hotspot_proxy_v2.bool_from_env(""))
self.assertFalse(hotspot_proxy_v2.bool_from_env(None))
finally:
self._restore_modules(saved_modules)
def _install_mysql_stub(self):
mysql_module = types.ModuleType("mysql")
connector_module = types.ModuleType("mysql.connector")
class ConnectorError(Exception):
pass
connector_module.Error = ConnectorError
connector_module.errorcode = types.SimpleNamespace(
ER_ACCESS_DENIED_ERROR=1045,
ER_BAD_DB_ERROR=1049,
)
mysql_module.connector = connector_module
sys.modules["mysql"] = mysql_module
sys.modules["mysql.connector"] = connector_module
def _install_proxy_stubs(self):
stubbed = ["dmr_utils3", "dmr_utils3.utils", "Pyro5", "Pyro5.api"]
saved_modules = {name: sys.modules.get(name) for name in stubbed + ["hotspot_proxy_v2"]}
dmr_utils3_module = types.ModuleType("dmr_utils3")
dmr_utils3_utils_module = types.ModuleType("dmr_utils3.utils")
dmr_utils3_utils_module.int_id = lambda value: int.from_bytes(value, "big")
dmr_utils3_module.utils = dmr_utils3_utils_module
pyro5_module = types.ModuleType("Pyro5")
pyro5_api_module = types.ModuleType("Pyro5.api")
pyro5_api_module.Proxy = object
pyro5_module.api = pyro5_api_module
sys.modules["dmr_utils3"] = dmr_utils3_module
sys.modules["dmr_utils3.utils"] = dmr_utils3_utils_module
sys.modules["Pyro5"] = pyro5_module
sys.modules["Pyro5.api"] = pyro5_api_module
sys.modules.pop("hotspot_proxy_v2", None)
return saved_modules
def _restore_modules(self, saved_modules):
for name, module in saved_modules.items():
if module is None:
sys.modules.pop(name, None)
else:
sys.modules[name] = module
class _FakeCursor:
def __init__(self):
self.executed = None
self.closed = False
def execute(self, statement, params):
self.executed = (statement, params)
def close(self):
self.closed = True
class _FakeDB:
def __init__(self):
self.cursor_obj = _FakeCursor()
self.committed = False
def is_connected(self):
return True
def cursor(self):
return self.cursor_obj
def commit(self):
self.committed = True
if __name__ == "__main__":
unittest.main()

@ -0,0 +1,32 @@
import ast
import pathlib
import unittest
ROOT = pathlib.Path(__file__).resolve().parents[1]
def load_bridge_helper(name):
source = (ROOT / "bridge.py").read_text()
module = ast.parse(source)
for node in module.body:
if isinstance(node, ast.FunctionDef) and node.name == name:
namespace = {}
exec(compile(ast.Module([node], []), "bridge.py", "exec"), namespace)
return namespace[name]
raise AssertionError(f"bridge.py helper not found: {name}")
class BridgeBackportTests(unittest.TestCase):
def test_dmrd_seq_delta_is_modulo_256(self):
dmrd_seq_delta = load_bridge_helper("dmrd_seq_delta")
self.assertIsNone(dmrd_seq_delta(1, False))
self.assertEqual(dmrd_seq_delta(2, 1), 1)
self.assertEqual(dmrd_seq_delta(0, 255), 1)
self.assertEqual(dmrd_seq_delta(2, 255), 3)
self.assertEqual(dmrd_seq_delta(250, 2), 248)
if __name__ == "__main__":
unittest.main()

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff
Loading…
Cancel
Save

Powered by TurnKey Linux.