rewrite of adn peer server using clean architecture

pull/4/head
Rodrigo Pérez 7 months ago
commit 30b9ef697f

95
.gitignore vendored

@ -0,0 +1,95 @@
# -----------------------------------------------------------------------------
# ADN DMR Peer Server – production config, secrets, data, build, IDE/OS
# -----------------------------------------------------------------------------
# Production config (use adn-server.example.yaml as template)
adn-server.yaml
*.local.yaml
*.local.yml
# Local/helper scripts (not for prod repo; users use their own run method)
scripts/run.sh
# Env and secrets
.env
.env.local
.env.*.local
!.env.example
user_passwords.json
encryption_key.secret
*.secret
keys.json
# -----------------------------------------------------------------------------
# Data and runtime (keep data/ via .gitkeep, do not commit contents)
# -----------------------------------------------------------------------------
data/*.json
data/*.pkl
data/*.tsv
data/*.bak
!data/.gitkeep
json/*
!json/.gitkeep
docs/
Audio/
# -----------------------------------------------------------------------------
# Python
# -----------------------------------------------------------------------------
__pycache__/
*.py[cod]
*$py.class
*.so
.Python
build/
develop-eggs/
dist/
eggs/
.eggs/
lib/
lib64/
parts/
sdist/
var/
wheels/
*.egg-info/
.installed.cfg
*.egg
.eggs/
# Virtual envs
.venv/
venv/
env/
ENV/
# -----------------------------------------------------------------------------
# Logs and runtime
# -----------------------------------------------------------------------------
*.log
log/
tmp/
temp/
*.tmp
*.temp
.cache/
# -----------------------------------------------------------------------------
# IDE and OS
# -----------------------------------------------------------------------------
.idea/
.vscode/
.cursor/
*.swp
*.swo
*~
.DS_Store
Thumbs.db
# -----------------------------------------------------------------------------
# Misc
# -----------------------------------------------------------------------------
*.bak
*.backup
*.orig
logo.png

@ -0,0 +1,30 @@
# ADN DMR Peer Server
Clean Architecture rewrite of the ADN DMR conference bridge. Same behaviour as the original server; configuration is YAML.
## License
GPL v3. Derived from ADN DMR Server / HBlink.
## Requirements
- Python 3.10+
- Dependencies: `pip install -r requirements.txt`
## Configuration
Copy `adn-server.example.yaml` to `adn-server.yaml` and edit with your settings. Production config is not committed.
## Run
```bash
pip install -r requirements.txt
python adn-server.py
```
Options:
```bash
python adn-server.py -c /path/to/adn-server.yaml
python adn-server.py --logging DEBUG
```

@ -0,0 +1,166 @@
# ADN DMR Peer Server - example configuration (YAML at project root)
# Copy to adn-server.yaml and set real secrets there (never commit adn-server.yaml).
GLOBAL:
PATH: ./
PING_TIME: 10
MAX_MISSED: 3
USE_ACL: true
REG_ACL: PERMIT:ALL
SUB_ACL: DENY:1
TGID_TS1_ACL: PERMIT:ALL
TGID_TS2_ACL: PERMIT:ALL
GEN_STAT_BRIDGES: true
ALLOW_NULL_PASSPHRASE: true
ANNOUNCEMENT_LANGUAGES: ""
SERVER_ID: 73010
DATA_GATEWAY: false
VALIDATE_SERVER_IDS: true
URL_SECURITY: "<set-in-adn-server.yaml>"
PORT_SECURITY: "<set-in-adn-server.yaml>"
PASS_SECURITY: "<set-in-adn-server.yaml>"
USERS_PASS: user_passwords.json
HASH_ENCRYPT: encryption_key.secret
REPORTS:
REPORT: true
REPORT_INTERVAL: 60
REPORT_PORT: 4321
REPORT_CLIENTS: "127.0.0.1"
LOGGER:
LOG_FILE: /var/log/adn-server/adn-server.log
LOG_HANDLERS: file-timed
LOG_LEVEL: INFO
LOG_NAME: adn-server
ALIASES:
TRY_DOWNLOAD: true
PATH: ./
PEER_FILE: peer_ids.json
SUBSCRIBER_FILE: subscriber_ids.json
TGID_FILE: talkgroup_ids.json
PEER_URL: https://adn.systems/files/peer_ids.json
SUBSCRIBER_URL: https://adn.systems/files/subscriber_ids.json
TGID_URL: https://adn.systems/files/talkgroup_ids.json
LOCAL_SUBSCRIBER_FILE: local_subcriber_ids.json
STALE_DAYS: 1
SUB_MAP_FILE: ""
SERVER_ID_URL: https://adn.systems/files/server_ids.tsv
SERVER_ID_FILE: server_ids.tsv
CHECKSUM_URL: https://adn.systems/files/file_checksums.json
CHECKSUM_FILE: file_checksums.json
KEYS_FILE: keys.json
ALLSTAR:
ENABLED: false
USER: llcgi
PASS: "<set-in-adn-server.yaml>"
SERVER: my.asl.server
PORT: 5038
NODE: "0000"
# Systems: MASTER, PEER, OPENBRIDGE. Names match legacy [SYSTEM], [D-APRS], [ECHO], [OBP-*].
SYSTEMS:
SYSTEM:
MODE: MASTER
ENABLED: true
REPEAT: true
MAX_PEERS: 2
EXPORT_AMBE: false
IP: 127.0.0.1
PORT: 56400
PASSPHRASE: "<set-in-adn-server.yaml>"
GROUP_HANGTIME: 5
USE_ACL: true
REG_ACL: DENY:1
SUB_ACL: DENY:1
TGID_TS1_ACL: PERMIT:ALL
TGID_TS2_ACL: PERMIT:ALL
DEFAULT_UA_TIMER: 60
SINGLE_MODE: false
VOICE_IDENT: false
TS1_STATIC: ""
TS2_STATIC: ""
DEFAULT_REFLECTOR: 0
ANNOUNCEMENT_LANGUAGE: es_ES
GENERATOR: 102
ALLOW_UNREG_ID: false
PROXY_CONTROL: false
OVERRIDE_IDENT_TG: ""
D-APRS:
MODE: MASTER
ENABLED: true
REPEAT: false
MAX_PEERS: 2
EXPORT_AMBE: false
IP: ""
PORT: 52555
PASSPHRASE: "<set-in-adn-server.yaml>"
GROUP_HANGTIME: 0
USE_ACL: true
REG_ACL: DENY:1
SUB_ACL: DENY:1
TGID_TS1_ACL: PERMIT:ALL
TGID_TS2_ACL: PERMIT:ALL
DEFAULT_UA_TIMER: 10
SINGLE_MODE: false
VOICE_IDENT: false
TS1_STATIC: ""
TS2_STATIC: ""
DEFAULT_REFLECTOR: 0
ANNOUNCEMENT_LANGUAGE: es_ES
GENERATOR: 2
ALLOW_UNREG_ID: true
PROXY_CONTROL: false
OVERRIDE_IDENT_TG: ""
ECHO:
MODE: PEER
ENABLED: true
LOOSE: false
EXPORT_AMBE: false
IP: 127.0.0.1
PORT: 54917
MASTER_IP: 127.0.0.1
MASTER_PORT: 54915
PASSPHRASE: "<set-in-adn-server.yaml>"
CALLSIGN: ECHO
RADIO_ID: 9990
RX_FREQ: 449000000
TX_FREQ: 444000000
TX_POWER: 25
COLORCODE: 1
SLOTS: 1
LATITUDE: "00.0000"
LONGITUDE: "000.0000"
HEIGHT: 0
LOCATION: "9990 Parrot"
DESCRIPTION: ECHO
URL: adn.systems
SOFTWARE_ID: "20170620"
PACKAGE_ID: MMDVM_ADN-Systems
GROUP_HANGTIME: 5
OPTIONS: ""
USE_ACL: true
SUB_ACL: DENY:1
TGID_TS1_ACL: PERMIT:ALL
TGID_TS2_ACL: PERMIT:ALL
ANNOUNCEMENT_LANGUAGE: en_GB
OBP-TEST:
MODE: OPENBRIDGE
ENABLED: false
IP: ""
PORT: 62044
NETWORK_ID: 1
PASSPHRASE: "<set-in-adn-server.yaml>"
TARGET_IP: ""
TARGET_PORT: 62044
USE_ACL: true
SUB_ACL: DENY:1
TGID_ACL: DENY:0-82,92-199,800-899,9990-9999,730999
RELAX_CHECKS: true
ENHANCED_OBP: true
PROTO_VER: 2

@ -0,0 +1,30 @@
#!/usr/bin/env python3
# ADN DMR Peer Server - launcher script (like monitor/monitor.py)
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# GPLv3. Derived from ADN DMR Server / HBlink.
"""
Run the ADN DMR Peer Server from the project root.
python adn-server.py
python adn-server.py -c adn-server.yaml
python adn-server.py --logging DEBUG
Config default: adn-server.yaml in this directory.
"""
from __future__ import annotations
import sys
from pathlib import Path
_ROOT = Path(__file__).resolve().parent
if str(_ROOT) not in sys.path:
sys.path.insert(0, str(_ROOT))
if str(_ROOT / "src") not in sys.path:
sys.path.insert(0, str(_ROOT / "src"))
from adn_server.main import main
if __name__ == "__main__":
main()

@ -0,0 +1,31 @@
# ADN DMR Peer Server
# GPL v3. Derived from ADN DMR Server / HBlink.
[build-system]
requires = ["setuptools>=61", "wheel"]
build-backend = "setuptools.build_meta"
[project]
name = "adn-server"
version = "0.1.0"
description = "ADN DMR Peer Server"
readme = "README.md"
license = { text = "GPL-3.0-or-later" }
requires-python = ">=3.10"
dependencies = [
"Twisted>=22.0",
"dmr_utils3>=0.1.19",
"pyyaml>=6.0",
"setproctitle>=1.3",
"bitarray>=2.0",
"cryptography>=41.0",
]
[project.optional-dependencies]
dev = ["pytest>=7", "pytest-cov"]
[tool.setuptools.packages.find]
where = ["src"]
[project.scripts]
adn-server = "adn_server.main:main"

@ -0,0 +1,9 @@
# ADN DMR Peer Server - dependencies
# Python 3.10+
Twisted>=22.0.0
dmr_utils3>=0.1.19
pyyaml>=6.0
setproctitle>=1.3.0
bitarray>=2.0.0
cryptography>=41.0.0

@ -0,0 +1,25 @@
# ADN DMR Peer Server
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
# 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, see <http://www.gnu.org/licenses/>.
#
# Derived from: ADN DMR Server / HBlink (bridge_master.py, hblink.py):
# Copyright (C) 2026 Joaquin Madrid Belando, EA5GVK <ea5gvk@gmail.com>
# Copyright (C) 2025 Esteban Mackay, HP3ICC <setcom40@gmail.com>
# Copyright (C) 2025 Bruno Farias, CS8ABG <cs8abg@gmail.com>
# Copyright (C) 2020-2023 Simon Adlem, G7RZU <g7rzu@gb7fr.org.uk>
# Copyright (C) 2016-2019 Cortney T. Buffington, N0MJS <n0mjs@me.com>
# Original works and this derivative are under GPLv3.
"""ADN DMR Peer Server — conference bridge (rewrite of bridge_master)."""

@ -0,0 +1,33 @@
# ADN DMR Peer Server - application layer
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
from .ports import (
ConfigLoader,
AliasLoader,
SubMapStore,
KeysStore,
ReportSender,
BridgeRouter,
VoiceProvider,
SecurityDownloader,
)
from .bridge_use_cases import BridgeUseCases
from .ident_use_cases import IdentUseCases
from .voice_use_cases import VoiceUseCases
from .reporting_use_cases import ReportingUseCases
__all__ = [
"ConfigLoader",
"AliasLoader",
"SubMapStore",
"KeysStore",
"ReportSender",
"BridgeRouter",
"VoiceProvider",
"SecurityDownloader",
"BridgeUseCases",
"IdentUseCases",
"VoiceUseCases",
"ReportingUseCases",
]

File diff suppressed because it is too large Load Diff

@ -0,0 +1,139 @@
# ADN DMR Peer Server - voice ident use cases
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink (threadIdent, ident). GPLv3.
"""Voice ident: periodic ident on MASTER systems (VOICE_IDENT, slot 2 idle 30s)."""
from __future__ import annotations
import logging
import re
import time
from typing import Any, Callable
from ..domain import bytes_3, int_id
from ..infrastructure.hbp_constants import HBPF_SLT_VTERM
logger = logging.getLogger(__name__)
def _alias_tg(dst_id: bytes, config: dict[str, Any]) -> str:
"""Resolve TGID to alias for logging; fallback to numeric."""
tg_ids = config.get("_TG_IDS", {})
idx = int_id(dst_id)
return tg_ids.get(idx, str(idx))
class IdentUseCases:
"""Run voice ident for MASTER systems with VOICE_IDENT (legacy threadIdent/ident)."""
def __init__(
self,
config: dict[str, Any],
voice_use_cases: Any,
audio_path: str,
get_protocols: Callable[[], dict[str, Any]],
call_from_reactor: Callable[..., None],
) -> None:
self._config = config
self._voice = voice_use_cases
self._audio_path = audio_path
self._get_protocols = get_protocols
self._call_from_reactor = call_from_reactor
def run_ident(self) -> None:
"""Run ident once (legacy ident()). Call from thread; uses call_from_reactor to send packets."""
systems_cfg = self._config.get("SYSTEMS", {})
protocols = self._get_protocols()
ann_lang = (self._config.get("VOICE", {}).get("ANNOUNCEMENT_LANGUAGES") or "").strip()
if not ann_lang:
return
words_by_lang = self._voice.get_ambe_words(ann_lang, self._audio_path)
if not words_by_lang:
return
for system in list(systems_cfg.keys()):
sys_cfg = systems_cfg.get(system, {})
if sys_cfg.get("MODE") != "MASTER":
continue
if not sys_cfg.get("VOICE_IDENT"):
continue
_lang = sys_cfg.get("ANNOUNCEMENT_LANGUAGE", "en_GB")
if _lang not in words_by_lang:
continue
words = words_by_lang[_lang]
max_peers = int(sys_cfg.get("MAX_PEERS", 1))
if max_peers > 1:
logger.debug("(IDENT) %s System has MAX_PEERS > 1, skipping", system)
continue
_callsign = None
peers = sys_cfg.get("PEERS", {})
for _peerid in peers:
peer_cfg = peers.get(_peerid, {})
if isinstance(peer_cfg, dict) and peer_cfg.get("CALLSIGN"):
cs = peer_cfg["CALLSIGN"]
_callsign = cs.decode("utf-8", errors="replace") if isinstance(cs, bytes) else cs
break
if not _callsign:
logger.debug("(IDENT) %s System has no peers or no recorded callsign, skipping", system)
continue
protocol = protocols.get(system)
if not protocol or not getattr(protocol, "STATUS", None):
continue
_slot = protocol.STATUS.get(2)
if not _slot:
continue
rx_type = _slot.get("RX_TYPE")
tx_type = _slot.get("TX_TYPE")
tx_time = _slot.get("TX_TIME", 0)
rx_time = _slot.get("RX_TIME", 0)
now = time.time()
if (rx_type != HBPF_SLT_VTERM or tx_type != HBPF_SLT_VTERM or
now - tx_time <= 30 or now - rx_time <= 30):
continue
_all_call = bytes_3(16777215)
_source_id = bytes_3(5000)
_dst_id = b""
override_tg = sys_cfg.get("OVERRIDE_IDENT_TG")
if override_tg is not None and int(override_tg) > 0 and int(override_tg) < 16777215:
_dst_id = bytes_3(int(override_tg))
else:
_dst_id = _all_call
logger.info(
"(%s) %s System idle. Sending voice ident to TG %s",
system, _callsign, _alias_tg(_dst_id, self._config),
)
silence = words.get("silence")
if not silence:
continue
_say = [silence, silence, silence, words.get("this-is") or silence]
for _ in range(6):
_say.append(silence)
_systemcs = re.sub(r"\W+", "", _callsign).upper()
for character in _systemcs:
_say.append(words.get(character) or silence)
_say.append(silence)
for _ in range(5):
_say.append(silence)
_say.append(words.get("adn") or silence)
server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00")
if not isinstance(server_id, bytes):
server_id = bytes_3(int(server_id))
speech = self._voice.pkt_gen(_source_id, _dst_id, server_id, 1, _say)
time.sleep(1)
_slot = protocol.STATUS.get(2)
if not _slot:
continue
_next_time = time.time()
for pkt in speech:
_next_time += 0.058
_delay = _next_time - time.time()
if _delay > 0.001:
time.sleep(_delay)
self._call_from_reactor(protocol.send_voice_packet, pkt, _source_id, _dst_id, _slot)

@ -0,0 +1,143 @@
# ADN DMR Peer Server - application ports
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Abstract ports (interfaces) for infrastructure adapters."""
from __future__ import annotations
from abc import ABC, abstractmethod
from typing import Any
class ConfigLoader(ABC):
"""Load and validate config (YAML at project root). Returns same semantic structure as legacy INI."""
@abstractmethod
def load(self, path: str | None = None) -> dict[str, Any]:
"""Load config; path None = default (project root adn-server.yaml)."""
...
@abstractmethod
def reload_voice_config(self, config: dict[str, Any]) -> None:
"""Reload voice/announcement config from main config file (optional second arg: path)."""
...
class AliasLoader(ABC):
"""Load and refresh alias dicts (peer_ids, subscriber_ids, talkgroup_ids, server_ids, checksums)."""
@abstractmethod
def load_aliases(self, config: dict[str, Any]) -> tuple[
dict[int, str], # peer_ids
dict[int, str], # subscriber_ids
dict[int, str], # talkgroup_ids
dict[int, str], # local_subscriber_ids
dict[str, str], # server_ids
dict[str, str], # checksums
]:
"""Build alias dicts from config (downloads + files). Same as legacy mk_aliases."""
...
class SubMapStore(ABC):
"""Persist and load SUB_MAP (bytes_3(peer) -> (callsign, slot, time))."""
@abstractmethod
def load(self, path: str) -> dict[bytes, tuple[str, int, float]]:
"""Load SUB_MAP from pickle file."""
...
@abstractmethod
def save(self, path: str, sub_map: dict[bytes, tuple[str, int, float]]) -> None:
"""Save SUB_MAP to pickle file."""
...
class KeysStore(ABC):
"""Load/save keys JSON (e.g. system API key)."""
@abstractmethod
def load(self, path: str) -> dict[str, Any]:
"""Load keys from JSON."""
...
@abstractmethod
def save(self, path: str, keys: dict[str, Any]) -> None:
"""Save keys to JSON."""
...
class ReportSender(ABC):
"""Send config and bridge state to report TCP clients (CONFIG_SND, BRIDGE_SND, BRDG_EVENT)."""
@abstractmethod
def send_config(self, systems: dict[str, Any]) -> None:
"""Send CONFIG_SND (pickle systems)."""
...
@abstractmethod
def send_bridge(self, bridges: dict[str, Any]) -> None:
"""Send BRIDGE_SND (pickle bridges)."""
...
@abstractmethod
def send_bridge_event(self, event: str) -> None:
"""Send BRDG_EVENT (opcode + event string)."""
...
class BridgeRouter(ABC):
"""Query and update BRIDGES (conference bridge state). Used by rule_timer, make_single_bridge, etc."""
@abstractmethod
def get_bridges(self) -> dict[str, list[dict[str, Any]]]:
"""Return current BRIDGES dict (key = TGID or #reflector)."""
...
@abstractmethod
def set_bridges(self, bridges: dict[str, list[dict[str, Any]]]) -> None:
"""Replace BRIDGES (e.g. after rule_timer or make_single_bridge)."""
...
@abstractmethod
def acl_check(self, id_bytes_or_int: bytes | int, acl: tuple[bool, list[tuple[int, int]]]) -> bool:
"""Check ID against ACL; return True if permitted."""
...
class VoiceProvider(ABC):
"""AMBE voice packets, file announcements, TTS. Used for announcements and playback."""
@abstractmethod
def get_ambe_words(self, languages: str, audio_path: str) -> dict[str, dict[str, Any]]:
"""Load AMBE words by language (readAMBE.readfiles)."""
...
@abstractmethod
def pkt_gen(self, rf_src: bytes, dst_id: bytes, peer: bytes, slot: int, phrase: list[Any]) -> Any:
"""Generate HBP voice packets for a phrase (generator). Legacy mk_voice.pkt_gen."""
...
@abstractmethod
def ensure_tts_ambe(self, text: str, lang: str, out_path: str, config: dict[str, Any]) -> str | None:
"""TTS to AMBE file; return path or None. Legacy tts_engine.ensure_tts_ambe."""
...
def read_single_file(self, audio_path: str, lang: str, file_number: str) -> list:
"""Read one AMBE file (e.g. ondemand/{file_number}.ambe). Legacy readSingleFile for playFileOnRequest."""
return []
class SecurityDownloader(ABC):
"""Periodic security downloads (passwords, encryption). Legacy security_downloader."""
@abstractmethod
def init_downloads(self, config: dict[str, Any]) -> None:
"""One-time init (e.g. create dirs)."""
...
@abstractmethod
def periodic_download(self, config: dict[str, Any]) -> None:
"""Periodic password/encryption download."""
...

@ -0,0 +1,50 @@
# ADN DMR Peer Server - reporting use cases
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Reporting: send config/bridge to TCP clients, KA reporting. Orchestrates ReportSender."""
from __future__ import annotations
import logging
import time
from typing import Any
from .ports import ReportSender
logger = logging.getLogger(__name__)
class ReportingUseCases:
"""Use cases for TCP report server and keepalive reporting."""
def __init__(self, report_sender: ReportSender, config: dict[str, Any]) -> None:
self._sender = report_sender
self._config = config
def send_config(self, systems: dict[str, Any]) -> None:
"""Send CONFIG_SND to all report clients."""
self._sender.send_config(systems)
def send_bridge(self, bridges: dict[str, Any]) -> None:
"""Send BRIDGE_SND to all report clients."""
self._sender.send_bridge(bridges)
def send_bridge_event(self, event: str) -> None:
"""Send BRDG_EVENT to all report clients."""
self._sender.send_bridge_event(event)
def ka_reporting_loop(self) -> None:
"""Legacy kaReporting (60s): check OBP keepalive status and log warnings for stale connections."""
logger.debug("(ROUTER) KeepAlive reporting loop started")
systems_cfg = self._config.get("SYSTEMS", {})
now = time.time()
for system_name, sys_cfg in systems_cfg.items():
if sys_cfg.get("MODE") == "OPENBRIDGE" and sys_cfg.get("ENHANCED_OBP"):
if "_bcka" not in sys_cfg:
logger.warning("(ROUTER) not sending to system %s as KeepAlive never seen", system_name)
elif sys_cfg["_bcka"] < now - 60:
logger.warning(
"(ROUTER) not sending to system %s as last KeepAlive was %s seconds ago",
system_name, int(now - sys_cfg["_bcka"]),
)

@ -0,0 +1,613 @@
# ADN DMR Peer Server - voice use cases
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Voice/AMBE/TTS: scheduled announcements, TTS announcements, playback. Orchestrates VoiceProvider."""
from __future__ import annotations
import logging
import os
import time
from datetime import datetime
from typing import Any, Callable
from ..domain import bytes_3
from ..infrastructure.hbp_constants import HBPF_SLT_VTERM
from ..infrastructure.voice.tts_engine import ensure_tts_ambe as tts_ensure_tts_ambe
from .ports import VoiceProvider
logger = logging.getLogger(__name__)
_FRAME_INTERVAL = 0.058
_ANNOUNCEMENT_EXCLUDED = ("ECHO", "D-APRS")
_BROADCAST_GAP = 1.5
class VoiceUseCases:
"""Use cases for voice announcements and TTS."""
def __init__(
self,
voice_provider: VoiceProvider,
config: dict[str, Any],
get_protocols: Callable[[], dict[str, Any]] | None = None,
call_from_reactor: Callable[..., None] | None = None,
audio_path: str | None = None,
get_bridges: Callable[[], dict[str, list[dict[str, Any]]]] | None = None,
call_later: Callable[..., Any] | None = None,
start_looping_call: Callable[[Callable[[], None], float, bool], Any] | None = None,
defer_to_thread: Callable[..., Any] | None = None,
) -> None:
self._voice = voice_provider
self._config = config
self._get_protocols = get_protocols
self._call_from_reactor = call_from_reactor
self._audio_path = audio_path or ""
self._get_bridges = get_bridges
self._call_later = call_later
self._start_looping_call = start_looping_call
self._defer_to_thread = defer_to_thread
self._ann_tasks: dict[int, Any] = {}
self._tts_tasks: dict[int, Any] = {}
self._announcement_running: dict[int, bool] = {1: False, 2: False, 3: False, 4: False}
self._tts_running: dict[int, bool] = {1: False, 2: False, 3: False, 4: False}
self._announcement_last_hour: dict[int, int] = {1: -1, 2: -1, 3: -1, 4: -1}
self._tts_last_hour: dict[int, int] = {1: -1, 2: -1, 3: -1, 4: -1}
self._config_file_mtime: float = 0.0
self._broadcast_queue: list[dict[str, Any]] = []
self._broadcast_active: bool = False
def get_ambe_words(self, languages: str, audio_path: str) -> dict[str, dict[str, Any]]:
"""Load AMBE words for given languages (legacy readAMBE.readfiles)."""
return self._voice.get_ambe_words(languages, audio_path)
def pkt_gen(self, rf_src: bytes, dst_id: bytes, peer: bytes, slot: int, phrase: list[Any]) -> Any:
"""Generate HBP voice packets for phrase (legacy mk_voice.pkt_gen)."""
return self._voice.pkt_gen(rf_src, dst_id, peer, slot, phrase)
def _build_announcement_targets(
self, tg_int: int, tg_str: str, label: str
) -> list[dict[str, Any]]:
"""MASTER systems with active bridge for tg_int and idle slot (legacy target list)."""
targets: list[dict[str, Any]] = []
protocols = self._get_protocols() if self._get_protocols else {}
systems_cfg = self._config.get("SYSTEMS", {})
bridges = self._get_bridges() if self._get_bridges else {}
bridge_entries = bridges.get(tg_str, [])
for sys_name in list(protocols.keys()):
if sys_name in _ANNOUNCEMENT_EXCLUDED or any(
sys_name.startswith(ex + "-") for ex in _ANNOUNCEMENT_EXCLUDED
):
continue
if sys_name not in systems_cfg or systems_cfg[sys_name].get("MODE") != "MASTER":
continue
if not systems_cfg[sys_name].get("PEERS"):
continue
has_peers = any(
systems_cfg[sys_name]["PEERS"].get(pid, {}).get("CALLSIGN")
for pid in systems_cfg[sys_name]["PEERS"]
)
if not has_peers or sys_name not in protocols:
continue
sys_obj = protocols[sys_name]
if not getattr(sys_obj, "STATUS", None):
continue
active_slots = [
be["TS"] for be in bridge_entries
if be.get("SYSTEM") == sys_name and be.get("ACTIVE") and be.get("TS") is not None
]
active_slots = list(dict.fromkeys(active_slots))
if not active_slots:
continue
for ts in active_slots:
slot_index = 2 if ts == 2 else 1
slot = sys_obj.STATUS.get(slot_index)
if not slot:
continue
rx_type = slot.get("RX_TYPE")
tx_type = slot.get("TX_TYPE")
if (rx_type != HBPF_SLT_VTERM) or (tx_type != HBPF_SLT_VTERM):
logger.debug("(%s) System %s TS%s busy, skipping", label, sys_name, ts)
continue
targets.append({"sys_obj": sys_obj, "name": sys_name, "slot": slot, "ts": ts})
return targets
def _send_filtered_by_tg(
self, sys_obj: Any, pkt: bytes, tg: int, ts: int, bridges: dict[str, list[dict[str, Any]]]
) -> int:
"""Return -1 if sent, 0 if TG/TS not active (legacy _sendFilteredByTG)."""
tg_str = str(tg)
for be in bridges.get(tg_str, []):
if be.get("SYSTEM") == getattr(sys_obj, "_system", None) and be.get("TS") == ts and be.get("ACTIVE"):
sys_obj.send_system(pkt)
return -1
return 0
def _enqueue_broadcast(
self, _type: str, targets: list[dict[str, Any]], pkts_by_ts: dict[int, list[bytes]],
source_id: bytes, dst_id: bytes, tg: int, num: int, label: str,
) -> None:
self._broadcast_queue.append({
'type': _type, 'targets': targets, 'pkts_by_ts': pkts_by_ts,
'source_id': source_id, 'dst_id': dst_id, 'tg': tg, 'num': num, 'label': label,
})
_pos = len(self._broadcast_queue)
if self._broadcast_active:
logger.info('(%s) Enqueued broadcast (position %s in queue)', label, _pos)
else:
self._start_next_broadcast()
def _start_next_broadcast(self) -> None:
if not self._broadcast_queue:
self._broadcast_active = False
return
self._broadcast_active = True
_item = self._broadcast_queue.pop(0)
_type = _item['type']
_label = _item['label']
logger.info('(%s) Starting broadcast from queue (%s remaining)', _label, len(self._broadcast_queue))
if self._call_later:
if _type == 'ann':
self._call_later(0.5, self._announcement_send_broadcast, _item['targets'], _item['pkts_by_ts'], 0, _item['source_id'], _item['dst_id'], _item['tg'], _item['num'], _label, None)
elif _type == 'tts':
self._call_later(0.5, self._tts_send_broadcast, _item['targets'], _item['pkts_by_ts'], 0, _item['source_id'], _item['dst_id'], _item['tg'], _item['num'], _label, None)
def _broadcast_finished(self) -> None:
if self._broadcast_queue:
logger.info('(QUEUE) Broadcast finished, next in %.1fs (%s queued)', _BROADCAST_GAP, len(self._broadcast_queue))
if self._call_later:
self._call_later(_BROADCAST_GAP, self._start_next_broadcast)
else:
self._broadcast_active = False
logger.info('(QUEUE) Broadcast finished, queue empty')
def _announcement_send_broadcast(
self,
targets: list[dict[str, Any]],
pkts_by_ts: dict[int, list[bytes]],
pkt_idx: int,
source_id: bytes,
dst_id: bytes,
tg: int,
ann_num: int,
label: str,
next_time: float | None = None,
) -> None:
"""Send one batch of packets; schedule next via call_later (legacy _announcementSendBroadcast)."""
total = len(pkts_by_ts.get(1, []))
if pkt_idx >= total:
for t in targets:
try:
obj = t.get("sys_obj")
if getattr(obj, "STATUS", None):
for sid in list(obj.STATUS.keys()):
if sid not in (1, 2):
del obj.STATUS[sid]
except Exception:
pass
self._announcement_running[ann_num] = False
logger.info("(%s) Broadcast complete: %s packets sent to %s targets", label, total, len(targets))
self._broadcast_finished()
return
bridges = self._get_bridges() if self._get_bridges else {}
now = time.time()
for t in targets:
try:
sys_obj = t["sys_obj"]
slot = t["slot"]
t_ts = t["ts"]
pkt = pkts_by_ts[t_ts][pkt_idx]
stream_id = pkt[16:20]
if stream_id not in sys_obj.STATUS:
sys_obj.STATUS[stream_id] = {
"START": now,
"CONTENTION": False,
"RFS": source_id,
"TGID": dst_id,
"LAST": now,
}
slot["TX_TGID"] = dst_id
else:
sys_obj.STATUS[stream_id]["LAST"] = now
slot["TX_TIME"] = now
self._send_filtered_by_tg(sys_obj, pkt, tg, t_ts, bridges)
except Exception as e:
logger.error("(%s) Error sending packet %s to %s/TS%s: %s", label, pkt_idx, t.get("name"), t.get("ts"), e)
if next_time is None:
next_time = now + _FRAME_INTERVAL
else:
next_time = next_time + _FRAME_INTERVAL
delay = max(0.001, next_time - time.time())
if self._call_later:
self._call_later(
delay,
self._announcement_send_broadcast,
targets,
pkts_by_ts,
pkt_idx + 1,
source_id,
dst_id,
tg,
ann_num,
label,
next_time,
)
def scheduled_announcement(self, ann_num: int = 1, _retry: int = 0) -> None:
"""Run one scheduled file announcement (legacy scheduledAnnouncement 1–4)."""
g = self._config.get("VOICE", {})
prefix = "ANNOUNCEMENT" if ann_num == 1 else "ANNOUNCEMENT{}".format(ann_num)
label = "ANNOUNCEMENT" if ann_num == 1 else "ANNOUNCEMENT-{}".format(ann_num)
if not g.get("{}_ENABLED".format(prefix)):
return
if self._announcement_running.get(ann_num):
if _retry == 0:
logger.debug("(%s) Previous announcement still running, skipping", label)
return
mode = g.get("{}_MODE".format(prefix), "interval")
if mode == "hourly" and _retry == 0:
now = datetime.now()
if now.minute != 0:
return
if self._announcement_last_hour.get(ann_num) == now.hour:
return
self._announcement_last_hour[ann_num] = now.hour
if self._broadcast_active and _retry < 60:
if _retry == 0:
logger.debug("(%s) Broadcast queue busy, deferring prep", label)
if self._call_later:
self._call_later(3.0 + ann_num * 0.5, self.scheduled_announcement, ann_num, _retry + 1)
return
_file = g.get("{}_FILE".format(prefix), "")
_tg = int(g.get("{}_TG".format(prefix), 0))
_lang = g.get("{}_LANGUAGE".format(prefix), "en_GB")
if not _file or not _tg:
return
_dst_id = bytes_3(_tg)
_source_id = bytes_3(5000)
server_id = g.get("SERVER_ID", b"\x00\x00\x00\x00")
if not isinstance(server_id, bytes):
server_id = bytes_3(int(server_id))
logger.info("(%s) Playing file: %s to TG %s (both TS, mode: %s, lang: %s)", label, _file, _tg, mode, _lang)
try:
_say = self.read_single_file(self._audio_path, _lang, str(_file))
except Exception as e:
logger.warning("(%s) Cannot read AMBE file: Audio/%s/ondemand/%s.ambe: %s", label, _lang, _file, e)
return
if not _say:
logger.warning("(%s) AMBE file empty or not found: %s/ondemand/%s.ambe", label, _lang, _file)
return
tg_str = str(_tg)
targets = self._build_announcement_targets(_tg, tg_str, label)
if not targets:
logger.info("(%s) No systems with active bridge for TG %s to send to", label, _tg)
return
_say_list = [_say]
pkts_by_ts = {
1: list(self.pkt_gen(_source_id, _dst_id, server_id, 0, _say_list)),
2: list(self.pkt_gen(_source_id, _dst_id, server_id, 1, _say_list)),
}
ts1_count = sum(1 for t in targets if t["ts"] == 1)
ts2_count = sum(1 for t in targets if t["ts"] == 2)
sys_names = ", ".join("{}/TS{}".format(t["name"], t["ts"]) for t in targets[:8])
if len(targets) > 8:
sys_names += ", ... +{}".format(len(targets) - 8)
logger.info(
"(%s) Broadcasting %s packets to %s targets (TS1:%s TS2:%s): %s",
label, len(pkts_by_ts[1]), len(targets), ts1_count, ts2_count, sys_names,
)
self._announcement_running[ann_num] = True
self._enqueue_broadcast('ann', targets, pkts_by_ts, _source_id, _dst_id, _tg, ann_num, label)
def scheduled_tts_announcement(self, tts_num: int = 1, _retry: int = 0) -> None:
"""Run one scheduled TTS announcement (legacy scheduledTTSAnnouncement 1–4)."""
g = self._config.get("VOICE", {})
prefix = "TTS_ANNOUNCEMENT{}".format(tts_num)
label = "TTS-{}".format(tts_num)
if not g.get("{}_ENABLED".format(prefix), False):
return
if self._tts_running.get(tts_num):
if _retry == 0:
logger.debug("(%s) Previous TTS announcement still running, skipping", label)
return
mode = g.get("{}_MODE".format(prefix), "interval")
if mode == "hourly" and _retry == 0:
now = datetime.now()
if now.minute != 0:
return
if self._tts_last_hour.get(tts_num) == now.hour:
return
self._tts_last_hour[tts_num] = now.hour
if self._broadcast_active and _retry < 60:
if _retry == 0:
logger.debug("(%s) Broadcast queue busy, deferring TTS prep", label)
if self._call_later:
self._call_later(3.0 + tts_num * 0.5, self.scheduled_tts_announcement, tts_num, _retry + 1)
return
_file = g.get("{}_FILE".format(prefix), "")
_tg = int(g.get("{}_TG".format(prefix), 0))
_lang = g.get("{}_LANGUAGE".format(prefix), "en_GB")
self._tts_running[tts_num] = True
logger.info("(%s) Starting TTS conversion in background thread for %s", label, _file)
if self._defer_to_thread:
d = self._defer_to_thread(tts_ensure_tts_ambe, self._config, tts_num, self._audio_path)
d.addCallback(self._tts_conversion_done, tts_num, _file, _tg, _lang, mode, label)
d.addErrback(self._tts_conversion_error, tts_num, label)
else:
try:
ambe_path = tts_ensure_tts_ambe(self._config, tts_num, self._audio_path)
self._tts_conversion_done(ambe_path, tts_num, _file, _tg, _lang, mode, label)
except Exception as e:
self._tts_conversion_error(e, tts_num, label)
def _tts_conversion_done(
self, ambe_path: str | None, tts_num: int, _file: str, _tg: int, _lang: str, mode: str, label: str, _retry: int = 0
) -> None:
"""After TTS conversion: broadcast like scheduled_announcement (legacy _ttsConversionDone)."""
if not ambe_path:
self._tts_running[tts_num] = False
logger.warning("(%s) No AMBE file available for TTS announcement %s", label, _file)
return
if self._broadcast_active and _retry < 60:
if _retry == 0:
logger.debug('(%s) Broadcast queue busy, deferring TTS packet prep', 'TTS-{}'.format(tts_num))
if self._call_later:
self._call_later(3.0 + tts_num * 0.5, self._tts_conversion_done, ambe_path, tts_num, _file, _tg, _lang, mode, label, _retry + 1)
return
logger.info("(%s) Playing TTS file: %s to TG %s (both TS, mode: %s, lang: %s)", label, _file, _tg, mode, _lang)
_dst_id = bytes_3(_tg)
_source_id = bytes_3(5000)
server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00")
if not isinstance(server_id, bytes):
server_id = bytes_3(int(server_id))
_file_base = _file.replace(".ambe", "")
_say = self.read_single_file(self._audio_path, _lang, _file_base)
if not _say:
logger.warning("(%s) Cannot read AMBE file: %s", label, ambe_path)
self._tts_running[tts_num] = False
return
tg_str = str(_tg)
targets = self._build_announcement_targets(_tg, tg_str, label)
if not targets:
self._tts_running[tts_num] = False
logger.info("(%s) No systems with active bridge for TG %s to send to", label, _tg)
return
_say_list = [_say]
pkts_by_ts = {
1: list(self.pkt_gen(_source_id, _dst_id, server_id, 0, _say_list)),
2: list(self.pkt_gen(_source_id, _dst_id, server_id, 1, _say_list)),
}
logger.info("(%s) Broadcasting %s packets to %s targets", label, len(pkts_by_ts[1]), len(targets))
self._enqueue_broadcast('tts', targets, pkts_by_ts, _source_id, _dst_id, _tg, tts_num, label)
def _tts_conversion_error(self, failure: Any, tts_num: int, label: str) -> None:
self._tts_running[tts_num] = False
try:
msg = failure.getErrorMessage()
except Exception:
msg = str(failure)
logger.error("(%s) TTS conversion error: %s", label, msg)
def _tts_send_broadcast(
self,
targets: list[dict[str, Any]],
pkts_by_ts: dict[int, list[bytes]],
pkt_idx: int,
source_id: bytes,
dst_id: bytes,
tg: int,
tts_num: int,
label: str,
next_time: float | None = None,
) -> None:
"""Same as _announcement_send_broadcast but clears _tts_running (legacy _ttsSendBroadcast)."""
total = len(pkts_by_ts.get(1, []))
if pkt_idx >= total:
for t in targets:
try:
obj = t.get("sys_obj")
if getattr(obj, "STATUS", None):
for sid in list(obj.STATUS.keys()):
if sid not in (1, 2):
del obj.STATUS[sid]
except Exception:
pass
self._tts_running[tts_num] = False
logger.info("(%s) Broadcast complete: %s packets sent to %s targets", label, total, len(targets))
self._broadcast_finished()
return
bridges = self._get_bridges() if self._get_bridges else {}
now = time.time()
for t in targets:
try:
sys_obj = t["sys_obj"]
slot = t["slot"]
t_ts = t["ts"]
pkt = pkts_by_ts[t_ts][pkt_idx]
stream_id = pkt[16:20]
if stream_id not in sys_obj.STATUS:
sys_obj.STATUS[stream_id] = {
"START": now,
"CONTENTION": False,
"RFS": source_id,
"TGID": dst_id,
"LAST": now,
}
slot["TX_TGID"] = dst_id
else:
sys_obj.STATUS[stream_id]["LAST"] = now
slot["TX_TIME"] = now
self._send_filtered_by_tg(sys_obj, pkt, tg, t_ts, bridges)
except Exception as e:
logger.error("(%s) Error sending packet %s to %s: %s", label, pkt_idx, t.get("name"), e)
if next_time is None:
next_time = now + _FRAME_INTERVAL
else:
next_time = next_time + _FRAME_INTERVAL
delay = max(0.001, next_time - time.time())
if self._call_later:
self._call_later(
delay,
self._tts_send_broadcast,
targets,
pkts_by_ts,
pkt_idx + 1,
source_id,
dst_id,
tg,
tts_num,
label,
next_time,
)
def read_single_file(self, audio_path: str, lang: str, file_number: str) -> list:
"""Read one AMBE file (e.g. ondemand/{file_number}.ambe). Legacy readSingleFile."""
return self._voice.read_single_file(audio_path, lang, file_number)
def play_file_on_request(self, file_number: str, system: str) -> None:
"""Play AMBE file on request (legacy playFileOnRequest). TG 9991-9999 triggers this."""
if not self._get_protocols or not self._call_from_reactor or not self._audio_path:
return
protocol = self._get_protocols().get(system)
if not protocol or not getattr(protocol, "STATUS", None):
return
sys_cfg = self._config.get("SYSTEMS", {}).get(system, {})
lang = sys_cfg.get("ANNOUNCEMENT_LANGUAGE", "en_GB")
pairs = self.read_single_file(self._audio_path, lang, file_number)
if not pairs:
logger.warning("(%s) AMBE file not found or empty: %s/ondemand/%s.ambe", system, lang, file_number)
return
logger.info("(%s) Playing on-demand AMBE file: %s (ID: %s)", system, file_number, file_number)
time.sleep(1)
_say = [pairs]
server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00")
if not isinstance(server_id, bytes):
server_id = bytes_3(int(server_id))
speech = self.pkt_gen(bytes_3(5000), bytes_3(9), server_id, 1, _say)
_slot = protocol.STATUS.get(2)
if not _slot:
return
_source_id = bytes_3(5000)
_dst_id = bytes_3(9)
_next_time = time.time()
_pkt_count = 0
for pkt in speech:
_next_time += 0.058
delay = _next_time - time.time()
if delay > 0.001:
time.sleep(delay)
self._call_from_reactor(protocol.send_voice_packet, pkt, _source_id, _dst_id, _slot)
_pkt_count += 1
logger.info("(%s) On-demand playback complete: %s (%d packets)", system, file_number, _pkt_count)
def disconnected_voice(self, system: str) -> None:
"""Send 'disconnected' / 'linked to reflector' voice (legacy disconnectedVoice). Run from thread."""
if not self._get_protocols or not self._call_from_reactor or not self._audio_path:
return
protocol = self._get_protocols().get(system)
if not protocol or not getattr(protocol, "STATUS", None):
return
ann_lang = (self._config.get("VOICE", {}).get("ANNOUNCEMENT_LANGUAGES") or "").strip()
if not ann_lang:
return
words_by_lang = self.get_ambe_words(ann_lang, self._audio_path)
sys_cfg = self._config.get("SYSTEMS", {}).get(system, {})
_lang = sys_cfg.get("ANNOUNCEMENT_LANGUAGE", "en_GB")
if _lang not in words_by_lang:
return
words = words_by_lang[_lang]
silence = words.get("silence")
if not silence:
return
_say = [silence, silence]
default_refl = int(sys_cfg.get("DEFAULT_REFLECTOR", 0))
if default_refl > 0:
_say.append(silence)
_say.append(words.get("linkedto") or silence)
_say.append(silence)
_say.append(words.get("to") or silence)
_say.append(silence)
_say.append(silence)
for digit in str(default_refl):
_say.append(words.get(digit) or silence)
_say.append(silence)
else:
_say.append(words.get("notlinked") or silence)
_say.append(silence)
server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00")
if not isinstance(server_id, bytes):
server_id = bytes_3(int(server_id))
speech = self.pkt_gen(bytes_3(5000), bytes_3(9), server_id, 1, _say)
time.sleep(1)
_slot = protocol.STATUS.get(2)
if not _slot:
return
logger.debug("(%s) Sending disconnected voice", system)
_next_time = time.time()
for pkt in speech:
_next_time += 0.058
_delay = _next_time - time.time()
if _delay > 0.001:
time.sleep(_delay)
self._call_from_reactor(protocol.send_voice_packet, pkt, bytes_3(5000), bytes_3(9), _slot)
logger.debug("(%s) disconnected voice thread end", system)
def check_voice_config_reload(self, config_file_path: str | None = None) -> None:
"""Check main config file mtime and reload if changed (15s loop). Start/stop announcement LoopingCalls."""
if config_file_path and os.path.isfile(config_file_path):
try:
mtime = os.path.getmtime(config_file_path)
except OSError:
return
if mtime == self._config_file_mtime:
return
self._config_file_mtime = mtime
logger.info("(VOICE-RELOAD) config file change detected, reloading configuration...")
g = self._config.get("VOICE", {})
if not self._start_looping_call:
return
for ann_num in range(1, 5):
prefix = "ANNOUNCEMENT" if ann_num == 1 else "ANNOUNCEMENT{}".format(ann_num)
label = "ANNOUNCEMENT" if ann_num == 1 else "ANNOUNCEMENT-{}".format(ann_num)
enabled = g.get("{}_ENABLED".format(prefix), False)
if ann_num in self._ann_tasks:
try:
if getattr(self._ann_tasks[ann_num], "running", False):
self._ann_tasks[ann_num].stop()
except Exception:
pass
del self._ann_tasks[ann_num]
logger.info("(VOICE-RELOAD) %s stopped", label)
if enabled:
mode = g.get("{}_MODE".format(prefix), "interval")
interval = 30.0 if mode == "hourly" else float(g.get("{}_INTERVAL".format(prefix), 60))
lc = self._start_looping_call(lambda an=ann_num: self.scheduled_announcement(an), interval, False)
self._ann_tasks[ann_num] = lc
logger.info(
"(VOICE-RELOAD) %s enabled - mode: %s, file: %s, TG: %s",
label, mode, g.get("{}_FILE".format(prefix)), g.get("{}_TG".format(prefix)),
)
for tts_num in range(1, 5):
prefix = "TTS_ANNOUNCEMENT{}".format(tts_num)
label = "TTS-{}".format(tts_num)
enabled = g.get("{}_ENABLED".format(prefix), False)
if tts_num in self._tts_tasks:
try:
if getattr(self._tts_tasks[tts_num], "running", False):
self._tts_tasks[tts_num].stop()
except Exception:
pass
del self._tts_tasks[tts_num]
logger.info("(VOICE-RELOAD) %s stopped", label)
if enabled:
mode = g.get("{}_MODE".format(prefix), "interval")
interval = 30.0 if mode == "hourly" else float(g.get("{}_INTERVAL".format(prefix), 60))
lc = self._start_looping_call(lambda tn=tts_num: self.scheduled_tts_announcement(tn), interval, False)
self._tts_tasks[tts_num] = lc
logger.info(
"(VOICE-RELOAD) %s enabled - mode: %s, file: %s, TG: %s",
label, mode, g.get("{}_FILE".format(prefix)), g.get("{}_TG".format(prefix)),
)
if config_file_path and self._config_file_mtime:
logger.info("(VOICE-RELOAD) config reload completed")

@ -0,0 +1,33 @@
# ADN DMR Peer Server - domain layer
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
from .entities import BridgeEntry, StreamState, SystemConfig
from .value_objects import DmrId, TgId, Slot, CallType, bytes_3, bytes_4, int_id, ID_MIN, ID_MAX, PEER_MAX
from .errors import DomainError, ConfigError, ACLError
from .result import Result, Success, Fail, is_fail, is_ok, unwrap_or
__all__ = [
"BridgeEntry",
"StreamState",
"SystemConfig",
"DmrId",
"TgId",
"Slot",
"CallType",
"bytes_3",
"bytes_4",
"int_id",
"ID_MIN",
"ID_MAX",
"PEER_MAX",
"DomainError",
"ConfigError",
"ACLError",
"Result",
"Success",
"Fail",
"is_fail",
"is_ok",
"unwrap_or",
]

@ -0,0 +1,66 @@
# ADN DMR Peer Server - domain entities
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Domain entities: bridge entry, stream state, system config (mirror legacy BRIDGES/STATUS/CONFIG)."""
from __future__ import annotations
from dataclasses import dataclass, field
from typing import Any
from .value_objects import Slot
@dataclass
class BridgeEntry:
"""One system/slot entry in a bridge (legacy BRIDGES[tgid][i])."""
SYSTEM: str
TS: Slot
TGID: bytes # 3 bytes
ACTIVE: bool
TIMEOUT: float | str # seconds or '' for STAT
TO_TYPE: str # 'ON', 'OFF', 'STAT', 'NONE'
OFF: list[bytes]
ON: list[bytes]
RESET: list[Any]
TIMER: float
@dataclass
class StreamState:
"""Per-slot stream state (legacy STATUS[slot] and per-stream state)."""
RX_START: float = 0.0
TX_START: float = 0.0
RX_SEQ: int = 0
RX_RFS: bytes = b"\x00"
TX_RFS: bytes = b"\x00"
RX_PEER: bytes = b"\x00"
TX_PEER: bytes = b"\x00"
RX_STREAM_ID: bytes = b"\x00"
TX_STREAM_ID: bytes = b"\x00"
RX_TGID: bytes = b"\x00\x00\x00"
TX_TGID: bytes = b"\x00\x00\x00"
RX_TIME: float = 0.0
TX_TIME: float = 0.0
RX_TYPE: int = 0
TX_TYPE: int = 0
RX_LC: bytes = b"\x00"
TX_H_LC: bytes = b"\x00"
TX_T_LC: bytes = b"\x00"
TX_EMB_LC: dict[int, bytes] = field(default_factory=dict)
lastSeq: bool = False
lastData: bool = False
packets: int = 0
crcs: set[Any] = field(default_factory=set)
@dataclass
class SystemConfig:
"""Minimal system config for domain (MODE, name, etc.). Full config lives in infrastructure."""
name: str
MODE: str # MASTER, PEER, OPENBRIDGE
ENABLED: bool = True

@ -0,0 +1,35 @@
# ADN DMR Peer Server - domain errors
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Domain errors for the ADN DMR Peer Server."""
class DomainError(Exception):
"""Base exception for domain/application."""
pass
class ConfigError(DomainError):
"""Configuration loading or validation error."""
pass
class ACLError(DomainError):
"""ACL parse or evaluation error."""
pass
class AliasError(DomainError):
"""Alias load or resolve error."""
pass
class ReportProtocolError(DomainError):
"""Report protocol decode or opcode error."""
pass

@ -0,0 +1,45 @@
# ADN DMR Peer Server - result type
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Result type for functional error handling (Success/Fail)."""
from __future__ import annotations
from dataclasses import dataclass
from typing import Generic, TypeVar
T = TypeVar("T")
E = TypeVar("E", bound=Exception)
@dataclass(frozen=True, slots=True)
class Success(Generic[T]):
"""Successful result wrapping a value."""
value: T
@dataclass(frozen=True, slots=True)
class Fail(Generic[E]):
"""Failed result wrapping an error."""
error: E
Result = Success[T] | Fail[E]
def is_fail(r: Result[T, E]) -> bool:
"""Return True if result is Fail."""
return isinstance(r, Fail)
def is_ok(r: Result[T, E]) -> bool:
"""Return True if result is Success."""
return isinstance(r, Success)
def unwrap_or(r: Result[T, E], default: T) -> T:
"""Return value if Success, else default."""
return r.value if isinstance(r, Success) else default

@ -0,0 +1,63 @@
# ADN DMR Peer Server - value objects
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Value objects: DMR IDs, slot, call type, etc. (mirror legacy const and types)."""
from __future__ import annotations
from dataclasses import dataclass
from typing import Literal
# Legacy ID_MIN=1, ID_MAX=16776415, PEER_MAX=4294967295
ID_MIN = 1
ID_MAX = 16776415
PEER_MAX = 4294967295
@dataclass(frozen=True, slots=True)
class DmrId:
"""DMR subscriber/peer/talkgroup numeric ID."""
value: int
def __int__(self) -> int:
return self.value
@dataclass(frozen=True, slots=True)
class TgId:
"""Talkgroup ID (same numeric space as DmrId for group calls)."""
value: int
def __int__(self) -> int:
return self.value
Slot = Literal[1, 2]
"""Timeslot 1 or 2."""
CallType = Literal["unit", "group"]
"""Unit (private) or group call."""
def bytes_3(i: int) -> bytes:
"""Encode int as 3-byte big-endian (legacy bytes_3)."""
return (int(i) & 0xFFFFFF).to_bytes(3, "big")
def bytes_4(i: int) -> bytes:
"""Encode int as 4-byte big-endian (legacy bytes_4, salt, peer_id in packets)."""
return (int(i) & 0xFFFFFFFF).to_bytes(4, "big")
def int_id(val: bytes | int) -> int:
"""Decode 3- or 4-byte big-endian to int (legacy int_id)."""
if isinstance(val, int):
return val
if len(val) >= 4:
return int.from_bytes(val[:4], "big")
if len(val) == 3:
return int.from_bytes(val, "big")
return 0

@ -0,0 +1,8 @@
# ADN DMR Peer Server - infrastructure layer
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
from .config_loader import YamlConfigLoader
from .logging_config import setup_logging
__all__ = ["YamlConfigLoader", "setup_logging"]

@ -0,0 +1,43 @@
# ADN DMR Peer Server - bridge router implementation
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""In-memory BRIDGES and ACL check (legacy acl_check)."""
from __future__ import annotations
from typing import Any
from ..application.ports import BridgeRouter
def _int_id(val: bytes | int) -> int:
if isinstance(val, int):
return val
if len(val) >= 4:
return int.from_bytes(val[:4], "big")
if len(val) == 3:
return int.from_bytes(val, "big")
return 0
class InMemoryBridgeRouter(BridgeRouter):
"""Holds BRIDGES dict; implements acl_check like legacy."""
def __init__(self) -> None:
self._bridges: dict[str, list[dict[str, Any]]] = {}
def get_bridges(self) -> dict[str, list[dict[str, Any]]]:
return self._bridges
def set_bridges(self, bridges: dict[str, list[dict[str, Any]]]) -> None:
self._bridges = bridges
def acl_check(self, id_bytes_or_int: bytes | int, acl: tuple[bool, list[tuple[int, int]]]) -> bool:
"""Legacy acl_check: (action, ranges). If id in any range return action else not action."""
action, ranges = acl
i = _int_id(id_bytes_or_int)
for lo, hi in ranges:
if lo <= i <= hi:
return action
return not action

@ -0,0 +1,117 @@
# ADN DMR Peer Server - YAML config loader
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Load config from YAML at project root; same semantic structure as legacy INI."""
from __future__ import annotations
import os
import sys
from pathlib import Path
from typing import Any
import yaml
from ..domain import ID_MAX, ID_MIN, PEER_MAX
from ..domain.errors import ConfigError
from . import logging_config
def acl_build(acl_str: str | None, max_id: int) -> tuple[bool, list[tuple[int, int]]]:
"""Build ACL from string e.g. 'DENY:1-5,3120101'. Returns (action, [(lo, hi), ...])."""
if not acl_str:
return (True, [(ID_MIN, max_id)])
parts = acl_str.split(":", 1)
if len(parts) != 2:
return (True, [(ID_MIN, max_id)])
action = parts[0].strip().upper() == "PERMIT"
acl: list[tuple[int, int]] = []
for entry in parts[1].strip().split(","):
entry = entry.strip()
if entry == "ALL":
acl.append((ID_MIN, max_id))
break
if "-" in entry:
start_s, end_s = entry.split("-", 1)
start, end = int(start_s.strip()), int(end_s.strip())
if not (ID_MIN <= start <= max_id) and not (ID_MIN <= end <= max_id):
raise ConfigError(f"ACL range out of bounds: {entry}")
acl.append((start, end))
else:
i = int(entry)
if not (ID_MIN <= i <= max_id):
raise ConfigError(f"ACL id out of bounds: {entry}")
acl.append((i, i))
return (action, acl)
def process_acls(config: dict[str, Any]) -> None:
"""Inject processed ACLs into CONFIG (GLOBAL and per SYSTEM). Mutates config."""
from ..domain import PEER_MAX
g = config.get("GLOBAL", {})
g["REG_ACL"] = acl_build(g.get("REG_ACL", "PERMIT:ALL"), PEER_MAX)
for key in ("SUB_ACL", "TG1_ACL", "TG2_ACL"):
acl_key = "TGID_TS1_ACL" if key == "TG1_ACL" else ("TGID_TS2_ACL" if key == "TG2_ACL" else key)
g[key] = acl_build(g.get(acl_key, g.get(key, "PERMIT:ALL")), ID_MAX)
for system_name, sys_cfg in config.get("SYSTEMS", {}).items():
if sys_cfg.get("MODE") == "MASTER":
sys_cfg["REG_ACL"] = acl_build(sys_cfg.get("REG_ACL", "PERMIT:ALL"), PEER_MAX)
for key in ("SUB_ACL", "TG1_ACL", "TG2_ACL"):
acl_key = "TGID_TS1_ACL" if key == "TG1_ACL" else ("TGID_TS2_ACL" if key == "TG2_ACL" else key)
sys_cfg[key] = acl_build(sys_cfg.get(acl_key, sys_cfg.get(key, "PERMIT:ALL")), ID_MAX)
class YamlConfigLoader:
"""Load config from YAML file; same semantics as legacy build_config."""
def __init__(self, project_root: str | Path = ".") -> None:
self._root = Path(project_root).resolve()
def load(self, path: str | None = None) -> dict[str, Any]:
"""Load config. path=None uses project root adn-server.yaml."""
if path is None:
path = str(self._root / "adn-server.yaml")
if not os.path.isfile(path):
raise ConfigError(f"Config file not found: {path}")
with open(path, "r", encoding="utf-8") as f:
data = yaml.safe_load(f)
if not data or not isinstance(data, dict):
raise ConfigError("Invalid YAML or empty config")
# Normalize to same top-level keys as legacy
config: dict[str, Any] = {
"GLOBAL": data.get("GLOBAL", {}),
"VOICE": data.get("VOICE", {}),
"REPORTS": data.get("REPORTS", {}),
"LOGGER": data.get("LOGGER", {}),
"ALIASES": data.get("ALIASES", {}),
"ALLSTAR": data.get("ALLSTAR", {}),
"SYSTEMS": data.get("SYSTEMS", {}),
}
# Ensure REPORT_CLIENTS is list
if "REPORT_CLIENTS" in config["REPORTS"] and isinstance(config["REPORTS"]["REPORT_CLIENTS"], str):
config["REPORTS"]["REPORT_CLIENTS"] = [
x.strip() for x in config["REPORTS"]["REPORT_CLIENTS"].split(",")
]
process_acls(config)
return config
def reload_voice_config(self, config: dict[str, Any], config_path: str | None = None) -> None:
"""Reload GLOBAL from main config file (voice, announcements, TTS live there). Updates config in place."""
if not config_path or not os.path.isfile(config_path):
return
try:
with open(config_path, "r", encoding="utf-8") as f:
data = yaml.safe_load(f)
except (OSError, yaml.YAMLError):
return
if not data or not isinstance(data, dict):
return
new_global = data.get("GLOBAL", {})
if new_global:
config["GLOBAL"].update(new_global)
process_acls(config)
new_voice = data.get("VOICE", {})
if new_voice:
config.setdefault("VOICE", {}).update(new_voice)

@ -0,0 +1,54 @@
# ADN DMR Peer Server - HBP protocol constants
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink (const.py). GPLv3.
"""Homebrew protocol opcodes and frame types (legacy const.py)."""
# DMR
DMR = b"DMR"
DMRD = b"DMRD"
DMRE = b"DMRE"
DMRF = b"DMRF"
DMRA = b"DMRA"
# Master/peer
RPTL = b"RPTL"
RPTPING = b"RPTPING"
RPTACK = b"RPTACK"
RPTCL = b"RPTCL"
RPTK = b"RPTK"
RPTC = b"RPTC"
RPTP = b"RPTP"
RPTA = b"RPTA"
RPTO = b"RPTO"
MSTCL = b"MSTCL"
MSTNAK = b"MSTNAK"
MSTPONG = b"MSTPONG"
MSTN = b"MSTN" # peer receives MSTNAK as MSTN (4-char command)
MSTP = b"MSTP" # peer receives MSTPONG as MSTP
MSTC = b"MSTC" # peer receives MSTCL as MSTC
# OpenBridge
EOBP = b"EOBP"
BC = b"BC"
BCKA = b"BCKA"
BCSQ = b"BCSQ"
BCST = b"BCST"
BCVE = b"BCVE"
# Proxy
PRIN = b"PRIN"
PRBL = b"PRBL"
# Frame types (bits)
HBPF_VOICE = 0x0
HBPF_VOICE_SYNC = 0x1
HBPF_DATA_SYNC = 0x2
HBPF_SLT_VHEAD = 0x1
HBPF_SLT_VTERM = 0x2
VER = 5
PROTO_VER = 5
# Legacy const.py: stream timeout (seconds) for contention
STREAM_TO = 0.36

@ -0,0 +1,45 @@
# ADN DMR Peer Server - logging setup
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Configure logging (same format/handlers as legacy log.config_logging)."""
from __future__ import annotations
import logging
import sys
from functools import partial, partialmethod
from typing import Any
def setup_logging(log_config: dict[str, Any]) -> logging.Logger:
"""Configure logging from CONFIG['LOGGER']. Returns root logger."""
level = getattr(logging, (log_config.get("LOG_LEVEL", "INFO")).upper(), logging.INFO)
log_file = log_config.get("LOG_FILE", "/dev/null")
handlers_cfg = log_config.get("LOG_HANDLERS", "console-timed").strip().split(",")
log_name = log_config.get("LOG_NAME", "ADN")
logging.TRACE = 5
logging.addLevelName(logging.TRACE, "TRACE")
logging.Logger.trace = partialmethod(logging.Logger.log, logging.TRACE)
logging.trace = partial(logging.log, logging.TRACE)
handlers: list[logging.Handler] = []
if "console-timed" in handlers_cfg or "console" in handlers_cfg:
h = logging.StreamHandler()
h.setFormatter(logging.Formatter("%(levelname)s %(asctime)s %(message)s"))
handlers.append(h)
if ("file-timed" in handlers_cfg or "file" in handlers_cfg) and log_file and log_file != "/dev/null":
try:
h = logging.FileHandler(log_file, encoding="utf-8")
h.setFormatter(logging.Formatter("%(levelname)s %(asctime)s %(message)s"))
handlers.append(h)
except OSError as e:
# Fallback: warn to stderr if file cannot be opened (e.g. permission, missing dir)
sys.stderr.write("(LOGGER) Could not open log file %s: %s\n" % (log_file, e))
# force=True (Python 3.8+) so our handlers replace any already set by other libs (e.g. Twisted)
logging.basicConfig(level=level, handlers=handlers or [logging.NullHandler()], force=True)
logger = logging.getLogger(log_name)
logger.setLevel(level)
return logger

@ -0,0 +1,8 @@
# ADN DMR Peer Server - persistence
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
from .sub_map_store import PickleSubMapStore
from .keys_store import JsonKeysStore
__all__ = ["PickleSubMapStore", "JsonKeysStore"]

@ -0,0 +1,201 @@
# ADN DMR Peer Server - alias loader
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Load alias dicts (peer_ids, subscriber_ids, talkgroup_ids, etc.). Legacy mk_aliases."""
from __future__ import annotations
import csv
import hashlib
import json
import logging
import ssl
import time
from pathlib import Path
from typing import Any
from urllib.request import urlopen
from ...application.ports import AliasLoader
logger = logging.getLogger(__name__)
def try_download(path: Path, file_name: str, url: str, stale_sec: float) -> str:
"""Legacy try_download: download file from url if missing or older than stale_sec. Returns result message."""
if not url:
return f"ID ALIAS MAPPER: '{file_name}' URL empty, not downloaded"
full = path / file_name
now = time.time()
file_exists = full.is_file()
if file_exists:
file_old = (full.stat().st_mtime + stale_sec) < now
else:
file_old = True
if not file_old and file_exists:
return f"ID ALIAS MAPPER: '{file_name}' is current, not downloaded"
try:
ctx = ssl.create_default_context()
ctx.check_hostname = False
ctx.verify_mode = ssl.CERT_NONE
with urlopen(url, context=ctx, timeout=30) as response:
data = response.read()
except OSError as e:
return f"ID ALIAS MAPPER: '{file_name}' could not be downloaded due to an IOError: {e}"
if not data or data == b"{}":
return f"ID ALIAS MAPPER: '{file_name}' file not written because downloaded data is empty"
try:
full.parent.mkdir(parents=True, exist_ok=True)
full.write_bytes(data)
except OSError as e:
return f"ID ALIAS mapper '{file_name}' file could not be written: {e}"
return f"ID ALIAS MAPPER: '{file_name}' successfully downloaded"
def _blake2bsum(file_path: Path) -> str:
"""Blake2b hex digest of file (legacy blake2bsum)."""
h = hashlib.blake2b()
with open(file_path, "rb") as f:
for chunk in iter(lambda: f.read(4096), b""):
h.update(chunk)
return h.hexdigest()
class DefaultAliasLoader(AliasLoader):
"""Load aliases from JSON files and optional downloads. Legacy mk_aliases."""
def load_aliases(
self,
config: dict[str, Any],
) -> tuple[
dict[int, str],
dict[int, str],
dict[int, str],
dict[int, str],
dict[str, str],
dict[str, str],
]:
"""Build alias dicts. Same order as legacy mk_aliases."""
aliases = config.get("ALIASES", {})
path = Path(aliases.get("PATH", "./data/")).resolve()
stale_sec = float(aliases.get("STALE_TIME", aliases.get("STALE_DAYS", 1) * 86400))
if aliases.get("TRY_DOWNLOAD"):
if aliases.get("CHECKSUM_FILE") and aliases.get("CHECKSUM_URL"):
result = try_download(path, aliases["CHECKSUM_FILE"], aliases.get("CHECKSUM_URL", ""), stale_sec)
logger.info("(ALIAS) %s", result)
for key, url_key in [
("PEER_FILE", "PEER_URL"),
("SUBSCRIBER_FILE", "SUBSCRIBER_URL"),
("TGID_FILE", "TGID_URL"),
("SERVER_ID_FILE", "SERVER_ID_URL"),
]:
url = aliases.get(url_key)
if url and aliases.get(key):
result = try_download(path, aliases[key], url, stale_sec)
logger.info("(ALIAS) %s", result)
checksums = self._load_checksums(path, aliases.get("CHECKSUM_FILE"))
peer_file = aliases.get("PEER_FILE", "peer_ids.json")
sub_file = aliases.get("SUBSCRIBER_FILE", "subscriber_ids.json")
tgid_file = aliases.get("TGID_FILE", "talkgroup_ids.json")
server_file = aliases.get("SERVER_ID_FILE", "server_ids.tsv")
peer_ids = self._load_id_json_verified(path / peer_file, checksums.get("peer_ids"), "peer_ids")
subscriber_ids = self._load_id_json_verified(path / sub_file, checksums.get("subscriber_ids"), "subscriber_ids")
talkgroup_ids = self._load_id_json_verified(path / tgid_file, checksums.get("talkgroup_ids"), "talkgroup_ids")
local_subscriber_ids = self._load_id_json(
path / aliases.get("LOCAL_SUBSCRIBER_FILE", "subscriber_ids.json")
)
server_ids = self._load_server_tsv_verified(path, server_file, checksums.get("server_ids"))
if server_ids:
logger.info("(ALIAS) ID ALIAS MAPPER: server_ids dictionary is available")
return (peer_ids, subscriber_ids, talkgroup_ids, local_subscriber_ids, server_ids, checksums)
def _load_checksums(self, path: Path, file_name: str | None) -> dict[str, str]:
"""Load checksum JSON (legacy load_json of CHECKSUM_FILE). Keys e.g. peer_ids, subscriber_ids, talkgroup_ids, server_ids."""
if not file_name:
return {}
full = path / file_name
if not full.is_file():
return {}
try:
with open(full, "r", encoding="utf-8") as f:
data = json.load(f)
except (json.JSONDecodeError, OSError) as e:
logger.error("(ALIAS) ID ALIAS MAPPER: Cannot load checksums: %s", e)
return {}
return dict(data) if isinstance(data, dict) else {}
def _load_server_tsv(self, path: Path, file_name: str) -> dict[str, str]:
"""Legacy mk_server_dict: TSV with 'OPB Net ID' -> 'Country'."""
full = path / file_name
if not full.is_file():
return {}
try:
with open(full, "r", newline="", encoding="utf-8") as f:
reader = csv.DictReader(f, dialect="excel-tab")
out: dict[str, str] = {}
for row in reader:
net_id = row.get("OPB Net ID", "").strip()
country = row.get("Country", "").strip()
if net_id:
out[net_id] = country
return out
except Exception as err:
logger.warning("(ALIAS) ID ALIAS MAPPER: %s could not be read: %s", file_name, err)
return {}
def _load_id_json_verified(
self, file_path: Path, expected_checksum: str | None, name: str
) -> dict[int, str]:
"""Load ID JSON; if expected_checksum given, verify blake2b first (legacy)."""
if not file_path.is_file():
return {}
if expected_checksum:
try:
if _blake2bsum(file_path) != expected_checksum:
logger.error("(ALIAS) ID ALIAS MAPPER: problem with blake2bsum of %s file. not updating.", name)
return {}
except Exception as e:
logger.error("(ALIAS) ID ALIAS MAPPER: problem with blake2bsum of %s file: %s", name, e)
return {}
return self._load_id_json(file_path)
def _load_id_json(self, file_path: Path) -> dict[int, str]:
"""Load JSON with 'id' -> 'callsign' structure; return {int(id): callsign}."""
if not file_path.is_file():
return {}
try:
with open(file_path, "r", encoding="utf-8") as f:
data = json.load(f)
except (json.JSONDecodeError, OSError):
return {}
if not isinstance(data, dict) or "count" in data:
if isinstance(data, dict) and "count" in data:
del data["count"]
out: dict[int, str] = {}
if isinstance(data, dict):
for key, val in data.items():
if isinstance(val, list):
for record in val:
if isinstance(record, dict) and "id" in record and "callsign" in record:
try:
out[int(record["id"])] = str(record["callsign"])
except (ValueError, TypeError):
pass
return out
def _load_server_tsv_verified(
self, path: Path, file_name: str, expected_checksum: str | None
) -> dict[str, str]:
"""Load server_ids TSV; if expected_checksum given, verify blake2b first (legacy)."""
full = path / file_name
if not full.is_file():
return {}
if expected_checksum:
try:
if _blake2bsum(full) != expected_checksum:
logger.error("(ALIAS) ID ALIAS MAPPER: problem with blake2bsum of server_ids file: not updating.")
return {}
except Exception as e:
logger.error("(ALIAS) ID ALIAS MAPPER: problem with blake2bsum of server_ids file: %s", e)
return {}
return self._load_server_tsv(path, file_name)

@ -0,0 +1,35 @@
# ADN DMR Peer Server - keys JSON store
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Load/save keys from JSON (legacy utils.load_json/save_json)."""
from __future__ import annotations
import json
from pathlib import Path
from typing import Any
from ...application.ports import KeysStore
class JsonKeysStore(KeysStore):
"""Persist keys as JSON."""
def load(self, path: str) -> dict[str, Any]:
"""Load keys from JSON file."""
p = Path(path)
if not p.is_file():
return {}
try:
with open(p, "r", encoding="utf-8") as f:
return json.load(f)
except (json.JSONDecodeError, OSError):
return {}
def save(self, path: str, keys: dict[str, Any]) -> None:
"""Save keys to JSON file."""
p = Path(path)
p.parent.mkdir(parents=True, exist_ok=True)
with open(p, "w", encoding="utf-8") as f:
json.dump(keys, f, indent=2)

@ -0,0 +1,35 @@
# ADN DMR Peer Server - SUB_MAP persistence
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Load/save SUB_MAP (pickle): bytes_3(peer) -> (callsign, slot, time)."""
from __future__ import annotations
import pickle
from pathlib import Path
from typing import Any
from ...application.ports import SubMapStore
class PickleSubMapStore(SubMapStore):
"""Persist SUB_MAP as pickle (legacy compatible)."""
def load(self, path: str) -> dict[bytes, tuple[str, int, float]]:
"""Load SUB_MAP from pickle file."""
p = Path(path)
if not p.is_file():
return {}
try:
with open(p, "rb") as f:
return pickle.load(f)
except (pickle.PickleError, OSError):
return {}
def save(self, path: str, sub_map: dict[bytes, tuple[str, int, float]]) -> None:
"""Save SUB_MAP to pickle file."""
p = Path(path)
p.parent.mkdir(parents=True, exist_ok=True)
with open(p, "wb") as f:
pickle.dump(sub_map, f)

@ -0,0 +1,7 @@
# ADN DMR Peer Server - security
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
from .password_download import DefaultSecurityDownloader, StubSecurityDownloader
__all__ = ["DefaultSecurityDownloader", "StubSecurityDownloader"]

@ -0,0 +1,70 @@
# ADN DMR Peer Server - password encryption (legacy password_crypto)
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server. GPLv3.
"""Fernet-based encrypt/decrypt for stored passwords. Legacy password_crypto."""
from __future__ import annotations
import os
from typing import Optional
try:
from cryptography.fernet import Fernet
except ImportError:
Fernet = None # type: ignore[misc, assignment]
def get_or_create_key(key_path: str) -> bytes:
"""Load encryption key from file or create and save one."""
if Fernet is None:
raise RuntimeError("cryptography package required for password_crypto")
if os.path.exists(key_path):
with open(key_path, "rb") as f:
return f.read()
key = Fernet.generate_key()
os.makedirs(os.path.dirname(key_path), exist_ok=True)
with open(key_path, "wb") as f:
f.write(key)
return key
def get_fernet(key_path: str = "config/encryption_key.secret") -> "Fernet":
"""Return Fernet instance using key at key_path."""
if Fernet is None:
raise RuntimeError("cryptography package required for password_crypto")
key = get_or_create_key(key_path)
return Fernet(key)
def decrypt_password(encrypted_password: Optional[str], key_path: str = "config/encryption_key.secret") -> Optional[str]:
"""Decrypt a stored password; return as-is if empty or on failure."""
if not encrypted_password:
return encrypted_password
try:
fernet = get_fernet(key_path)
decrypted = fernet.decrypt(encrypted_password.encode("utf-8"))
return decrypted.decode("utf-8")
except Exception:
return encrypted_password
def encrypt_password(password: Optional[str], key_path: str = "config/encryption_key.secret") -> Optional[str]:
"""Encrypt a password for storage."""
if not password:
return password
fernet = get_fernet(key_path)
encrypted = fernet.encrypt(password.encode("utf-8"))
return encrypted.decode("utf-8")
def is_encrypted(value: Optional[str], key_path: str = "config/encryption_key.secret") -> bool:
"""Return True if value looks like Fernet-encrypted data."""
if not value:
return False
try:
fernet = get_fernet(key_path)
fernet.decrypt(value.encode("utf-8"))
return True
except Exception:
return False

@ -0,0 +1,228 @@
# ADN DMR Peer Server - security downloads (legacy security_downloader)
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server. GPLv3.
"""Central security server: download encryption key and user passwords. Legacy init_security_downloads, periodic_password_download."""
from __future__ import annotations
import logging
import os
import shutil
import socket
import tempfile
import time
from typing import Any
from urllib.parse import quote
from urllib.request import Request, urlopen
from urllib.error import HTTPError, URLError
from ...application.ports import SecurityDownloader
logger = logging.getLogger(__name__)
DOWNLOAD_INTERVAL_PASSWORDS = 300
_last_passwords_download = 0.0
_last_passwords_size = 0
_last_passwords_content: bytes | None = None
def _resolve_hostname(hostname: str, timeout: int = 10) -> str | None:
try:
old_timeout = socket.getdefaulttimeout()
socket.setdefaulttimeout(timeout)
ip = socket.gethostbyname(hostname)
socket.setdefaulttimeout(old_timeout)
logger.debug("(SECURITY) Resolved %s to %s", hostname, ip)
return ip
except socket.gaierror as e:
logger.error("(SECURITY) DNS resolution failed for %s: %s", hostname, e)
return None
except Exception as e:
logger.error("(SECURITY) Unexpected error resolving %s: %s", hostname, e)
return None
def _build_download_url(config: dict[str, Any], filename: str) -> tuple[str | None, str | None]:
g = config.get("GLOBAL", {})
url_security = (g.get("URL_SECURITY") or "").strip()
port_security = (g.get("PORT_SECURITY") or "").strip()
pass_security = (g.get("PASS_SECURITY") or "").strip()
if not url_security or not port_security or not pass_security:
return None, None
try:
socket.inet_aton(url_security)
host = url_security
except OSError:
host = _resolve_hostname(url_security)
if not host:
logger.error("(SECURITY) Could not resolve hostname: %s", url_security)
return None, None
url = f"http://{host}:{port_security}/descargar?pass={quote(pass_security, safe='')}&file={filename}"
return url, url_security
def _download_file_safely(
url: str, dest_path: str, timeout: int = 60
) -> bool:
try:
fd, temp_path = tempfile.mkstemp()
os.close(fd)
try:
logger.debug("(SECURITY) Attempting download from: %s", url)
req = Request(url)
req.add_header("User-Agent", "ADN-Systems-DMR/1.0")
with urlopen(req, timeout=timeout) as response:
content = response.read()
if len(content) == 0:
logger.warning("(SECURITY) Downloaded file is empty, keeping existing: %s", dest_path)
os.unlink(temp_path)
return False
with open(temp_path, "wb") as f:
f.write(content)
os.makedirs(os.path.dirname(dest_path), exist_ok=True)
shutil.move(temp_path, dest_path)
logger.info("(SECURITY) Successfully downloaded: %s (%d bytes)", dest_path, len(content))
return True
except HTTPError as e:
logger.error("(SECURITY) HTTP error downloading %s: %s (Code: %d)", dest_path, e, e.code)
if os.path.exists(temp_path):
os.unlink(temp_path)
return False
except URLError as e:
logger.error("(SECURITY) URL error downloading %s: %s", dest_path, e.reason)
if os.path.exists(temp_path):
os.unlink(temp_path)
return False
except socket.timeout:
logger.error("(SECURITY) Timeout downloading %s", dest_path)
if os.path.exists(temp_path):
os.unlink(temp_path)
return False
except Exception as e:
logger.error("(SECURITY) Unexpected error downloading %s: %s", dest_path, e)
return False
def _download_encryption_key(config: dict[str, Any], config_dir: str) -> bool:
g = config.get("GLOBAL", {})
hash_encrypt = (g.get("HASH_ENCRYPT") or "encryption_key.secret").strip()
dest_path = os.path.join(config_dir, hash_encrypt)
url, _ = _build_download_url(config, hash_encrypt)
if not url:
logger.debug("(SECURITY) Security server not configured, skipping encryption key download")
return False
logger.info("(SECURITY) Downloading encryption key from central server...")
return _download_file_safely(url, dest_path)
def _download_user_passwords(
config: dict[str, Any], data_dir: str, force: bool = False
) -> bool:
global _last_passwords_download, _last_passwords_size, _last_passwords_content
now = time.time()
if not force and (now - _last_passwords_download) < DOWNLOAD_INTERVAL_PASSWORDS:
return False
g = config.get("GLOBAL", {})
users_pass = (g.get("USERS_PASS") or "user_passwords.json").strip()
dest_path = os.path.join(data_dir, users_pass)
url, _ = _build_download_url(config, users_pass)
if not url:
logger.debug("(SECURITY) Security server not configured, skipping passwords download")
return False
try:
logger.debug("(SECURITY) Downloading passwords from: %s", url)
req = Request(url)
req.add_header("User-Agent", "ADN-Systems-DMR/1.0")
with urlopen(req, timeout=60) as response:
new_content = response.read()
new_size = len(new_content)
if new_size == 0:
logger.warning("(SECURITY) Downloaded passwords file is empty, keeping existing")
_last_passwords_download = now
return False
if _last_passwords_content is not None and new_content == _last_passwords_content:
logger.debug("(SECURITY) Passwords file unchanged, no update needed")
_last_passwords_download = now
return False
fd, temp_path = tempfile.mkstemp()
os.close(fd)
with open(temp_path, "wb") as f:
f.write(new_content)
os.makedirs(os.path.dirname(dest_path), exist_ok=True)
shutil.move(temp_path, dest_path)
_last_passwords_content = new_content
_last_passwords_size = new_size
_last_passwords_download = now
logger.info("(SECURITY) Successfully updated passwords file: %s (%d bytes)", dest_path, new_size)
return True
except HTTPError as e:
logger.error("(SECURITY) HTTP error downloading passwords: %s (Code: %d)", e, e.code)
_last_passwords_download = now
return False
except URLError as e:
logger.error("(SECURITY) URL error downloading passwords: %s", e.reason)
_last_passwords_download = now
return False
except socket.timeout:
logger.error("(SECURITY) Timeout downloading passwords")
_last_passwords_download = now
return False
except Exception as e:
logger.error("(SECURITY) Unexpected error downloading passwords: %s", e)
_last_passwords_download = now
return False
class DefaultSecurityDownloader(SecurityDownloader):
"""Legacy security_downloader: init_security_downloads + periodic_password_download."""
def __init__(self, project_root: str) -> None:
self._project_root = project_root
def _config_dir(self, config: dict[str, Any]) -> str:
return os.path.join(self._project_root, (config.get("GLOBAL", {}).get("CONFIG_PATH") or "config"))
def _data_dir(self, config: dict[str, Any]) -> str:
path = (config.get("ALIASES", {}).get("PATH") or "data").rstrip("/")
return os.path.join(self._project_root, path)
def init_downloads(self, config: dict[str, Any]) -> None:
"""One-time init: resolve hostname, download encryption key and passwords (force)."""
url_security = (config.get("GLOBAL", {}).get("URL_SECURITY") or "").strip()
if not url_security:
logger.info("(SECURITY) Central security server not configured")
return
port_security = (config.get("GLOBAL", {}).get("PORT_SECURITY") or "").strip()
logger.info("(SECURITY) Initializing centralized security downloads...")
logger.info("(SECURITY) Security server: %s:%s", url_security, port_security)
try:
socket.inet_aton(url_security)
logger.info("(SECURITY) Using IP address: %s", url_security)
except OSError:
resolved = _resolve_hostname(url_security)
if resolved:
logger.info("(SECURITY) Resolved hostname %s to IP: %s", url_security, resolved)
else:
logger.error("(SECURITY) Failed to resolve hostname: %s", url_security)
return
config_dir = self._config_dir(config)
data_dir = self._data_dir(config)
_download_encryption_key(config, config_dir)
_download_user_passwords(config, data_dir, force=True)
def periodic_download(self, config: dict[str, Any]) -> None:
"""Periodic password file download (every 5 min)."""
data_dir = self._data_dir(config)
_download_user_passwords(config, data_dir)
class StubSecurityDownloader(SecurityDownloader):
"""Stub: init and periodic_download do nothing."""
def init_downloads(self, config: dict[str, Any]) -> None:
pass
def periodic_download(self, config: dict[str, Any]) -> None:
pass

@ -0,0 +1,83 @@
# ADN DMR Peer Server - load and decrypt user passwords (legacy hblink load_user_passwords, get_user_password)
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Load user_passwords.json, decrypt with password_crypto; expose get_user_password(radio_id) for auth."""
from __future__ import annotations
import json
import logging
import os
import time
from typing import Any
logger = logging.getLogger(__name__)
USER_PASSWORDS_RELOAD_INTERVAL = 10.0
_last_load = 0.0
class UserPasswordsLoader:
"""Load and cache decrypted user passwords from GLOBAL.USERS_PASS (JSON with 'passwords' dict)."""
def __init__(self, project_root: str) -> None:
self._project_root = project_root
self._passwords: dict[str, str] = {}
self._config: dict[str, Any] = {}
self._config_dir = os.path.join(project_root, "config")
self._key_path = os.path.join(self._config_dir, "encryption_key.secret")
def load(self, config: dict[str, Any]) -> dict[str, str]:
"""Load user_passwords.json from data dir, decrypt each; return passwords dict. Legacy load_user_passwords."""
global _last_load
self._config = config
now = time.time()
if now - _last_load < USER_PASSWORDS_RELOAD_INTERVAL and self._passwords:
return self._passwords
data_dir = os.path.join(
self._project_root,
(config.get("ALIASES", {}).get("PATH") or "data").rstrip("/"),
)
users_pass = (config.get("GLOBAL", {}).get("USERS_PASS") or "user_passwords.json").strip()
path = os.path.join(data_dir, users_pass)
key_path = os.path.join(
self._project_root,
(config.get("GLOBAL", {}).get("CONFIG_PATH") or "config").rstrip("/"),
)
hash_encrypt = (config.get("GLOBAL", {}).get("HASH_ENCRYPT") or "encryption_key.secret").strip()
self._key_path = os.path.join(key_path, hash_encrypt)
self._passwords = {}
if not os.path.exists(path):
_last_load = now
return self._passwords
try:
with open(path, "r", encoding="utf-8") as f:
data = json.load(f)
encrypted = data.get("passwords", {})
from .password_crypto import decrypt_password
for radio_id, pwd in encrypted.items():
self._passwords[str(radio_id)] = decrypt_password(pwd, self._key_path) or ""
logger.debug("(AUTH) Loaded %d individual passwords from %s", len(self._passwords), path)
except (FileNotFoundError, json.JSONDecodeError, Exception) as e:
logger.warning("(AUTH) Could not load user passwords: %s", e)
_last_load = now
return self._passwords
def get_user_password(self, radio_id: int) -> bytes | None:
"""Return password for radio_id (for login auth); 7-char prefix match like legacy. Legacy get_user_password."""
if time.time() - _last_load >= USER_PASSWORDS_RELOAD_INTERVAL and self._config:
self.load(self._config)
radio_id_str = str(radio_id)
if not radio_id_str.isdigit():
return None
if radio_id_str in self._passwords:
pwd = self._passwords[radio_id_str]
return pwd.encode("utf-8") if isinstance(pwd, str) else pwd
if len(radio_id_str) == 9:
base_id = radio_id_str[:7]
if base_id in self._passwords:
logger.debug("(AUTH) Radio ID %s using base ID %s password", radio_id_str, base_id)
pwd = self._passwords[base_id]
return pwd.encode("utf-8") if isinstance(pwd, str) else pwd
return None

@ -0,0 +1,8 @@
# ADN DMR Peer Server - Twisted adapters
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
from .report_server import ReportServerFactory, REPORT_OPCODES
from .udp_hbp import HBPProtocolFactory
__all__ = ["ReportServerFactory", "REPORT_OPCODES", "HBPProtocolFactory"]

@ -0,0 +1,115 @@
# ADN DMR Peer Server - TCP report server
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Report server: CONFIG_SND, BRIDGE_SND, BRDG_EVENT (legacy reportFactory, bridgeReportFactory)."""
from __future__ import annotations
import logging
import pickle
from typing import Any
from twisted.internet.protocol import Factory
from twisted.protocols.basic import NetstringReceiver
logger = logging.getLogger(__name__)
# Same opcodes as legacy reporting_const
REPORT_OPCODES = {
"CONFIG_REQ": b"\x00",
"CONFIG_SND": b"\x01",
"BRIDGE_REQ": b"\x02",
"BRIDGE_SND": b"\x03",
"CONFIG_UPD": b"\x04",
"BRIDGE_UPD": b"\x05",
"LINK_EVENT": b"\x06",
"BRDG_EVENT": b"\x07",
}
class ReportProtocol(NetstringReceiver):
"""Single report client connection."""
def __init__(self, factory: "ReportServerFactory") -> None:
self._factory = factory
def connectionMade(self) -> None:
self._factory.clients.append(self)
peer = self.transport.getPeer() if self.transport else None
addr = f"{peer.host}:{peer.port}" if peer else "?"
logger.info("(REPORT) Client connected from %s (%s client(s))", addr, len(self._factory.clients))
# Send CONFIG_SND and BRIDGE_SND immediately so monitor gets systems/bridges without waiting for loop
self._factory._send_config_to(self)
self._factory._send_bridge_to(self)
def connectionLost(self, reason: Any = None) -> None:
if self in self._factory.clients:
self._factory.clients.remove(self)
logger.info("(REPORT) Client disconnected (%s client(s))", len(self._factory.clients))
def stringReceived(self, data: bytes) -> None:
if data[:1] == REPORT_OPCODES["CONFIG_REQ"]:
self._factory._send_config_to(self)
elif data[:1] == REPORT_OPCODES["BRIDGE_REQ"]:
self._factory._send_bridge_to(self)
else:
pass # unknown opcode
class ReportServerFactory(Factory):
"""Factory for report protocol; holds config and bridges and sends to clients."""
def __init__(self, config: dict[str, Any]) -> None:
self._config = config
self.clients: list[ReportProtocol] = []
self._systems: dict[str, Any] = {}
self._bridges: dict[str, Any] = {}
def buildProtocol(self, addr: Any) -> ReportProtocol | None:
allowed = self._config.get("REPORTS", {}).get("REPORT_CLIENTS", ["127.0.0.1"])
if isinstance(allowed, str):
allowed = [x.strip() for x in allowed.split(",")]
if "*" in allowed or (hasattr(addr, "host") and addr.host in allowed):
return ReportProtocol(self)
return None
def set_systems(self, systems: dict[str, Any]) -> None:
self._systems = systems
def set_bridges(self, bridges: dict[str, Any]) -> None:
self._bridges = bridges
def _send_config_to(self, client: ReportProtocol) -> None:
"""Send CONFIG_SND to a single client (e.g. on connect or CONFIG_REQ)."""
payload = pickle.dumps(self._systems, protocol=2)
msg = REPORT_OPCODES["CONFIG_SND"] + payload
client.sendString(msg)
logger.debug("(REPORT) Sent CONFIG_SND to client (%d systems)", len(self._systems))
def _send_bridge_to(self, client: ReportProtocol) -> None:
"""Send BRIDGE_SND to a single client (e.g. on connect or BRIDGE_REQ)."""
payload = pickle.dumps(self._bridges, protocol=2)
msg = REPORT_OPCODES["BRIDGE_SND"] + payload
client.sendString(msg)
logger.debug("(REPORT) Sent BRIDGE_SND to client (%d bridges)", len(self._bridges))
def send_config(self) -> None:
"""Send CONFIG_SND (pickle SYSTEMS) to all clients."""
n = len(self.clients)
logger.debug("(REPORT) Sending CONFIG_SND to %s client(s)", n)
for c in self.clients:
self._send_config_to(c)
def send_bridge(self) -> None:
"""Send BRIDGE_SND (pickle BRIDGES) to all clients."""
n = len(self.clients)
logger.debug("(REPORT) Sending BRIDGE_SND to %s client(s)", n)
for c in self.clients:
self._send_bridge_to(c)
def send_bridge_event(self, event: str) -> None:
"""Send BRDG_EVENT."""
msg = REPORT_OPCODES["BRDG_EVENT"] + event.encode("utf-8", errors="ignore")
for c in self.clients:
c.sendString(msg)

File diff suppressed because it is too large Load Diff

@ -0,0 +1,8 @@
# ADN DMR Peer Server - voice infrastructure
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
from .ambe_reader import DefaultVoiceProvider, ReadAMBE, StubVoiceProvider
from .pkt_gen import pkt_gen
__all__ = ["DefaultVoiceProvider", "ReadAMBE", "StubVoiceProvider", "pkt_gen"]

@ -0,0 +1,184 @@
# ADN DMR Peer Server - AMBE reader and default voice provider
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink (read_ambe.py, mk_voice). GPLv3.
"""AMBE word loading (readAMBE) and VoiceProvider that uses it + pkt_gen."""
from __future__ import annotations
import glob
import logging
import os
from itertools import islice
from pathlib import Path
from typing import Any, Iterator
from bitarray import bitarray
from ...application.ports import VoiceProvider
from .pkt_gen import pkt_gen as _pkt_gen
from .voice_map import VOICE_MAP
logger = logging.getLogger(__name__)
_AMBE_LENGTH = 9
# Silence burst pair (legacy default when no .ambe available)
SILENCE_PAIR = [
bitarray("101011000000101010100000010000000000001000000000000000000000010001000000010000000000100000000000100000000000"),
bitarray("001010110000001010101000000100000000000010000000000000000000000100010000000100000000001000000000001000000000"),
]
def _make_bursts(data: bitarray):
"""Yield 108-bit bursts from bitarray (legacy _make_bursts)."""
it = iter(data)
n = len(data)
for i in range(0, n, 108):
chunk = bitarray([k for k in islice(it, 108)])
if len(chunk) < 108:
chunk.extend([False] * (108 - len(chunk)))
yield chunk
class ReadAMBE:
"""Legacy readAMBE: load AMBE words by language from path (dir of .ambe or .indx+.ambe)."""
def __init__(self, lang: str, path: str | Path) -> None:
self.langcsv = lang
self.langs = [s.strip() for s in lang.split(",") if s.strip()]
self.path = Path(path) if not isinstance(path, Path) else path
def readfiles(self) -> dict[str, dict[str, list[list[bitarray]]]]:
"""Load words per language. Returns {lang: {voice_name: [[b0,b1], ...]}}."""
result: dict[str, dict[str, list[list[bitarray]]]] = {}
for _lang in self.langs:
_prefix = self.path / _lang
_wordBADict: dict[str, list[list[bitarray]]] = {}
if _prefix.is_dir():
for ambe_path in glob.glob(str(_prefix / "*.ambe")):
basename = os.path.basename(ambe_path)
voice_name, _ = basename.split(".", 1)
try:
with open(ambe_path, "rb") as f:
_wordBitarray = bitarray(endian="big")
_wordBitarray.frombytes(f.read())
except OSError:
continue
_wordBADict[voice_name] = self._pairs_from_bitarray(_wordBitarray)
_wordBADict.setdefault("silence", [SILENCE_PAIR])
result[_lang] = _wordBADict
else:
index_path = Path(str(_prefix) + ".indx")
ambe_path = Path(str(_prefix) + ".ambe")
if not index_path.is_file() or not ambe_path.is_file():
result[_lang] = {"silence": [SILENCE_PAIR]}
continue
indexDict: dict[str, list[int]] = {}
with open(index_path, "r", encoding="utf-8") as index:
for line in index:
parts = line.split()
if len(parts) >= 3:
voice_name, start, length = parts[0], int(parts[1]), int(parts[2])
indexDict[voice_name] = [start * _AMBE_LENGTH, length * _AMBE_LENGTH]
try:
with open(ambe_path, "rb") as ambe:
for voice_name, (start, length) in indexDict.items():
ambe.seek(start)
_wordBitarray = bitarray(endian="big")
_wordBitarray.frombytes(ambe.read(length))
_wordBADict[voice_name] = self._pairs_from_bitarray(_wordBitarray)
except OSError:
result[_lang] = {"silence": [SILENCE_PAIR]}
continue
_wordBADict.setdefault("silence", [SILENCE_PAIR])
result[_lang] = _wordBADict
return result
def _pairs_from_bitarray(self, _wordBitarray: bitarray) -> list[list[bitarray]]:
"""Convert bitarray to list of [burst0, burst1] pairs (108 bits each)."""
_wordBA: list[list[bitarray]] = []
pairs = 1
_lastburst: bitarray | None = None
for _burst in _make_bursts(_wordBitarray):
if pairs == 2 and _lastburst is not None:
_wordBA.append([_lastburst, _burst])
_lastburst = None
pairs = 1
else:
pairs = 2
_lastburst = _burst
return _wordBA
def readSingleFile(self, filename: str) -> list[list[bitarray]]:
"""Read one .ambe file; return list of [b0, b1] pairs (legacy readSingleFile)."""
full = self.path / filename if not filename.startswith("/") else Path(filename)
if not full.is_file():
full = self.path / filename
if not full.is_file():
return []
try:
with open(full, "rb") as ambe:
_wordBitarray = bitarray(endian="big")
_wordBitarray.frombytes(ambe.read())
except OSError:
return []
return self._pairs_from_bitarray(_wordBitarray)
class DefaultVoiceProvider(VoiceProvider):
"""Voice provider using ReadAMBE and pkt_gen (legacy readAMBE + mk_voice.pkt_gen)."""
def __init__(self) -> None:
self._words: dict[str, dict[str, Any]] = {}
def get_ambe_words(self, languages: str, audio_path: str) -> dict[str, dict[str, Any]]:
"""Load AMBE words (readAMBE.readfiles). Apply i18n voiceMap per lang. Cached per (languages, audio_path)."""
key = f"{languages}:{audio_path}"
if key not in self._words:
reader = ReadAMBE(languages, audio_path)
result = reader.readfiles()
for lang, words in result.items():
_map = VOICE_MAP.get(lang, {})
for mapword, mapped in _map.items():
if mapped in words:
words[mapword] = words[mapped]
logger.info("(AMBE) for language %s, read %s words into voice dict", lang, len(words) - 1)
self._words[key] = result
return self._words[key]
def read_single_file(self, audio_path: str, lang: str, file_number: str) -> list:
"""Read one .ambe file (e.g. {lang}/ondemand/{file_number}.ambe). Legacy readSingleFile."""
rel = os.path.join(lang, "ondemand", f"{file_number}.ambe")
reader = ReadAMBE(lang, audio_path)
return reader.readSingleFile(rel)
def pkt_gen(
self, rf_src: bytes, dst_id: bytes, peer: bytes, slot: int, phrase: list[Any]
) -> Iterator[bytes]:
"""Generate HBP voice packets for phrase. Legacy mk_voice.pkt_gen."""
return _pkt_gen(rf_src, dst_id, peer, slot, phrase)
def ensure_tts_ambe(self, text: str, lang: str, out_path: str, config: dict[str, Any]) -> str | None:
"""Return out_path if .ambe file exists (cached). Full TTS conversion is in tts_engine.ensure_tts_ambe."""
if out_path and os.path.isfile(out_path):
return out_path
return None
class StubVoiceProvider(VoiceProvider):
"""Stub: get_ambe_words returns empty; pkt_gen returns empty iterator; read_single_file returns []."""
def get_ambe_words(self, languages: str, audio_path: str) -> dict[str, dict[str, Any]]:
return {}
def pkt_gen(
self, rf_src: bytes, dst_id: bytes, peer: bytes, slot: int, phrase: list[Any]
) -> Iterator[bytes]:
return iter([])
def ensure_tts_ambe(self, text: str, lang: str, out_path: str, config: dict[str, Any]) -> str | None:
return None
def read_single_file(self, audio_path: str, lang: str, file_number: str) -> list:
return []

@ -0,0 +1,96 @@
# ADN DMR Peer Server - voice packet generator (legacy mk_voice.pkt_gen)
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink (mk_voice.py). GPLv3.
"""Generate HBP DMRD voice packets for a phrase. Legacy mk_voice.pkt_gen."""
from __future__ import annotations
from random import randint
from typing import Any, Iterator
from bitarray import bitarray
from dmr_utils3 import bptc
from dmr_utils3.const import EMB, BS_DATA_SYNC, BS_VOICE_SYNC, LC_OPT, SLOT_TYPE
from dmr_utils3.utils import bytes_4
# Precalculated DMRD byte 15 (slot << 7 | this)
HEADBITS = 0b00100001
BURSTBITS = [0b00010000, 0b00000001, 0b00000010, 0b00000011, 0b00000100, 0b00000101]
TERMBITS = 0b00100010
NULL_EMB_LC = bitarray(endian="big")
NULL_EMB_LC.frombytes(b"\x00\x00\x00\x00")
TAIL = b"\x00\x00"
def pkt_gen(
rf_src: bytes,
dst_id: bytes,
peer: bytes,
slot: int,
phrase: list[Any],
) -> Iterator[bytes]:
"""Generate DMRD voice packets for phrase. Each word in phrase is list of [b0, b1] burst pairs (bitarray)."""
stream_id = bytes_4(randint(0x00, 0xFFFFFFFF))
sdp = rf_src + dst_id + peer
lc = LC_OPT + dst_id + rf_src
head_lc = bptc.encode_header_lc(lc)
head_lc = [head_lc[:98], head_lc[-98:]]
term_lc = bptc.encode_terminator_lc(lc)
term_lc = [term_lc[:98], term_lc[-98:]]
emb_lc = bptc.encode_emblc(lc)
embed = [
BS_VOICE_SYNC,
EMB["BURST_B"][:8] + emb_lc[1] + EMB["BURST_B"][-8:],
EMB["BURST_C"][:8] + emb_lc[2] + EMB["BURST_C"][-8:],
EMB["BURST_D"][:8] + emb_lc[3] + EMB["BURST_D"][-8:],
EMB["BURST_E"][:8] + emb_lc[4] + EMB["BURST_E"][-8:],
EMB["BURST_F"][:8] + NULL_EMB_LC + EMB["BURST_F"][-8:],
]
seq = 0
slot_byte = (slot << 7) & 0xFF
for _ in range(3):
pkt = (
b"DMRD"
+ bytes([seq])
+ sdp
+ bytes([slot_byte | HEADBITS])
+ stream_id
+ (head_lc[0] + SLOT_TYPE["VOICE_LC_HEAD"][:10] + BS_DATA_SYNC + SLOT_TYPE["VOICE_LC_HEAD"][-10:] + head_lc[1]).tobytes()
+ TAIL
)
seq = (seq + 1) % 0x100
yield pkt
for word in phrase:
for burst in range(len(word)):
b0, b1 = word[burst][0], word[burst][1]
pkt = (
b"DMRD"
+ bytes([seq])
+ sdp
+ bytes([slot_byte | BURSTBITS[burst % 6]])
+ stream_id
+ (b0 + embed[burst % 6] + b1).tobytes()
+ TAIL
)
seq = (seq + 1) % 0x100
yield pkt
pkt = (
b"DMRD"
+ bytes([seq])
+ sdp
+ bytes([slot_byte | TERMBITS])
+ stream_id
+ (term_lc[0] + SLOT_TYPE["VOICE_LC_TERM"][:10] + BS_DATA_SYNC + SLOT_TYPE["VOICE_LC_TERM"][-10:] + term_lc[1]).tobytes()
+ TAIL
)
yield pkt

@ -0,0 +1,112 @@
# ADN DMR Peer Server - voice recording (legacy _handleRecording, _saveRecording)
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""Record DMR voice to AMBE file when RECORDING_ENABLED and packet matches RECORDING_TG / RECORDING_TIMESLOT."""
from __future__ import annotations
import logging
import os
import time
from typing import Any
from bitarray import bitarray
logger = logging.getLogger(__name__)
RECORDING_MAX_FRAMES = 2750
class RecordingHandler:
"""Legacy _handleRecording + _saveRecording: accumulate voice bursts, save to Audio/{lang}/ondemand/{file}.ambe."""
def __init__(self, config: dict[str, Any], project_root: str) -> None:
self._config = config
self._project_root = project_root
self._active = False
self._stream_id: bytes | None = None
self._bursts = bitarray(endian="big")
self._start_time = 0.0
self._frames = 0
self._rf_src: bytes | None = None
def handle_recording(
self,
dmrpkt: bytes,
frame_type: int,
dtype_vseq: int,
stream_id: bytes,
pkt_time: float,
rf_src: bytes,
int_dst_id: int,
slot: int,
) -> None:
"""Legacy _handleRecording: accumulate voice or start/end stream."""
g = self._config.get("VOICE", {})
if not g.get("RECORDING_ENABLED"):
return
if int_dst_id != g.get("RECORDING_TG", 0) or slot != g.get("RECORDING_TIMESLOT", 2):
return
if frame_type == 2 and dtype_vseq == 1: # HBPF_DATA_SYNC, HBPF_SLT_VHEAD
if self._active and self._stream_id != stream_id:
logger.info("(RECORDING) New transmission detected, saving previous recording (%d frames)", self._frames)
self._save_recording()
self._active = True
self._stream_id = stream_id
self._bursts = bitarray(endian="big")
self._start_time = pkt_time
self._frames = 0
self._rf_src = rf_src
logger.info("(RECORDING) Recording started - SUB: %s, TG: %s, TS: %s", int.from_bytes(rf_src, "big") if rf_src else 0, int_dst_id, slot)
return
if not self._active or self._stream_id != stream_id:
return
if frame_type in (0, 1): # HBPF_VOICE, HBPF_VOICE_SYNC
_bits_data = bitarray(endian="big")
_bits_data.frombytes(dmrpkt)
if len(_bits_data) >= 264:
self._bursts.extend(_bits_data[:108])
self._bursts.extend(_bits_data[156:264])
self._frames += 1
if self._frames >= RECORDING_MAX_FRAMES:
logger.info("(RECORDING) Max duration reached (%d frames), saving", self._frames)
self._save_recording()
return
if frame_type == 2 and dtype_vseq == 2: # HBPF_SLT_VTERM
self._save_recording()
return
def _save_recording(self) -> None:
"""Legacy _saveRecording: write bursts to Audio/{lang}/ondemand/{file}.ambe."""
if not self._active or self._frames == 0:
self._active = False
self._stream_id = None
logger.warning("(RECORDING) No frames recorded, discarding")
return
g = self._config.get("VOICE", {})
lang = g.get("RECORDING_LANGUAGE", "en_GB")
file_name = g.get("RECORDING_FILE", "recording")
audio_path = os.path.join(self._project_root, g.get("AUDIO_PATH", "Audio"))
out_dir = os.path.join(audio_path, lang, "ondemand")
os.makedirs(out_dir, exist_ok=True)
path = os.path.join(out_dir, file_name + ".ambe")
try:
with open(path, "wb") as f:
f.write(self._bursts.tobytes())
except OSError as e:
logger.warning("(RECORDING) Could not save: %s", e)
self._active = False
self._stream_id = None
self._bursts = bitarray(endian="big")
self._frames = 0
self._rf_src = None
return
duration = time.time() - self._start_time
_id = int.from_bytes(self._rf_src, "big") if self._rf_src else 0
logger.info("(RECORDING) Recording saved: %s (%d frames, %.1f seconds, SUB: %s)", path, self._frames, duration, _id)
self._active = False
self._stream_id = None
self._bursts = bitarray(endian="big")
self._frames = 0
self._rf_src = None

@ -0,0 +1,17 @@
# ADN DMR Peer Server - TTS to AMBE stub
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""TTS to AMBE file (legacy tts_engine.ensure_tts_ambe). Stub until ported."""
from __future__ import annotations
from typing import Any
def ensure_tts_ambe(text: str, lang: str, out_path: str, config: dict[str, Any]) -> str | None:
"""
Generate TTS audio and convert to AMBE file; return output path or None.
Legacy: tts_engine.ensure_tts_ambe. Full implementation to be ported when needed.
"""
return None

@ -0,0 +1,323 @@
# ADN DMR Peer Server - TTS engine (legacy tts_engine.py)
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server tts_engine. GPLv3.
"""
TTS Engine: convert .txt to .ambe for DMR.
Pipeline: .txt -> gTTS -> .mp3 -> ffmpeg -> .wav (8kHz mono 16-bit) -> vocoder/AMBEServer -> .ambe
"""
from __future__ import annotations
import logging
import os
import socket
import struct
import subprocess
import wave
from typing import Any
logger = logging.getLogger(__name__)
_LANG_MAP = {
"es_ES": "es", "en_GB": "en", "en_US": "en", "fr_FR": "fr",
"de_DE": "de", "it_IT": "it", "pt_PT": "pt", "pt_BR": "pt",
"pl_PL": "pl", "nl_NL": "nl", "da_DK": "da", "sv_SE": "sv",
"no_NO": "no", "el_GR": "el", "th_TH": "th", "cy_GB": "cy",
"ca_ES": "ca", "gl_ES": "gl", "eu_ES": "eu",
}
DV3K_START_BYTE = 0x61
DV3K_TYPE_CONTROL = 0x00
DV3K_TYPE_AMBE = 0x01
DV3K_TYPE_AUDIO = 0x02
DV3K_AMBE_FIELD_ID = 0x01
DV3K_AUDIO_FIELD_ID = 0x00
DV3K_SAMPLES_PER_FRAME = 160
DV3K_RATET_DMR = bytes([0x61, 0x00, 0x02, 0x00, 0x09, 0x21])
DV3K_PRODID_REQ = bytes([0x61, 0x00, 0x01, 0x00, 0x30])
def _get_tts_lang(announcement_language: str) -> str:
if announcement_language in _LANG_MAP:
return _LANG_MAP[announcement_language]
return announcement_language[:2] if len(announcement_language) >= 2 else announcement_language
def _generate_tts_audio(text: str, lang: str, mp3_path: str) -> bool:
try:
from gtts import gTTS
except ImportError:
logger.error("(TTS) gTTS not installed. Run: pip install gTTS")
return False
try:
tts = gTTS(text=text, lang=lang, slow=False)
tts.save(mp3_path)
logger.info("(TTS) TTS audio generated: %s", mp3_path)
return True
except Exception as e:
logger.error("(TTS) Error generating TTS audio: %s", e)
return False
def _convert_to_wav(mp3_path: str, wav_path: str, volume_db: int = 0, speed: float = 1.0) -> bool:
speed = max(0.5, min(2.0, speed))
_filters: list[str] = []
if speed != 1.0:
_filters.append('atempo={:.2f}'.format(speed))
logger.info('(TTS) Aplicando velocidad: x%.2f', speed)
if volume_db != 0:
_filters.append('volume={}dB'.format(volume_db))
logger.info("(TTS) Applying volume adjustment: %ddB", volume_db)
cmd = ["ffmpeg", "-y", "-i", mp3_path, "-ar", "8000", "-ac", "1", "-sample_fmt", "s16"]
if _filters:
cmd += ["-af", ",".join(_filters)]
cmd += ["-f", "wav", wav_path]
try:
result = subprocess.run(cmd, capture_output=True, timeout=60)
if result.returncode != 0:
logger.error("(TTS) ffmpeg error: %s", (result.stderr or b"").decode("utf-8", errors="ignore")[:500])
return False
logger.info("(TTS) Audio converted to 8kHz mono WAV: %s", wav_path)
return True
except FileNotFoundError:
logger.error("(TTS) ffmpeg not found. Install ffmpeg on the system")
return False
except subprocess.TimeoutExpired:
logger.error("(TTS) ffmpeg conversion timeout")
return False
except Exception as e:
logger.error("(TTS) Audio conversion error: %s", e)
return False
def _encode_ambe_vocoder(wav_path: str, ambe_path: str, vocoder_cmd: str) -> bool:
if not vocoder_cmd:
return False
cmd = vocoder_cmd.replace("{wav}", wav_path).replace("{ambe}", ambe_path)
try:
result = subprocess.run(cmd, shell=True, capture_output=True, timeout=120)
if result.returncode != 0:
logger.error("(TTS) Vocoder error: %s", (result.stderr or b"").decode("utf-8", errors="ignore")[:500])
return False
if not os.path.isfile(ambe_path):
logger.error("(TTS) Vocoder did not produce AMBE file: %s", ambe_path)
return False
logger.info("(TTS) Audio encoded to AMBE via external vocoder: %s", ambe_path)
return True
except subprocess.TimeoutExpired:
logger.error("(TTS) AMBE encoding timeout")
return False
except Exception as e:
logger.error("(TTS) Error running vocoder: %s", e)
return False
def _build_audio_packet(pcm_samples: list[int]) -> bytes:
payload = struct.pack("BB", DV3K_AUDIO_FIELD_ID, len(pcm_samples))
for sample in pcm_samples:
payload += struct.pack(">h", sample)
header = bytes([DV3K_START_BYTE]) + struct.pack(">HB", len(payload), DV3K_TYPE_AUDIO)
return header + payload
def _parse_ambe_response(data: bytes) -> bytes | None:
if len(data) < 4 or data[0] != DV3K_START_BYTE:
return None
_payload_len = struct.unpack(">H", data[1:3])[0]
_pkt_type = data[3]
if _pkt_type == DV3K_TYPE_AMBE and len(data) > 5:
_field_id = data[4]
if _field_id == DV3K_AMBE_FIELD_ID:
_num_bits = data[5]
_num_bytes = (_num_bits + 7) // 8
return data[6 : 6 + _num_bytes]
return None
def _encode_ambe_ambeserver(wav_path: str, ambe_path: str, host: str, port: int) -> bool:
host = host.strip().strip('"').strip("'")
logger.info("(TTS-AMBESERVER) Connecting to AMBEServer %s:%d", host, port)
try:
resolved = socket.gethostbyname(host)
if resolved != host:
logger.info("(TTS-AMBESERVER) Host %s resolved to %s", host, resolved)
host = resolved
except socket.gaierror as e:
logger.error('(TTS-AMBESERVER) Cannot resolve host "%s": %s', host, e)
return False
try:
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
sock.settimeout(5.0)
except Exception as e:
logger.error("(TTS-AMBESERVER) Error creating UDP socket: %s", e)
return False
try:
sock.sendto(DV3K_PRODID_REQ, (host, port))
data, _ = sock.recvfrom(1024)
if data[0] != DV3K_START_BYTE:
logger.error("(TTS-AMBESERVER) Invalid response from AMBEServer")
sock.close()
return False
logger.info("(TTS-AMBESERVER) AMBEServer connected")
except socket.timeout:
logger.error("(TTS-AMBESERVER) Timeout connecting to AMBEServer %s:%d", host, port)
sock.close()
return False
except Exception as e:
logger.error("(TTS-AMBESERVER) Connection error: %s", e)
sock.close()
return False
try:
sock.sendto(DV3K_RATET_DMR, (host, port))
data, _ = sock.recvfrom(1024)
if data[0] != DV3K_START_BYTE:
logger.error("(TTS-AMBESERVER) Error setting RATET DMR")
sock.close()
return False
logger.info("(TTS-AMBESERVER) RATET DMR configured")
except (socket.timeout, Exception):
sock.close()
return False
try:
wf = wave.open(wav_path, "rb")
except Exception as e:
logger.error("(TTS-AMBESERVER) Error opening WAV: %s", e)
sock.close()
return False
if wf.getsampwidth() != 2 or wf.getnchannels() != 1:
logger.error("(TTS-AMBESERVER) WAV must be mono 16-bit PCM")
wf.close()
sock.close()
return False
_total_frames = wf.getnframes()
_sample_rate = wf.getframerate()
logger.info("(TTS-AMBESERVER) WAV: %d samples, %d Hz", _total_frames, _sample_rate)
_raw_frames = wf.readframes(_total_frames)
wf.close()
_samples = list(struct.unpack("<" + "h" * _total_frames, _raw_frames))
_ambe_frames: list[bytes] = []
for i in range(0, len(_samples), DV3K_SAMPLES_PER_FRAME):
_chunk = _samples[i : i + DV3K_SAMPLES_PER_FRAME]
if len(_chunk) < DV3K_SAMPLES_PER_FRAME:
_chunk = _chunk + [0] * (DV3K_SAMPLES_PER_FRAME - len(_chunk))
_audio_pkt = _build_audio_packet(_chunk)
try:
sock.sendto(_audio_pkt, (host, port))
data, _ = sock.recvfrom(1024)
_ambe_data = _parse_ambe_response(data)
if _ambe_data is not None:
_ambe_frames.append(_ambe_data)
except (socket.timeout, Exception):
pass
sock.close()
if not _ambe_frames:
logger.error("(TTS-AMBESERVER) No AMBE frames received")
return False
try:
with open(ambe_path, "wb") as f:
for frame in _ambe_frames:
f.write(frame)
except Exception as e:
logger.error("(TTS-AMBESERVER) Error writing AMBE: %s", e)
return False
logger.info("(TTS-AMBESERVER) Encoding completed: %s", ambe_path)
return True
def _cleanup(files: list[str]) -> None:
for f in files:
try:
if os.path.isfile(f):
os.remove(f)
except Exception:
pass
def text_to_ambe(
txt_path: str,
ambe_path: str,
language: str,
vocoder_cmd: str,
ambeserver_host: str = "",
ambeserver_port: int = 2460,
volume_db: int = 0,
speed: float = 1.0,
) -> bool:
"""Convert .txt to .ambe (gTTS -> mp3 -> ffmpeg -> wav -> vocoder/AMBEServer)."""
if not os.path.isfile(txt_path):
logger.warning("(TTS) Text file not found: %s", txt_path)
return False
if os.path.isfile(ambe_path):
if os.path.getmtime(ambe_path) > os.path.getmtime(txt_path):
logger.info("(TTS) Using cached AMBE (newer than .txt): %s", ambe_path)
return True
with open(txt_path, "r", encoding="utf-8") as f:
text = f.read().strip()
if not text:
logger.warning("(TTS) Text file is empty: %s", txt_path)
return False
logger.info("(TTS) Converting text to AMBE: %s (%d chars, language: %s)", txt_path, len(text), language)
_dir = os.path.dirname(ambe_path)
if _dir:
os.makedirs(_dir, exist_ok=True)
_base = os.path.splitext(ambe_path)[0]
_mp3_path = _base + ".mp3"
_wav_path = _base + ".wav"
_tts_lang = _get_tts_lang(language)
if not _generate_tts_audio(text, _tts_lang, _mp3_path):
return False
if not _convert_to_wav(_mp3_path, _wav_path, volume_db, speed):
_cleanup([_mp3_path])
return False
_encoded = False
if ambeserver_host:
logger.info("(TTS) Using AMBEServer %s:%d", ambeserver_host, ambeserver_port)
_encoded = _encode_ambe_ambeserver(_wav_path, ambe_path, ambeserver_host, ambeserver_port)
if not _encoded:
logger.warning("(TTS) AMBEServer failed, trying external vocoder...")
if not _encoded and vocoder_cmd:
logger.info("(TTS) Using external vocoder")
_encoded = _encode_ambe_vocoder(_wav_path, ambe_path, vocoder_cmd)
if not _encoded:
logger.warning("(TTS) Could not encode to AMBE. Configure TTS_AMBESERVER_HOST or TTS_VOCODER_CMD.")
_cleanup([_mp3_path, _wav_path])
return False
_cleanup([_mp3_path, _wav_path])
logger.info("(TTS) Conversion completed: %s -> %s", txt_path, ambe_path)
return True
def ensure_tts_ambe(config: dict[str, Any], tts_num: int, audio_path: str) -> str | None:
"""Ensure .ambe exists for TTS_ANNOUNCEMENT{tts_num}; create from .txt if needed. Returns path or None."""
g = config.get("VOICE", {})
_prefix = "TTS_ANNOUNCEMENT{}".format(tts_num)
if not g.get("{}_ENABLED".format(_prefix), False):
return None
_file = g.get("{}_FILE".format(_prefix), "")
_lang = g.get("{}_LANGUAGE".format(_prefix), "en_GB")
if not _file:
return None
_txt_path = os.path.join(audio_path, _lang, "ondemand", _file + ".txt")
_ambe_path = os.path.join(audio_path, _lang, "ondemand", _file + ".ambe")
if os.path.isfile(_ambe_path):
if not os.path.isfile(_txt_path):
logger.info("(TTS-%d) Using existing AMBE file (no .txt): %s", tts_num, _ambe_path)
return _ambe_path
if os.path.getmtime(_ambe_path) > os.path.getmtime(_txt_path):
logger.debug("(TTS-%d) Using cached AMBE: %s", tts_num, _ambe_path)
return _ambe_path
if not os.path.isfile(_txt_path):
logger.warning("(TTS-%d) Text file not found: %s", tts_num, _txt_path)
return None
_vocoder_cmd = g.get("TTS_VOCODER_CMD", "")
_ambeserver_host = g.get("TTS_AMBESERVER_HOST", "").strip()
_ambeserver_port = int(g.get("TTS_AMBESERVER_PORT", 2460))
_volume_db = int(g.get("TTS_VOLUME", -3))
_speed = float(g.get("TTS_SPEED", 1.0))
if text_to_ambe(_txt_path, _ambe_path, _lang, _vocoder_cmd, _ambeserver_host, _ambeserver_port, _volume_db, _speed):
return _ambe_path
if os.path.isfile(_ambe_path):
logger.warning("(TTS-%d) Using previous AMBE (conversion failed): %s", tts_num, _ambe_path)
return _ambe_path
return None

@ -0,0 +1,81 @@
# ADN DMR Peer Server - i18n voice key mapping
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server i8n_voice_map.py (Simon Adlem, G7RZU). GPLv3.
"""Map logical keys (e.g. 'A') to AMBE file names (e.g. 'alpha') per language."""
VOICE_MAP: dict[str, dict[str, str]] = {
"en_GB": {
"A": "alpha", "B": "bravo", "C": "charlie", "D": "delta", "E": "echo",
"F": "foxtrot", "G": "golf", "H": "hotel", "I": "india", "J": "juliet",
"K": "kilo", "L": "lima", "M": "mike", "N": "november", "O": "oscar",
"P": "papa", "Q": "quebec", "R": "romeo", "S": "sierra", "T": "tango",
"U": "uniform", "V": "victor", "W": "whiskey", "X": "x-ray", "Y": "yankee",
"Z": "zulu", "to": "silence", "notlinked": "not-linked", "linkedto": "linked-to",
},
"en_GB_2": {
"A": "alpha", "B": "bravo", "C": "charlie", "D": "delta", "E": "echo",
"F": "foxtrot", "G": "golf", "H": "hotel", "I": "india", "J": "juliet",
"K": "kilo", "L": "lima", "M": "mike", "N": "november", "O": "oscar",
"P": "papa", "Q": "quebec", "R": "romeo", "S": "sierra", "T": "tango",
"U": "uniform", "V": "victor", "W": "whiskey", "X": "x-ray", "Y": "yankee",
"Z": "zulu", "to": "silence", "notlinked": "not-linked", "linkedto": "linked-to",
},
"cy_GB": {
"A": "alpha", "B": "bravo", "C": "charlie", "D": "delta", "E": "echo",
"F": "foxtrot", "G": "golf", "H": "hotel", "I": "india", "J": "juliet",
"K": "kilo", "L": "lima", "M": "mike", "N": "november", "O": "oscar",
"P": "papa", "Q": "quebec", "R": "romeo", "S": "sierra", "T": "tango",
"U": "uniform", "V": "victor", "W": "whiskey", "X": "x-ray", "Y": "yankee",
"Z": "zulu", "to": "silence", "notlinked": "not-linked", "linkedto": "linked-to",
"allstar-link-mode": "alpha",
},
"en_US": {
"to": "2", "adn": "silence", "this-is": "silence", "allstar-link-mode": "alpha",
},
"es_ES": {
"0": "zero", "1": "one", "2": "two", "3": "three", "4": "four",
"5": "five", "6": "six", "7": "seven", "8": "eight", "9": "nine",
"A": "alfa", "B": "bravo", "C": "charlie", "D": "delta", "E": "echo",
"F": "foxtrot", "G": "golf", "H": "hotel", "I": "india", "J": "juliet",
"K": "kilo", "L": "lima", "M": "mike", "N": "november", "O": "oscar",
"P": "papa", "Q": "quebec", "R": "romeo", "S": "sierra", "T": "tango",
"U": "uniform", "V": "victor", "W": "whiskey", "X": "x-ray", "Y": "yankee",
"Z": "zulu", "to": "silence", "notlinked": "not-linked", "linkedto": "linked-to",
"allstar-link-mode": "alfa",
},
"fr_FR": {
"A": "alpha", "B": "bravo", "C": "charlie", "D": "delta", "E": "echo",
"F": "foxtrot", "G": "golf", "H": "hotel", "I": "india", "J": "juliet",
"K": "kilo", "L": "lima", "M": "mike", "N": "november", "O": "oscar",
"P": "papa", "Q": "quebec", "R": "romeo", "S": "sierra", "T": "tango",
"U": "uniform", "V": "victor", "W": "whiskey", "X": "x-ray", "Y": "yankee",
"Z": "zulu", "to": "silence", "notlinked": "not-linked", "linkedto": "linked-to",
"allstar-link-mode": "alpha",
},
"pt_PT": {
"A": "alpha", "B": "bravo", "C": "charlie", "D": "delta", "E": "echo",
"F": "foxtrot", "G": "golf", "H": "hotel", "I": "india", "J": "juliet",
"K": "kilo", "L": "lima", "M": "mike", "N": "november", "O": "oscar",
"P": "papa", "Q": "quebec", "R": "romeo", "S": "sierra", "T": "tango",
"U": "uniform", "V": "victor", "W": "whiskey", "X": "x-ray", "Y": "yankee",
"Z": "zulu", "to": "silence", "notlinked": "not-linked", "linkedto": "linked-to",
"allstar-link-mode": "alpha",
},
"el_GR": {
"A": "alpha", "B": "bravo", "C": "charlie", "D": "delta", "E": "echo",
"F": "foxtrot", "G": "golf", "H": "hotel", "I": "india", "J": "juliet",
"K": "kilo", "L": "lima", "M": "mike", "N": "november", "O": "oscar",
"P": "papa", "Q": "quebec", "R": "romeo", "S": "sierra", "T": "tango",
"U": "uniform", "V": "victor", "W": "whiskey", "X": "x-ray", "Y": "yankee",
"Z": "zulu", "to": "silence", "notlinked": "not-linked", "allstar-link-mode": "alpha",
},
"de_DE": {"to": "silence", "allstar-link-mode": "A"},
"dk_DK": {"to": "silence", "adn": "silence", "this-is": "silence", "allstar-link-mode": "A"},
"it_IT": {"to": "silence", "adn": "silence", "this-is": "silence", "allstar-link-mode": "A"},
"no_NO": {"to": "silence", "adn": "silence", "this-is": "silence", "allstar-link-mode": "A"},
"pl_PL": {"to": "silence", "adn": "silence", "this-is": "silence", "allstar-link-mode": "A"},
"se_SE": {"to": "silence", "adn": "silence", "this-is": "silence", "allstar-link-mode": "A"},
"CW": {"to": "silence", "adn": "silence", "this-is": "silence", "linkedto": "silence", "allstar-link-mode": "T"},
"th_TH": {"to": "silence", "allstar-link-mode": "A"},
}

@ -0,0 +1,553 @@
# ADN DMR Peer Server - entrypoint
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
# Derived from ADN DMR Server / HBlink. GPLv3.
"""
ADN DMR Peer Server entrypoint.
Run: python -m adn_server.main [-c adn-server.yaml] [--logging LEVEL]
Config default: adn-server.yaml at project root.
"""
from __future__ import annotations
import argparse
import copy
import logging
import os
import signal
import sys
import time
from pathlib import Path
from typing import Any
# Ensure package is on path when run as __main__
_ROOT = Path(__file__).resolve().parent.parent.parent
if str(_ROOT) not in sys.path:
sys.path.insert(0, str(_ROOT))
from twisted.internet import reactor, task, threads
from .domain import bytes_3
from .infrastructure import YamlConfigLoader, setup_logging
from .infrastructure.persistence import PickleSubMapStore
from .infrastructure.persistence.keys_store import JsonKeysStore
from .infrastructure.persistence.alias_loader import DefaultAliasLoader
from .infrastructure.twisted_adapters.report_server import ReportServerFactory
from .infrastructure.twisted_adapters.udp_hbp import HBPProtocolFactory
from .infrastructure.bridge_router_impl import InMemoryBridgeRouter
from .infrastructure.voice import DefaultVoiceProvider, StubVoiceProvider
from .infrastructure.security.password_download import DefaultSecurityDownloader, StubSecurityDownloader
from .infrastructure.security.user_passwords_loader import UserPasswordsLoader
from .infrastructure.voice.recording import RecordingHandler
from .application import (
BridgeUseCases,
IdentUseCases,
VoiceUseCases,
ReportingUseCases,
ReportSender,
VoiceProvider,
SecurityDownloader,
)
class ReportSenderAdapter(ReportSender):
"""Adapt ReportServerFactory to ReportSender port."""
def __init__(self, factory: ReportServerFactory) -> None:
self._factory = factory
def send_config(self, systems) -> None:
self._factory.set_systems(systems)
self._factory.send_config()
def send_bridge(self, bridges) -> None:
self._factory.set_bridges(bridges)
self._factory.send_bridge()
def send_bridge_event(self, event: str) -> None:
self._factory.send_bridge_event(event)
def _make_echo_bridges() -> dict:
"""Initial BRIDGES for ECHO peer (legacy make_bridges 9990)."""
now = time.time()
timeout_sec = 2 * 60
return {
"9990": [
{
"SYSTEM": "ECHO",
"TS": 2,
"TGID": bytes_3(9990),
"ACTIVE": True,
"TIMEOUT": timeout_sec,
"TO_TYPE": "NONE",
"ON": [],
"OFF": [],
"RESET": [],
"TIMER": now + timeout_sec,
}
]
}
def _looping_errback(logger: logging.Logger, failure):
"""Errback for LoopingCalls (legacy loopingErrHandle)."""
logger.error("(GLOBAL) Unhandled error in timed loop: %s", failure.getTraceback())
def _expand_generator(config: dict, logger: logging.Logger) -> None:
"""Replace MASTER systems with GENERATOR > 1 by SYSTEM-0, SYSTEM-1, ... (legacy generator)."""
systems = config.get("SYSTEMS", {})
to_remove: list[str] = []
new_systems: dict = {}
for system_name, sys_cfg in list(systems.items()):
if not sys_cfg.get("ENABLED", True):
continue
if sys_cfg.get("MODE") != "MASTER":
continue
generator = int(sys_cfg.get("GENERATOR", 1))
if generator <= 1:
continue
for count in range(generator):
new_name = f"{system_name}-{count}"
new_cfg = copy.deepcopy(sys_cfg)
base_port = int(new_cfg.get("PORT", 56400))
new_cfg["PORT"] = base_port + count
new_cfg["_default_options"] = "SINGLE={};DEFAULT_UA_TIMER={};VOICE={};LANG={}".format(
int(new_cfg.get("SINGLE_MODE", False)),
new_cfg.get("DEFAULT_UA_TIMER", 60),
int(new_cfg.get("VOICE_IDENT", False)),
new_cfg.get("ANNOUNCEMENT_LANGUAGE", "en_GB"),
)
new_systems[new_name] = new_cfg
logger.debug("(GLOBAL) Generator - generated system %s", new_name)
to_remove.append(system_name)
for name in to_remove:
systems.pop(name, None)
for name, cfg in new_systems.items():
systems[name] = cfg
def _ensure_system_runtime_config(config: dict) -> None:
"""Ensure MASTER has PEERS and PEER has STATS (legacy config.py runtime state)."""
for name, sys_cfg in config.get("SYSTEMS", {}).items():
if sys_cfg.get("MODE") == "MASTER":
sys_cfg.setdefault("PEERS", {})
elif sys_cfg.get("MODE") == "PEER":
sys_cfg.setdefault("STATS", {
"CONNECTION": "NO",
"CONNECTED": None,
"PINGS_SENT": 0,
"PINGS_ACKD": 0,
"NUM_OUTSTANDING": 0,
"PING_OUTSTANDING": False,
"LAST_PING_TX_TIME": 0,
"LAST_PING_ACK_TIME": 0,
})
def _normalize_peer_config(config: dict) -> None:
"""Convert PEER systems from YAML to legacy format: MASTER_SOCKADDR, RADIO_ID/CALLSIGN/OPTIONS as bytes (config.py)."""
import socket
for name, sys_cfg in config.get("SYSTEMS", {}).items():
if sys_cfg.get("MODE") != "PEER":
continue
master_ip_str = str(sys_cfg.get("MASTER_IP", "127.0.0.1"))
master_port = int(sys_cfg.get("MASTER_PORT", 56400))
try:
resolved_ip = socket.gethostbyname(master_ip_str)
except OSError:
resolved_ip = master_ip_str
sys_cfg["_MASTER_IP"] = master_ip_str
sys_cfg["MASTER_IP"] = resolved_ip
sys_cfg["MASTER_PORT"] = master_port
sys_cfg["MASTER_SOCKADDR"] = (resolved_ip, master_port)
radio_id = int(sys_cfg.get("RADIO_ID", 0))
sys_cfg["RADIO_ID"] = (radio_id & 0xFFFFFFFF).to_bytes(4, "big")
for field, length in [
("CALLSIGN", 8), ("RX_FREQ", 9), ("TX_FREQ", 9), ("TX_POWER", 2), ("COLORCODE", 2),
("LATITUDE", 8), ("LONGITUDE", 9), ("HEIGHT", 3), ("LOCATION", 20), ("DESCRIPTION", 19),
("SLOTS", 1), ("URL", 124), ("SOFTWARE_ID", 40), ("PACKAGE_ID", 40),
]:
val = sys_cfg.get(field, "")
if isinstance(val, (int, float)):
val = str(val)
b = val.encode("utf-8") if isinstance(val, str) else val
if field == "CALLSIGN":
sys_cfg[field] = b.ljust(length)[:length]
elif field in ("RX_FREQ", "TX_FREQ", "LATITUDE", "LONGITUDE", "LOCATION", "DESCRIPTION", "URL", "SOFTWARE_ID", "PACKAGE_ID"):
sys_cfg[field] = b.ljust(length)[:length]
else:
sys_cfg[field] = b.rjust(length, b"0")[:length] if length <= 3 else b.ljust(length)[:length]
opt = sys_cfg.get("OPTIONS", "")
sys_cfg["OPTIONS"] = opt.encode("utf-8") if isinstance(opt, str) else (opt or b"")
passphrase = sys_cfg.get("PASSPHRASE", "")
sys_cfg["PASSPHRASE"] = passphrase.encode("utf-8") if isinstance(passphrase, str) else (passphrase or b"")
sys_cfg.setdefault("LOOSE", False)
stats = sys_cfg.get("STATS", {})
stats["DNS_TIME"] = time.time()
def _normalize_obp_config(config: dict) -> None:
"""Normalize OPENBRIDGE systems and GLOBAL SERVER_ID (legacy config.py)."""
import socket
g = config.setdefault("GLOBAL", {})
sid = g.get("SERVER_ID", 0)
g["SERVER_ID"] = (int(sid) & 0xFFFFFFFF).to_bytes(4, "big") if not isinstance(sid, bytes) else sid
for name, sys_cfg in config.get("SYSTEMS", {}).items():
if sys_cfg.get("MODE") != "OPENBRIDGE":
continue
net_id = int(sys_cfg.get("NETWORK_ID", 0))
sys_cfg["NETWORK_ID"] = (net_id & 0xFFFFFFFF).to_bytes(4, "big")
target_ip = str(sys_cfg.get("TARGET_IP", ""))
target_port = int(sys_cfg.get("TARGET_PORT", 62044))
if target_ip:
try:
resolved = socket.gethostbyname(target_ip)
sys_cfg["TARGET_IP"] = resolved
sys_cfg["TARGET_SOCK"] = (resolved, target_port)
except OSError:
sys_cfg["TARGET_IP"] = None
sys_cfg["TARGET_SOCK"] = (None, target_port)
else:
sys_cfg["TARGET_IP"] = None
sys_cfg["TARGET_SOCK"] = (None, target_port)
# Legacy config.py 359: VER from PROTO_VER (OPENBRIDGE uses VER for send_system and receive check)
ver = int(sys_cfg.get("PROTO_VER", sys_cfg.get("VER", 5)))
if ver in (0, 2, 3) or ver > 5:
ver = 5
sys_cfg["VER"] = ver
# Legacy config.py OPENBRIDGE: PASSPHRASE padded to 20 bytes with nulls (BLAKE2b/HMAC key)
p = sys_cfg.get("PASSPHRASE") or b""
if isinstance(p, str):
p = p.strip().encode("utf-8")
else:
p = p or b""
sys_cfg["PASSPHRASE"] = (p + b"\x00" * 20)[:20]
sys_cfg.setdefault("RELAX_CHECKS", True)
sys_cfg.setdefault("ENHANCED_OBP", True)
if "TG1_ACL" not in sys_cfg and "TGID_ACL" in sys_cfg:
sys_cfg["TG1_ACL"] = sys_cfg["TGID_ACL"]
def main() -> None:
parser = argparse.ArgumentParser(description="ADN DMR Peer Server")
parser.add_argument("-c", "--config", dest="CONFIG_FILE", default=None, help="Path to adn-server.yaml")
parser.add_argument("--logging", dest="LOG_LEVEL", default=None, help="Override log level")
args = parser.parse_args()
project_root = os.path.dirname(os.path.abspath(__file__))
if project_root.endswith("/adn_server"):
project_root = str(Path(project_root).parent.parent)
config_path = args.CONFIG_FILE or os.path.join(project_root, "adn-server.yaml")
loader = YamlConfigLoader(project_root)
config = loader.load(config_path)
if args.LOG_LEVEL:
config.setdefault("LOGGER", {})["LOG_LEVEL"] = args.LOG_LEVEL
logger = setup_logging(config.get("LOGGER", {}))
logger.info("\n\nCopyright (c) 2026 Rodrigo Pérez, CE5RPY ce5rpy@qmd.cl")
logger.info("\n\nCopyright (c) 2026 Joaquin Madrid Belando, EA5GVK ea5gvk@gmail.com")
logger.info("\nCopyright (c) 2024-2026 Esteban Mackay, HP3ICC setcom40@gmail.com")
logger.info("\nCopyright (c) 2020-2023 Simon G7RZU simon@gb7fr.org.uk")
logger.info("\nCopyright (c) 2013, 2014, 2015, 2016, 2018, 2019\n\tThe Regents of the K0USY Group. All rights reserved.")
logger.debug("\n\n(GLOBAL) Logging system started, anything from here on gets logged")
# Aliases
alias_loader = DefaultAliasLoader()
peer_ids, subscriber_ids, talkgroup_ids, local_subscriber_ids, server_ids, checksums = (
alias_loader.load_aliases(config)
)
config["_SUB_IDS"] = subscriber_ids
config["_PEER_IDS"] = peer_ids
config["_TG_IDS"] = talkgroup_ids
config["_LOCAL_SUBSCRIBER_IDS"] = local_subscriber_ids
config["_SERVER_IDS"] = server_ids
config["CHECKSUMS"] = checksums
# SUB_MAP (shared mutable; used by SubMapTrimmer and shutdown)
aliases_cfg = config.get("ALIASES", {})
data_path = (aliases_cfg.get("PATH") or ".").rstrip("/")
sub_map_file = aliases_cfg.get("SUB_MAP_FILE") or "sub_map.pkl"
sub_map_path = os.path.join(project_root, data_path, sub_map_file)
sub_map_store = PickleSubMapStore()
sub_map = sub_map_store.load(sub_map_path)
config["_SUB_MAP"] = sub_map
# Generator: expand MASTER systems with GENERATOR > 1 into SYSTEM-0, SYSTEM-1, ... (legacy)
_expand_generator(config, logger)
_ensure_system_runtime_config(config)
_normalize_peer_config(config)
_normalize_obp_config(config)
# BRIDGES
bridge_router = InMemoryBridgeRouter()
systems_cfg = config.get("SYSTEMS", {})
if "ECHO" in systems_cfg and systems_cfg.get("ECHO", {}).get("MODE") == "PEER":
bridge_router.set_bridges(_make_echo_bridges())
else:
bridge_router.set_bridges({})
# Protocol registry for send_to_system (legacy: systems[name].send_system(packet))
protocols: dict[str, Any] = {}
report_factory = ReportServerFactory(config)
def send_to_system(system_name: str, packet: bytes, **kwargs: Any) -> None:
p = protocols.get(system_name)
if p is not None and hasattr(p, "send_system"):
p.send_system(packet, **kwargs)
def send_bcsq(system_name: str, tgid: bytes, stream_id: bytes) -> None:
"""Legacy: bridge calls send_bcsq (e.g. loop control first OBP). OBP protocol only."""
p = protocols.get(system_name)
if p is not None and hasattr(p, "_obp_send_bcsq"):
p._obp_send_bcsq(tgid, stream_id)
report_sender = ReportSenderAdapter(report_factory)
report_factory.set_systems(systems_cfg)
report_factory.set_bridges(bridge_router.get_bridges())
reporting_use_cases = ReportingUseCases(report_sender, config)
if config.get("GLOBAL", {}).get("URL_SECURITY", "").strip():
security = DefaultSecurityDownloader(project_root)
else:
security = StubSecurityDownloader()
security.init_downloads(config)
user_passwords_loader = UserPasswordsLoader(project_root)
user_passwords_loader.load(config)
recording_handler = RecordingHandler(config, project_root)
# Voice: AMBE words + pkt_gen (legacy readAMBE + mk_voice). Use Default when languages + audio path set.
audio_path = os.path.join(project_root, config.get("VOICE", {}).get("AUDIO_PATH", "Audio"))
ann_lang = (config.get("VOICE", {}).get("ANNOUNCEMENT_LANGUAGES") or "").strip()
if ann_lang and os.path.isdir(audio_path):
voice_provider = DefaultVoiceProvider()
else:
voice_provider = StubVoiceProvider()
def _start_voice_loop(fn, interval: float, now: bool):
lc = task.LoopingCall(fn)
d = lc.start(interval, now=now)
d.addErrback(_looping_errback, logger)
return lc
voice_use_cases = VoiceUseCases(
voice_provider,
config,
get_protocols=lambda: protocols,
call_from_reactor=reactor.callFromThread,
audio_path=audio_path,
get_bridges=bridge_router.get_bridges,
call_later=reactor.callLater,
start_looping_call=_start_voice_loop,
defer_to_thread=threads.deferToThread,
)
voice_use_cases.check_voice_config_reload(config_path)
ident_use_cases = IdentUseCases(
config,
voice_use_cases,
audio_path,
get_protocols=lambda: protocols,
call_from_reactor=reactor.callFromThread,
)
bridge_use_cases = BridgeUseCases(
bridge_router,
config,
send_to_system=send_to_system,
get_protocols=lambda: protocols,
report_factory=report_factory,
on_bridge_deactivated=lambda sys: reactor.callInThread(voice_use_cases.disconnected_voice, sys),
send_bcsq=send_bcsq,
)
bridge_use_cases.apply_startup_bridges()
report_factory.set_bridges(bridge_router.get_bridges())
report_factory.set_systems(config.get("SYSTEMS", {}))
# Report server (same order as legacy: log then listen)
if config.get("REPORTS", {}).get("REPORT", True):
logger.info("(REPORT) HBlink TCP reporting server configured")
port = config["REPORTS"].get("REPORT_PORT", 4321)
reactor.listenTCP(port, report_factory)
logger.info("(REPORT) Report server listening on TCP %s", port)
# Reporting loop (REPORT_INTERVAL) — same logs as legacy after send_config/send_bridge
def reporting_loop():
logger.debug("(REPORT) Periodic reporting loop started")
report_factory.set_systems(config.get("SYSTEMS", {}))
report_factory.set_bridges(bridge_router.get_bridges())
report_factory.send_config()
report_factory.send_bridge()
# Legacy: peer count and SUB_MAP count
systems_with_peers = sum(1 for s in config.get("SYSTEMS", {}) if config.get("SYSTEMS", {}).get(s, {}).get("PEERS"))
logger.info("(REPORT) %s systems have at least one peer", systems_with_peers)
logger.info("(REPORT) Subscriber Map has %s entries", len(config.get("_SUB_MAP", {})))
report_interval = config.get("REPORTS", {}).get("REPORT_INTERVAL", 60)
task.LoopingCall(reporting_loop).start(report_interval).addErrback(_looping_errback, logger)
# LoopingCalls (legacy intervals)
def _rule_timer_in_thread():
bridge_use_cases.rule_timer_loop()
reactor.callFromThread(report_factory.set_bridges, bridge_router.get_bridges())
reactor.callFromThread(report_factory.send_bridge)
task.LoopingCall(lambda: threads.deferToThread(_rule_timer_in_thread)).start(52).addErrback(_looping_errback, logger)
task.LoopingCall(bridge_use_cases.stream_trimmer_loop).start(5).addErrback(_looping_errback, logger)
task.LoopingCall(bridge_use_cases.bridge_reset_loop).start(6).addErrback(_looping_errback, logger)
if config.get("GLOBAL", {}).get("GEN_STAT_BRIDGES", False):
def _stat_trimmer_in_thread():
bridge_use_cases.stat_trimmer_loop()
reactor.callFromThread(report_factory.set_bridges, bridge_router.get_bridges())
reactor.callFromThread(report_factory.send_bridge)
task.LoopingCall(lambda: threads.deferToThread(_stat_trimmer_in_thread)).start(303).addErrback(_looping_errback, logger)
# KA Reporting (legacy kaReporting, 60s)
task.LoopingCall(lambda: threads.deferToThread(reporting_use_cases.ka_reporting_loop)).start(60).addErrback(_looping_errback, logger)
# bridgeDebug (legacy 66s) — remove invalid bridges, fix >1 active dial per MASTER
task.LoopingCall(
lambda: (bridge_use_cases.bridge_debug_loop(), report_factory.set_bridges(bridge_router.get_bridges()), report_factory.send_bridge())
).start(66).addErrback(_looping_errback, logger)
# Alias reload (STALE_DAYS -> seconds)
alias_interval = (aliases_cfg.get("STALE_DAYS") or 1) * 86400
def alias_reload_loop():
logger.debug("(ALIAS) starting alias thread")
try:
p, s, t, l, sv, ch = alias_loader.load_aliases(config)
config["_SUB_IDS"] = s
config["_PEER_IDS"] = p
config["_TG_IDS"] = t
config["_LOCAL_SUBSCRIBER_IDS"] = l
config["_SERVER_IDS"] = sv
config["CHECKSUMS"] = ch
except Exception as e:
logger.warning("(ALIAS) alias reload failed: %s", e)
task.LoopingCall(alias_reload_loop).start(alias_interval).addErrback(_looping_errback, logger)
# SubMapTrimmer (3600s) + save
def sub_map_trimmer_loop():
logger.debug("(SUBSCRIBER) Subscriber Map trimmer loop started")
now = time.time()
to_remove = [k for k, v in sub_map.items() if v[2] < (now - 86400)]
for k in to_remove:
sub_map.pop(k, None)
if aliases_cfg.get("SUB_MAP_FILE"):
try:
sub_map_store.save(sub_map_path, sub_map)
logger.info("(SUBSCRIBER) Writing SUB_MAP to disk")
except Exception as e:
logger.warning("(SUBSCRIBER) Cannot write SUB_MAP to file: %s", e)
task.LoopingCall(sub_map_trimmer_loop).start(3600).addErrback(_looping_errback, logger)
# Kill switch + shutdown (legacy kill_server every 5s; SIGTERM/SIGINT trigger _KILL_SERVER)
config.setdefault("GLOBAL", {})["_KILL_SERVER"] = False
keys_store = JsonKeysStore()
keys_path = os.path.join(project_root, data_path, aliases_cfg.get("KEYS_FILE") or "keys.json")
keys = keys_store.load(keys_path) if keys_path else {}
def kill_server_loop():
try:
if config.get("GLOBAL", {}).get("_KILL_SERVER"):
logger.info("(GLOBAL) SHUTDOWN: CONFBRIDGE IS TERMINATING - killserver called")
if reactor.running:
reactor.stop()
if aliases_cfg.get("SUB_MAP_FILE"):
sub_map_store.save(sub_map_path, sub_map)
try:
keys_store.save(keys_path, keys)
logger.info("(KEYS) saved system keys to keystore")
except Exception as e:
logger.error("(GLOBAL) Cannot save key file: %s", e)
except KeyError:
pass
task.LoopingCall(kill_server_loop).start(5).addErrback(_looping_errback, logger)
def shutdown_handler():
"""On reactor shutdown: save SUB_MAP and keys."""
if aliases_cfg.get("SUB_MAP_FILE"):
try:
sub_map_store.save(sub_map_path, config["_SUB_MAP"])
logger.info("(SUBSCRIBER) Writing SUB_MAP to disk (shutdown)")
except Exception as e:
logger.warning("(SUBSCRIBER) Cannot write SUB_MAP on shutdown: %s", e)
try:
keys_store.save(keys_path, keys)
logger.info("(KEYS) saved system keys to keystore (shutdown)")
except Exception as e:
logger.error("(GLOBAL) Cannot save key file on shutdown: %s", e)
def sig_handler(sig, frame):
logger.info("(GLOBAL) SHUTDOWN: CONFBRIDGE IS TERMINATING WITH SIGNAL %s", sig)
config["GLOBAL"]["_KILL_SERVER"] = True
if reactor.running:
reactor.stop()
signal.signal(signal.SIGTERM, sig_handler)
signal.signal(signal.SIGINT, sig_handler)
reactor.addSystemEventTrigger("before", "shutdown", shutdown_handler)
# Voice/announcement config reload (15s): re-read main YAML, start/stop announcement LoopingCalls on change
def voice_reload_loop():
try:
loader.reload_voice_config(config, config_path)
voice_use_cases.check_voice_config_reload(config_path)
except Exception as e:
logger.debug("(VOICE-RELOAD) %s", e)
task.LoopingCall(voice_reload_loop).start(15).addErrback(_looping_errback, logger)
logger.info("(VOICE-RELOAD) config file watch active (every 15 seconds)")
def security_loop():
security.periodic_download(config)
task.LoopingCall(security_loop).start(300).addErrback(_looping_errback, logger)
logger.info("(SECURITY) Periodic password download task started (every 5 minutes)")
# Ident (3600s): run ident in thread for MASTERs with VOICE_IDENT
def ident_loop():
logger.debug("(IDENT) starting ident thread")
reactor.callInThread(ident_use_cases.run_ident)
task.LoopingCall(ident_loop).start(3600).addErrback(_looping_errback, logger)
task.LoopingCall(bridge_use_cases.options_config_loop).start(26).addErrback(_looping_errback, logger)
task.LoopingCall(bridge_use_cases.log_connected_systems_and_tgs).start(60).addErrback(_looping_errback, logger)
task.LoopingCall(lambda: logger.debug("(ROUTER) KeepAlive reporting loop started")).start(60).addErrback(_looping_errback, logger)
# UDP listeners per system (same order as legacy: SYSTEM STARTING then instance created per system)
logger.info("(GLOBAL) ADN DMR Peer Server -- SYSTEM STARTING...")
for system_name, sys_cfg in systems_cfg.items():
if not sys_cfg.get("ENABLED", True):
continue
ip = sys_cfg.get("IP", "")
udp_port = sys_cfg.get("PORT", 56400)
protocol = HBPProtocolFactory(
system_name,
config,
report_factory,
router=bridge_router,
dmrd_received=bridge_use_cases.dmrd_received,
get_user_password_callback=user_passwords_loader.get_user_password,
on_play_file_request=voice_use_cases.play_file_on_request,
on_handle_recording=recording_handler.handle_recording,
on_in_band_signalling=bridge_use_cases.apply_in_band_signalling,
on_options_received=bridge_use_cases.options_config_for_system,
on_deactivate_dynamic_bridges=bridge_use_cases.deactivate_all_dynamic_bridges,
)
protocols[system_name] = protocol
reactor.listenUDP(udp_port, protocol, interface=ip or "0.0.0.0")
logger.debug("(GLOBAL) %s instance created: %s, %s", sys_cfg.get("MODE", "?"), system_name, protocol)
logger.info("(GLOBAL) UDP %s listening on %s:%s", system_name, ip or "*", udp_port)
logger.info("(GLOBAL) ADN DMR Peer Server started. Use adn-dmr-server as reference.")
reactor.run()
if __name__ == "__main__":
main()
Loading…
Cancel
Save

Powered by TurnKey Linux.