From 616d939f5b358e6544682371456729b6fb5ce937 Mon Sep 17 00:00:00 2001 From: yo Date: Sun, 20 Sep 2026 09:09:19 +0200 Subject: [PATCH] refactor(obp): lift ingress admission out of the hblink port into the domain MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ``_obp_datagram_received`` decided admission inline: ~110 lines of nested ifs repeating the same four-step shape (evaluate, log once per stream_id, quench, return) eight times for the ACLs alone, reading GLOBAL and SYSTEMS state at six different depths, and doing it twice over — once for DMRD v1 and once for DMRE v5. Nothing in it could be exercised without a reactor and a socket. The decisions now live in ``domain/mesh_admission.py`` as pure functions. Each rule takes values and returns a ``Rejection`` (reason, log line, whether to quench, whether to log once) or ``None`` to admit; ``admit_dmrd_v1`` and ``admit_dmre_v5`` chain them in the order the legacy handler applied them. The adapter keeps what is genuinely I/O: one ``_obp_admission_context()`` that reads config once per frame, and one ``_obp_reject()`` that applies the outcome. No behaviour change intended. Beyond the suite, a differential harness fed 9576 cryptographically valid frames (DMRD v1 and DMRE v5, sweeping destination TG, subscriber, bits byte, STUN, both ACL scopes, hop count, packet age, source server and VALIDATE_SERVER_IDS) through a protocol built from this commit and one built from develop, comparing delivery, quench and every log line: no behavioural difference. One deliberate deviation: the four DMRE "GLOBAL TG FILTER (local to ...)" log calls on develop pass three arguments to a message with two placeholders, so Python logs a formatting error instead of the drop (648 of the 9576 frames hit this). The message now carries the talkgroup. ``test_every_rejection_can_be_ formatted`` keeps the whole family honest. Also fixed by construction: ``reason`` gives each drop a stable handle, so a later engine can count or trace drops without parsing log text. Tests: 59 new unit tests, 100% coverage of the new module; full suite 849 passed, 2 skipped (the 2 failures are this machine's: kernel rmem_max and the version installed in the venv, both failing on develop as well). Co-Authored-By: Claude Opus 5 --- src/adn_server/domain/mesh_admission.py | 418 ++++++++++++++++++ .../twisted_adapters/udp_hbp.py | 278 +++++------- tests/domain/test_mesh_admission.py | 393 ++++++++++++++++ 3 files changed, 918 insertions(+), 171 deletions(-) create mode 100644 src/adn_server/domain/mesh_admission.py create mode 100644 tests/domain/test_mesh_admission.py diff --git a/src/adn_server/domain/mesh_admission.py b/src/adn_server/domain/mesh_admission.py new file mode 100644 index 0000000..64039a5 --- /dev/null +++ b/src/adn_server/domain/mesh_admission.py @@ -0,0 +1,418 @@ +# ADN DMR Peer Server - domain mesh admission +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""Admission rules for OpenBridge ingress, as pure decisions. + +Every rule answers one question — may this frame continue? — and returns a +``Rejection`` describing what to log and whether to quench the peer, or ``None`` +to let the frame through. Nothing here touches sockets, the reactor, the clock +or the live SYSTEMS config: the caller passes what the rule needs and applies +the outcome. That is what makes the gauntlet testable with plain values, and +what lets a later engine reuse the same decisions without the Twisted adapter. + +The wire formats live in ``infrastructure.mesh`` (obp_v1, dmre_v5); this module +only decides what to do with a frame once it has been decoded. +""" + +from __future__ import annotations + +import logging +from collections.abc import Callable, Iterable +from dataclasses import dataclass, field +from typing import Any + +from .value_objects import int_id + +# Talkgroups a mesh peer must never hand us: hblink's hard-coded table, kept as +# named constants so the ranges can be read (and one day configured) on sight. +TG_LOCAL_TO_REPEATER_MAX = 79 +TG_LOCAL_TO_SERVER = (9990, 9999) +TG_LOCAL_TO_SERVER_MAIN = (92, 199) +TG_LOCAL_TO_MCC = ((80, 89), (800, 899)) +TG_DATA_GATEWAY = 900999 + +MAX_HOPS = 10 +MAX_PACKET_AGE_S = 5.0 + +_AclCheck = Callable[[bytes, Any], bool] + + +@dataclass(frozen=True) +class CallAttributes: + """Slot and call classification carried by the DMRD bits byte.""" + + slot: int + call_type: str + frame_type: int + dtype_vseq: int + + +def call_attributes(bits: int) -> CallAttributes: + """Decode the bits byte of a DMRD/DMRE frame. + + ``vcsbk`` (a CSBK preamble) shares the data-sync pattern with voice, so it is + recognised before the group fallback — mis-classifying it is what turns a + private data preamble into a dynamic talkgroup on the far side. + """ + if bits & 0x40: + call_type = "unit" + elif (bits & 0x23) == 0x23: + call_type = "vcsbk" + else: + call_type = "group" + return CallAttributes( + slot=2 if (bits & 0x80) else 1, + call_type=call_type, + frame_type=(bits & 0x30) >> 4, + dtype_vseq=bits & 0xF, + ) + + +@dataclass(frozen=True) +class ObpFrame: + """The identity of one ingress frame, as the admission rules see it.""" + + system: str + stream_id: bytes + rf_src: bytes + dst_id: bytes + slot: int + call_type: str + + @property + def dst(self) -> int: + return int(int_id(self.dst_id)) + + +@dataclass(frozen=True) +class Rejection: + """Why a frame was dropped, and what the caller should do about it. + + ``message``/``args`` are kept apart so the caller can hand them straight to + ``logger.log`` and keep lazy formatting. ``reason`` is the stable handle: log + text may be reworded, ``reason`` is what metrics and tests match on. + """ + + reason: str + message: str + args: tuple[Any, ...] = () + level: int = logging.INFO + quench: bool = True + log_once: bool = True + + +@dataclass(frozen=True) +class AclRules: + """One ACL scope (GLOBAL or the system's own).""" + + enabled: bool = False + sub_acl: Any = (True, []) + tg1_acl: Any = (True, []) + + +@dataclass(frozen=True) +class MeshEnvelope: + """DMRE v5 envelope fields the admission rules look at. + + ``source_server_id`` is the same value as ``source_server``, kept in its wire + form because the alias lookup that validates it takes bytes. + """ + + source_server: int + hops: int + timestamp_ns: int = 0 + source_server_id: bytes = b"" + + +@dataclass(frozen=True) +class AdmissionContext: + """Everything outside the frame that the rules depend on.""" + + stunned: bool = False + acl_check: _AclCheck | None = None + global_rules: AclRules = field(default_factory=AclRules) + system_rules: AclRules = field(default_factory=AclRules) + server_id: int = 0 + validate_server_ids: bool = False + known_server_prefixes: Iterable[str] = () + resolve_server_id: Callable[[bytes], Any] | None = None # alias lookup, bytes in + + +def server_prefix(server_id: Any) -> int: + """First four digits of our SERVER_ID, whether it is stored as int or bytes.""" + if isinstance(server_id, bytes): + server_id = int.from_bytes(server_id, "big") + try: + return int(str(int(server_id))[:4]) + except (TypeError, ValueError): + return 0 + + +def check_network_id( + system: str, + stream_id: bytes, + *, + expected: bytes, + received: bytes, + dmre: bool = False, +) -> Rejection | None: + """The peer must send the NETWORK_ID we have configured for this bridge.""" + if expected == received: + return None + label = "OpenBridge DMRE discarded" if dmre else "OpenBridge packet discarded" + return Rejection( + reason="network-id-mismatch", + message="(%s) " + label + " because NETWORK_ID: %s Does not match sent Peer ID: %s", + args=(system, int_id(expected or b""), int_id(received)), + level=logging.ERROR, + quench=False, + ) + + +def check_slot(frame: ObpFrame) -> Rejection | None: + """DMRD v1 over OpenBridge is TS1 only.""" + if frame.slot == 1: + return None + return Rejection( + reason="not-slot-1", + message="(%s) OpenBridge packet discarded because it was not received on slot 1. SID: %s, TGID %s", + args=(frame.system, int_id(frame.rf_src), int_id(frame.dst_id)), + level=logging.ERROR, + quench=False, + log_once=False, + ) + + +def check_stun(frame: ObpFrame, *, stunned: bool) -> Rejection | None: + """A STUNned bridge accepts nothing until the operator lifts it.""" + if not stunned: + return None + return Rejection( + reason="stunned", + message="(%s) Bridge STUNned, discarding", + args=(frame.system,), + level=logging.WARNING, + quench=False, + ) + + +def check_packet_age(frame: ObpFrame, envelope: MeshEnvelope, *, now: float) -> Rejection | None: + """DMRE carries a timestamp; a late frame is a replay or a stalled path.""" + if envelope.timestamp_ns / 1_000_000_000 >= (now - MAX_PACKET_AGE_S): + return None + return Rejection( + reason="stale-packet", + message="(%s) Packet from server %s more than 5s old!, discarding", + args=(frame.system, envelope.source_server), + level=logging.WARNING, + ) + + +def check_source_server( + frame: ObpFrame, + envelope: MeshEnvelope, + ctx: AdmissionContext, +) -> Rejection | None: + """A DMRE source server is a 4-7 digit ID, known to us or a valid DMR ID.""" + digits = str(envelope.source_server) + if len(digits) < 4 or len(digits) > 7: + return Rejection( + reason="source-server-length", + message="(%s) Source Server should be between 4 and 7 digits, discarding Src: %s", + args=(frame.system, envelope.source_server), + level=logging.WARNING, + ) + if ctx.validate_server_ids and len(digits) in (4, 5) and digits[:4] not in ctx.known_server_prefixes: + return Rejection( + reason="source-server-unknown", + message="(%s) Source Server ID is 4 or 5 digits but not in list: %s", + args=(frame.system, envelope.source_server), + level=logging.WARNING, + ) + if len(digits) > 5 and ctx.resolve_server_id is not None: + if not ctx.resolve_server_id(envelope.source_server_id): + return Rejection( + reason="source-server-invalid", + message="(%s) Source Server 6 or 7 digits but not a valid DMR ID, discarding Src: %s", + args=(frame.system, envelope.source_server), + level=logging.WARNING, + ) + return None + + +def check_hops(frame: ObpFrame, envelope: MeshEnvelope) -> Rejection | None: + """Every mesh hop bumps the counter; past MAX_HOPS the frame is looping.""" + hops = envelope.hops + 1 + if hops <= MAX_HOPS: + return None + return Rejection( + reason="max-hops", + message="(%s) MAX HOPS exceed, dropping. Hops: %s, DST: %s, SRC: %s", + args=(frame.system, hops, frame.dst, envelope.source_server), + level=logging.DEBUG, + log_once=False, + ) + + +def _tg_filter(frame: ObpFrame, reason: str, scope: str) -> Rejection: + return Rejection( + reason=reason, + message="(%s) CALL DROPPED WITH STREAM ID %s ON TG %s BY GLOBAL TG FILTER (%s)", + args=(frame.system, int_id(frame.stream_id), frame.dst, scope), + ) + + +def check_tg_filter_v1(frame: ObpFrame) -> Rejection | None: + """DMRD v1 filter: one rule for every talkgroup that must stay local.""" + if frame.call_type == "unit": + return None + dst = frame.dst + if ( + dst <= TG_LOCAL_TO_REPEATER_MAX + or TG_LOCAL_TO_SERVER[0] <= dst <= TG_LOCAL_TO_SERVER[1] + or TG_LOCAL_TO_SERVER_MAIN[0] <= dst <= TG_LOCAL_TO_SERVER_MAIN[1] + or dst == TG_DATA_GATEWAY + ): + return Rejection( + reason="tg-filter", + message="(%s) CALL DROPPED WITH STREAM ID %s FROM SUBSCRIBER %s BY GLOBAL TG FILTER", + args=(frame.system, int_id(frame.stream_id), dst), + ) + return None + + +def check_tg_filter_v5( + frame: ObpFrame, + envelope: MeshEnvelope, + ctx: AdmissionContext, +) -> Rejection | None: + """DMRE v5 filter: same idea, but the last two rules depend on who sent it. + + A talkgroup local to a server, or to an MCC, is only refused when the source + server is *not* part of that server or that MCC. + """ + if frame.call_type == "unit": + return None + dst = frame.dst + if dst <= TG_LOCAL_TO_REPEATER_MAX: + return _tg_filter(frame, "tg-filter-repeater", "local to repeater") + if TG_LOCAL_TO_SERVER[0] <= dst <= TG_LOCAL_TO_SERVER[1] or dst == TG_DATA_GATEWAY: + return _tg_filter(frame, "tg-filter-server", "local to server") + source = str(envelope.source_server) + if TG_LOCAL_TO_SERVER_MAIN[0] <= dst <= TG_LOCAL_TO_SERVER_MAIN[1]: + if int(source[:4]) != ctx.server_id: + return _tg_filter(frame, "tg-filter-server-main", "local to server main ID") + return None + for low, high in TG_LOCAL_TO_MCC: + if low <= dst <= high and int(source[:3]) != int(str(ctx.server_id)[:3]): + return _tg_filter(frame, "tg-filter-mcc", "local to MCC") + return None + + +def check_acl_chain(frame: ObpFrame, ctx: AdmissionContext) -> Rejection | None: + """Subscriber and talkgroup ACLs, GLOBAL first and then the system's own.""" + acl_check = ctx.acl_check + if acl_check is None: + return None + if ctx.global_rules.enabled: + if not acl_check(frame.rf_src, ctx.global_rules.sub_acl): + return Rejection( + reason="global-sub-acl", + message="(%s) CALL DROPPED WITH STREAM ID %s ON TGID %s BY GLOBAL TS1 ACL", + args=(frame.system, int_id(frame.stream_id), int_id(frame.rf_src)), + ) + if frame.slot == 1 and not acl_check(frame.dst_id, ctx.global_rules.tg1_acl): + return Rejection( + reason="global-tg1-acl", + message="(%s) CALL DROPPED WITH STREAM ID %s ON TGID %s BY GLOBAL TS1 ACL", + args=(frame.system, int_id(frame.stream_id), int_id(frame.dst_id)), + ) + if ctx.system_rules.enabled: + if not acl_check(frame.rf_src, ctx.system_rules.sub_acl): + return Rejection( + reason="system-sub-acl", + message="(%s) CALL DROPPED WITH STREAM ID %s FROM SUBSCRIBER %s BY SYSTEM ACL", + args=(frame.system, int_id(frame.stream_id), int_id(frame.rf_src)), + ) + if not acl_check(frame.dst_id, ctx.system_rules.tg1_acl): + return Rejection( + reason="system-tg1-acl", + message="(%s) CALL DROPPED WITH STREAM ID %s ON TGID %s BY SYSTEM ACL", + args=(frame.system, int_id(frame.stream_id), int_id(frame.dst_id)), + ) + return None + + +def admit_dmrd_v1(frame: ObpFrame, ctx: AdmissionContext) -> Rejection | None: + """Full DMRD v1 gauntlet, in the order the legacy handler applied it.""" + return ( + check_slot(frame) + or check_stun(frame, stunned=ctx.stunned) + or check_tg_filter_v1(frame) + or check_acl_chain(frame, ctx) + ) + + +def admit_dmre_v5( + frame: ObpFrame, + envelope: MeshEnvelope, + ctx: AdmissionContext, + *, + now: float, +) -> Rejection | None: + """Full DMRE v5 gauntlet, in the order the legacy handler applied it.""" + return ( + check_stun(frame, stunned=ctx.stunned) + or check_packet_age(frame, envelope, now=now) + or check_source_server(frame, envelope, ctx) + or check_hops(frame, envelope) + or check_tg_filter_v5(frame, envelope, ctx) + or check_acl_chain(frame, ctx) + ) + + +__all__ = [ + "AclRules", + "AdmissionContext", + "CallAttributes", + "MAX_HOPS", + "MAX_PACKET_AGE_S", + "MeshEnvelope", + "ObpFrame", + "Rejection", + "TG_DATA_GATEWAY", + "TG_LOCAL_TO_MCC", + "TG_LOCAL_TO_REPEATER_MAX", + "TG_LOCAL_TO_SERVER", + "TG_LOCAL_TO_SERVER_MAIN", + "admit_dmrd_v1", + "admit_dmre_v5", + "call_attributes", + "check_acl_chain", + "check_hops", + "check_network_id", + "check_packet_age", + "check_slot", + "check_source_server", + "check_stun", + "check_tg_filter_v1", + "check_tg_filter_v5", + "server_prefix", +] diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index b0f31cc..dc3167a 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -79,6 +79,18 @@ from ...domain import bytes_3, bytes_4, int_id from ...domain.dmr import decode from ...domain.dmr.const import LC_OPT from ...domain.hbp_protocol import normalize_fixed_width_ascii, normalize_fixed_width_bytes +from ...domain.mesh_admission import ( + AclRules, + AdmissionContext, + MeshEnvelope, + ObpFrame, + Rejection, + admit_dmrd_v1, + admit_dmre_v5, + call_attributes, + check_network_id, + server_prefix, +) from ...domain.mesh_routing import MeshEgress, MeshIngress, PeerMeshConfig from ...domain.talker_alias import ( DMRA_PACKET_LEN, @@ -2142,6 +2154,38 @@ class HBPProtocol(DatagramProtocol): self._config["TARGET_PORT"] = p self._config["TARGET_SOCK"] = (h, p) + def _obp_admission_context(self) -> AdmissionContext: + """Snapshot of everything the admission rules need, read once per frame.""" + _global = self._CONFIG.get("GLOBAL", {}) + return AdmissionContext( + stunned="STUN" in self._CONFIG, + acl_check=self._router.acl_check if self._router else None, + global_rules=AclRules( + enabled=bool(_global.get("USE_ACL")), + sub_acl=_global.get("SUB_ACL", (True, [])), + tg1_acl=_global.get("TG1_ACL", (True, [])), + ), + system_rules=AclRules( + enabled=bool(self._config.get("USE_ACL")), + sub_acl=self._config.get("SUB_ACL", (True, [])), + tg1_acl=self._config.get("TG1_ACL", (True, [])), + ), + server_id=server_prefix(_global.get("SERVER_ID", 0)), + validate_server_ids=bool(_global.get("VALIDATE_SERVER_IDS")), + known_server_prefixes=self._CONFIG.get("_SERVER_IDS", set()), + resolve_server_id=self.validate_obp_source_server_id, + ) + + def _obp_reject(self, rejection: Rejection, _dst_id: bytes, _stream_id: bytes) -> None: + """Apply one admission decision: log it (once per stream) and quench the peer.""" + if rejection.log_once: + if _stream_id in self._laststrid: + return + self._laststrid.append(_stream_id) + logger.log(rejection.level, rejection.message, *rejection.args) + if rejection.quench: + self._obp_send_bcsq(_dst_id, _stream_id) + def _obp_send_bcsq(self, _tgid: bytes, _stream_id: bytes) -> None: """Legacy send_bcsq: BCSQ + tgid + stream_id + HMAC-SHA1. Uses TARGET_SOCK (IP only).""" _addr = self._config.get("TARGET_SOCK") @@ -2178,67 +2222,35 @@ class HBPProtocol(DatagramProtocol): _data = _ingress.voice_frame self._obp_sync_target_sock_from_peer(_sockaddr, _stream_id) _peer_id = _data[11:15] - if self._config.get("NETWORK_ID") != _peer_id: - if _stream_id not in self._laststrid: - logger.error("(%s) OpenBridge packet discarded because NETWORK_ID: %s Does not match sent Peer ID: %s", self._system, int_id(self._config.get("NETWORK_ID", b"")), int_id(_peer_id)) - self._laststrid.append(_stream_id) + _dst_id = _data[8:11] + _rejection = check_network_id( + self._system, + _stream_id, + expected=self._config.get("NETWORK_ID", b""), + received=_peer_id, + ) + if _rejection is not None: + self._obp_reject(_rejection, _dst_id, _stream_id) return _seq = _data[4] _rf_src = _data[5:8] - _dst_id = _data[8:11] - _bits = _data[15] - _slot = 2 if (_bits & 0x80) else 1 - if _bits & 0x40: - _call_type = "unit" - elif (_bits & 0x23) == 0x23: - _call_type = "vcsbk" - else: - _call_type = "group" - _frame_type = (_bits & 0x30) >> 4 - _dtype_vseq = _bits & 0xF - if _slot != 1: - logger.error("(%s) OpenBridge packet discarded because it was not received on slot 1. SID: %s, TGID %s", self._system, int_id(_rf_src), int_id(_dst_id)) - return - if "STUN" in self._CONFIG: - if _stream_id not in self._laststrid: - logger.warning("(%s) Bridge STUNned, discarding", self._system) - self._laststrid.append(_stream_id) + _attrs = call_attributes(_data[15]) + _slot = _attrs.slot + _call_type = _attrs.call_type + _frame_type = _attrs.frame_type + _dtype_vseq = _attrs.dtype_vseq + _frame = ObpFrame( + system=self._system, + stream_id=_stream_id, + rf_src=_rf_src, + dst_id=_dst_id, + slot=_slot, + call_type=_call_type, + ) + _rejection = admit_dmrd_v1(_frame, self._obp_admission_context()) + if _rejection is not None: + self._obp_reject(_rejection, _dst_id, _stream_id) return - _int_dst_id = int_id(_dst_id) - if _call_type != "unit": - if _int_dst_id <= 79 or (_int_dst_id >= 9990 and _int_dst_id <= 9999) or (_int_dst_id >= 92 and _int_dst_id <= 199) or _int_dst_id == 900999: - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s FROM SUBSCRIBER %s BY GLOBAL TG FILTER", self._system, int_id(_stream_id), _int_dst_id) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - _global = self._CONFIG.get("GLOBAL", {}) - if self._router and _global.get("USE_ACL"): - if not self._router.acl_check(_rf_src, _global.get("SUB_ACL", (True, []))): - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s ON TGID %s BY GLOBAL TS1 ACL", self._system, int_id(_stream_id), int_id(_rf_src)) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - if _slot == 1 and not self._router.acl_check(_dst_id, _global.get("TG1_ACL", (True, []))): - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s ON TGID %s BY GLOBAL TS1 ACL", self._system, int_id(_stream_id), int_id(_dst_id)) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - if self._router and self._config.get("USE_ACL"): - if not self._router.acl_check(_rf_src, self._config.get("SUB_ACL", (True, []))): - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s FROM SUBSCRIBER %s BY SYSTEM ACL", self._system, int_id(_stream_id), int_id(_rf_src)) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - if not self._router.acl_check(_dst_id, self._config.get("TG1_ACL", (True, []))): - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s ON TGID %s BY SYSTEM ACL", self._system, int_id(_stream_id), int_id(_dst_id)) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return if _call_type == "group" and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD: logger.info( "(%s) CALL RX (OBP) src %s -> TG %s slot %s", @@ -2306,128 +2318,52 @@ class HBPProtocol(DatagramProtocol): _stream_id = _data[16:20] self._obp_sync_target_sock_from_peer(_sockaddr, _stream_id) _peer_id = _data[11:15] - if self._config.get("NETWORK_ID") != _peer_id: - if _stream_id not in self._laststrid: - logger.error("(%s) OpenBridge DMRE discarded because NETWORK_ID: %s Does not match sent Peer ID: %s", self._system, int_id(self._config.get("NETWORK_ID", b"")), int_id(_peer_id)) - self._laststrid.append(_stream_id) + _dst_id = _data[8:11] + _rejection = check_network_id( + self._system, + _stream_id, + expected=self._config.get("NETWORK_ID", b""), + received=_peer_id, + dmre=True, + ) + if _rejection is not None: + self._obp_reject(_rejection, _dst_id, _stream_id) return _seq = _data[4] _rf_src = _data[5:8] - _dst_id = _data[8:11] - _int_dst_id = int_id(_dst_id) - _bits = _data[15] - _slot = 2 if (_bits & 0x80) else 1 + _attrs = call_attributes(_data[15]) + _slot = _attrs.slot if self._config.get("MODE") == "OPENBRIDGE": # Legacy bridge_master: OpenBridge streams are effectively TS1 (DMRD v1 rejects slot != 1). # DMRE can still carry TS2 in bits; BRIDGES use TS:1 for OBP — normalize before STATUS/dmrd. _slot = 1 - if _bits & 0x40: - _call_type = "unit" - elif (_bits & 0x23) == 0x23: - _call_type = "vcsbk" - else: - _call_type = "group" - _frame_type = (_bits & 0x30) >> 4 - _dtype_vseq = _bits & 0xF - if "STUN" in self._CONFIG: - if _stream_id not in self._laststrid: - logger.warning("(%s) Bridge STUNned, discarding", self._system) - self._laststrid.append(_stream_id) - return - _ts_sec = int.from_bytes(_timestamp, "big") / 1_000_000_000 - if _ts_sec < (time.time() - 5): - if _stream_id not in self._laststrid: - logger.warning("(%s) Packet from server %s more than 5s old!, discarding", self._system, int.from_bytes(_source_server, "big")) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - _src_srv_int = int.from_bytes(_source_server, "big") - _src_srv_str = str(_src_srv_int) - _src_srv_len = len(_src_srv_str) - if _src_srv_len < 4 or _src_srv_len > 7: - if _stream_id not in self._laststrid: - logger.warning("(%s) Source Server should be between 4 and 7 digits, discarding Src: %s", self._system, _src_srv_int) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - _global = self._CONFIG.get("GLOBAL", {}) - _server_ids = self._CONFIG.get("_SERVER_IDS", set()) - if _global.get("VALIDATE_SERVER_IDS") and _src_srv_len in (4, 5) and (_src_srv_str[:4] not in _server_ids): - if _stream_id not in self._laststrid: - logger.warning("(%s) Source Server ID is 4 or 5 digits but not in list: %s", self._system, _src_srv_int) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - if _src_srv_len > 5 and not self.validate_obp_source_server_id(_source_server): - if _stream_id not in self._laststrid: - logger.warning("(%s) Source Server 6 or 7 digits but not a valid DMR ID, discarding Src: %s", self._system, _src_srv_int) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - _inthops = (_hops if isinstance(_hops, int) else int.from_bytes(_hops, "big")) + 1 - if _inthops > 10: - logger.debug( - "(%s) MAX HOPS exceed, dropping. Hops: %s, DST: %s, SRC: %s", - self._system, - _inthops, - _int_dst_id, - _src_srv_int, - ) - self._obp_send_bcsq(_dst_id, _stream_id) + _call_type = _attrs.call_type + _frame_type = _attrs.frame_type + _dtype_vseq = _attrs.dtype_vseq + _frame = ObpFrame( + system=self._system, + stream_id=_stream_id, + rf_src=_rf_src, + dst_id=_dst_id, + slot=_slot, + call_type=_call_type, + ) + _envelope = MeshEnvelope( + source_server=int.from_bytes(_source_server, "big"), + hops=_hops if isinstance(_hops, int) else int.from_bytes(_hops, "big"), + timestamp_ns=int.from_bytes(_timestamp, "big"), + source_server_id=_source_server, + ) + _rejection = admit_dmre_v5( + _frame, + _envelope, + self._obp_admission_context(), + now=time.time(), + ) + if _rejection is not None: + self._obp_reject(_rejection, _dst_id, _stream_id) return - if _call_type != "unit": - if _int_dst_id <= 79: - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s BY GLOBAL TG FILTER (local to repeater)", self._system, int_id(_stream_id), _int_dst_id) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - if (_int_dst_id >= 9990 and _int_dst_id <= 9999) or _int_dst_id == 900999: - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s BY GLOBAL TG FILTER (local to server)", self._system, int_id(_stream_id), _int_dst_id) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - _sid = _global.get("SERVER_ID", 0) - _our_srv = int(str(_sid)[:4]) if isinstance(_sid, int) else int(str(int.from_bytes(_sid, "big"))[:4]) - if (_int_dst_id >= 92 and _int_dst_id <= 199) and int(_src_srv_str[:4]) != _our_srv: - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s BY GLOBAL TG FILTER (local to server main ID)", self._system, int_id(_stream_id), _int_dst_id) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - if ((_int_dst_id >= 80 and _int_dst_id <= 89) or (_int_dst_id >= 800 and _int_dst_id <= 899)) and int(_src_srv_str[:3]) != int(str(_our_srv)[:3]): - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s BY GLOBAL TG FILTER (local to MCC)", self._system, int_id(_stream_id), _int_dst_id) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - if _global.get("USE_ACL") and self._router: - if not self._router.acl_check(_rf_src, _global.get("SUB_ACL", (True, []))): - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s ON TGID %s BY GLOBAL TS1 ACL", self._system, int_id(_stream_id), int_id(_rf_src)) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - if _slot == 1 and not self._router.acl_check(_dst_id, _global.get("TG1_ACL", (True, []))): - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s ON TGID %s BY GLOBAL TS1 ACL", self._system, int_id(_stream_id), int_id(_dst_id)) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - if self._config.get("USE_ACL") and self._router: - if not self._router.acl_check(_rf_src, self._config.get("SUB_ACL", (True, []))): - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s FROM SUBSCRIBER %s BY SYSTEM ACL", self._system, int_id(_stream_id), int_id(_rf_src)) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return - if not self._router.acl_check(_dst_id, self._config.get("TG1_ACL", (True, []))): - if _stream_id not in self._laststrid: - logger.info("(%s) CALL DROPPED WITH STREAM ID %s ON TGID %s BY SYSTEM ACL", self._system, int_id(_stream_id), int_id(_dst_id)) - self._obp_send_bcsq(_dst_id, _stream_id) - self._laststrid.append(_stream_id) - return + _inthops = _envelope.hops + 1 self.note_dmrd_stream(_peer_id, _rf_src, _stream_id) if ( _call_type in ("group", "vcsbk") diff --git a/tests/domain/test_mesh_admission.py b/tests/domain/test_mesh_admission.py new file mode 100644 index 0000000..ead8857 --- /dev/null +++ b/tests/domain/test_mesh_admission.py @@ -0,0 +1,393 @@ +# ADN DMR Peer Server - tests domain mesh admission +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""OpenBridge admission rules: one frame in, one decision out.""" + +from __future__ import annotations + +import logging + +import pytest + +from adn_server.domain import bytes_3, bytes_4 +from adn_server.domain.mesh_admission import ( + AclRules, + AdmissionContext, + MeshEnvelope, + ObpFrame, + admit_dmrd_v1, + admit_dmre_v5, + call_attributes, + check_acl_chain, + check_hops, + check_network_id, + check_packet_age, + check_slot, + check_source_server, + check_stun, + check_tg_filter_v1, + check_tg_filter_v5, + server_prefix, +) + +_SYSTEM = "OBP-FR" +_STREAM = bytes_4(0xAABBCCDD) +_SRC = bytes_3(2130001) +_NOW = 1_800_000_000.0 + + +def _frame(dst: int = 214, *, slot: int = 1, call_type: str = "group") -> ObpFrame: + return ObpFrame( + system=_SYSTEM, + stream_id=_STREAM, + rf_src=_SRC, + dst_id=bytes_3(dst), + slot=slot, + call_type=call_type, + ) + + +def _envelope(*, source_server: int = 20840, hops: int = 0, fresh: bool = True) -> MeshEnvelope: + return MeshEnvelope( + source_server=source_server, + hops=hops, + timestamp_ns=int(_NOW * 1_000_000_000) if fresh else 0, + source_server_id=bytes_4(source_server), + ) + + +def _ctx(**kwargs) -> AdmissionContext: + return AdmissionContext(**kwargs) + + +# --- bits byte --------------------------------------------------------------- + + +@pytest.mark.parametrize( + ("bits", "slot", "call_type"), + [ + (0x00, 1, "group"), + (0x80, 2, "group"), + (0x40, 1, "unit"), + (0xC0, 2, "unit"), + (0x23, 1, "vcsbk"), + (0xE3, 2, "unit"), # the private bit wins over the CSBK pattern + ], +) +def test_call_attributes_classifies_the_bits_byte(bits: int, slot: int, call_type: str) -> None: + attrs = call_attributes(bits) + assert (attrs.slot, attrs.call_type) == (slot, call_type) + + +def test_call_attributes_splits_frame_type_and_sequence() -> None: + attrs = call_attributes(0x16) + assert attrs.frame_type == 1 + assert attrs.dtype_vseq == 6 + + +# --- identity ---------------------------------------------------------------- + + +def test_network_id_match_is_admitted() -> None: + assert check_network_id(_SYSTEM, _STREAM, expected=bytes_4(20840), received=bytes_4(20840)) is None + + +def test_network_id_mismatch_is_logged_but_not_quenched() -> None: + rejection = check_network_id(_SYSTEM, _STREAM, expected=bytes_4(20840), received=bytes_4(26811)) + assert rejection is not None + assert rejection.reason == "network-id-mismatch" + assert rejection.level == logging.ERROR + assert rejection.quench is False + assert "OpenBridge packet discarded" in rejection.message + + +def test_network_id_mismatch_names_dmre_when_asked() -> None: + rejection = check_network_id( + _SYSTEM, _STREAM, expected=bytes_4(20840), received=bytes_4(26811), dmre=True + ) + assert rejection is not None + assert "OpenBridge DMRE discarded" in rejection.message + + +# --- slot, stun -------------------------------------------------------------- + + +def test_slot_1_is_admitted_and_slot_2_is_not() -> None: + assert check_slot(_frame()) is None + rejection = check_slot(_frame(slot=2)) + assert rejection is not None + assert rejection.reason == "not-slot-1" + assert rejection.quench is False + assert rejection.log_once is False # legacy logs this one on every frame + + +def test_stunned_bridge_drops_without_quenching() -> None: + assert check_stun(_frame(), stunned=False) is None + rejection = check_stun(_frame(), stunned=True) + assert rejection is not None + assert rejection.reason == "stunned" + assert rejection.quench is False + + +# --- DMRE envelope ----------------------------------------------------------- + + +def test_fresh_packet_passes_and_old_one_is_dropped() -> None: + assert check_packet_age(_frame(), _envelope(), now=_NOW) is None + rejection = check_packet_age(_frame(), _envelope(), now=_NOW + 6) + assert rejection is not None + assert rejection.reason == "stale-packet" + + +def test_packet_without_timestamp_counts_as_stale() -> None: + """A DMRE frame with no trailer has no timestamp, and legacy drops it.""" + rejection = check_packet_age(_frame(), _envelope(fresh=False), now=_NOW) + assert rejection is not None + assert rejection.reason == "stale-packet" + + +@pytest.mark.parametrize("source_server", [123, 12345678]) +def test_source_server_must_be_4_to_7_digits(source_server: int) -> None: + rejection = check_source_server(_frame(), _envelope(source_server=source_server), _ctx()) + assert rejection is not None + assert rejection.reason == "source-server-length" + + +def test_short_source_server_must_be_a_known_server() -> None: + envelope = _envelope(source_server=2084) + ctx = _ctx(validate_server_ids=True, known_server_prefixes={"2131"}) + rejection = check_source_server(_frame(), envelope, ctx) + assert rejection is not None + assert rejection.reason == "source-server-unknown" + ctx_known = _ctx(validate_server_ids=True, known_server_prefixes={"2084"}) + assert check_source_server(_frame(), envelope, ctx_known) is None + + +def test_unknown_short_source_server_passes_when_validation_is_off() -> None: + envelope = _envelope(source_server=2084) + assert check_source_server(_frame(), envelope, _ctx(known_server_prefixes={"2131"})) is None + + +def test_long_source_server_must_resolve_to_a_dmr_id() -> None: + envelope = _envelope(source_server=2130001) + rejection = check_source_server(_frame(), envelope, _ctx(resolve_server_id=lambda _id: False)) + assert rejection is not None + assert rejection.reason == "source-server-invalid" + assert check_source_server(_frame(), envelope, _ctx(resolve_server_id=lambda _id: "C31AG")) is None + + +def test_long_source_server_is_looked_up_by_its_wire_bytes() -> None: + seen: list[bytes] = [] + envelope = _envelope(source_server=2130001) + check_source_server(_frame(), envelope, _ctx(resolve_server_id=lambda sid: seen.append(sid) or True)) + assert seen == [bytes_4(2130001)] + + +def test_hops_are_counted_and_capped() -> None: + assert check_hops(_frame(), _envelope(hops=8)) is None + rejection = check_hops(_frame(), _envelope(hops=10)) + assert rejection is not None + assert rejection.reason == "max-hops" + assert rejection.level == logging.DEBUG + assert rejection.log_once is False # legacy quenches every looping frame + + +# --- talkgroup filters ------------------------------------------------------- + + +@pytest.mark.parametrize("dst", [9, 79, 92, 199, 9990, 9999, 900999]) +def test_tg_filter_v1_keeps_local_talkgroups_off_the_mesh(dst: int) -> None: + rejection = check_tg_filter_v1(_frame(dst)) + assert rejection is not None + assert rejection.reason == "tg-filter" + assert rejection.quench is True + + +@pytest.mark.parametrize("dst", [80, 91, 200, 214, 9989, 10000]) +def test_tg_filter_v1_lets_ordinary_talkgroups_through(dst: int) -> None: + assert check_tg_filter_v1(_frame(dst)) is None + + +def test_tg_filter_v1_ignores_private_calls() -> None: + assert check_tg_filter_v1(_frame(9, call_type="unit")) is None + + +@pytest.mark.parametrize( + ("dst", "reason"), + [ + (9, "tg-filter-repeater"), + (79, "tg-filter-repeater"), + (9990, "tg-filter-server"), + (900999, "tg-filter-server"), + ], +) +def test_tg_filter_v5_drops_talkgroups_that_never_leave_home(dst: int, reason: str) -> None: + rejection = check_tg_filter_v5(_frame(dst), _envelope(), _ctx(server_id=2131)) + assert rejection is not None + assert rejection.reason == reason + + +def test_tg_filter_v5_allows_a_server_local_tg_from_that_server() -> None: + """92-199 belong to one server: only that server may bridge them.""" + ctx = _ctx(server_id=2084) + assert check_tg_filter_v5(_frame(100), _envelope(source_server=20840), ctx) is None + rejection = check_tg_filter_v5(_frame(100), _envelope(source_server=21310), ctx) + assert rejection is not None + assert rejection.reason == "tg-filter-server-main" + + +def test_tg_filter_v5_allows_an_mcc_tg_from_the_same_mcc() -> None: + ctx = _ctx(server_id=2084) + assert check_tg_filter_v5(_frame(85), _envelope(source_server=20851), ctx) is None + rejection = check_tg_filter_v5(_frame(850), _envelope(source_server=21310), ctx) + assert rejection is not None + assert rejection.reason == "tg-filter-mcc" + + +def test_tg_filter_v5_lets_ordinary_talkgroups_through() -> None: + assert check_tg_filter_v5(_frame(214), _envelope(), _ctx(server_id=2131)) is None + + +def test_tg_filter_v5_ignores_private_calls() -> None: + assert check_tg_filter_v5(_frame(9, call_type="unit"), _envelope(), _ctx(server_id=2131)) is None + + +@pytest.mark.parametrize(("value", "prefix"), [(21310, 2131), (bytes_4(21310), 2131), (0, 0), (None, 0)]) +def test_server_prefix_reads_int_or_bytes(value: object, prefix: int) -> None: + assert server_prefix(value) == prefix + + +# --- ACLs -------------------------------------------------------------------- + + +def _deny(*denied: bytes): + def _check(target: bytes, _acl: object) -> bool: + return target not in denied + + return _check + + +def test_acl_chain_is_skipped_without_a_router() -> None: + ctx = _ctx(global_rules=AclRules(enabled=True), system_rules=AclRules(enabled=True)) + assert check_acl_chain(_frame(), ctx) is None + + +def test_acl_chain_admits_when_every_rule_passes() -> None: + ctx = _ctx( + acl_check=_deny(), + global_rules=AclRules(enabled=True), + system_rules=AclRules(enabled=True), + ) + assert check_acl_chain(_frame(), ctx) is None + + +@pytest.mark.parametrize( + ("denied", "scope", "reason"), + [ + (_SRC, "global", "global-sub-acl"), + (bytes_3(214), "global", "global-tg1-acl"), + (_SRC, "system", "system-sub-acl"), + (bytes_3(214), "system", "system-tg1-acl"), + ], +) +def test_acl_chain_reports_which_rule_dropped_the_call(denied: bytes, scope: str, reason: str) -> None: + rules = AclRules(enabled=True) + ctx = _ctx( + acl_check=_deny(denied), + global_rules=rules if scope == "global" else AclRules(), + system_rules=rules if scope == "system" else AclRules(), + ) + rejection = check_acl_chain(_frame(), ctx) + assert rejection is not None + assert rejection.reason == reason + assert rejection.quench is True + + +def test_global_talkgroup_acl_only_applies_to_slot_1() -> None: + ctx = _ctx(acl_check=_deny(bytes_3(214)), global_rules=AclRules(enabled=True)) + assert check_acl_chain(_frame(slot=2), ctx) is None + + +# --- full gauntlets ---------------------------------------------------------- + + +def test_dmrd_v1_admits_an_ordinary_group_call() -> None: + ctx = _ctx(acl_check=_deny(), global_rules=AclRules(enabled=True)) + assert admit_dmrd_v1(_frame(214), ctx) is None + + +def test_dmrd_v1_reports_the_first_rule_that_fails() -> None: + """Order matters: a STUNned bridge is reported as such, not as a TG drop.""" + ctx = _ctx(stunned=True, acl_check=_deny(_SRC), global_rules=AclRules(enabled=True)) + rejection = admit_dmrd_v1(_frame(9), ctx) + assert rejection is not None + assert rejection.reason == "stunned" + + +def test_dmre_v5_admits_an_ordinary_group_call() -> None: + ctx = _ctx(server_id=2131, acl_check=_deny(), system_rules=AclRules(enabled=True)) + assert admit_dmre_v5(_frame(214), _envelope(), ctx, now=_NOW) is None + + +def test_dmre_v5_checks_the_envelope_before_the_talkgroup() -> None: + ctx = _ctx(server_id=2131) + rejection = admit_dmre_v5(_frame(9), _envelope(fresh=False), ctx, now=_NOW) + assert rejection is not None + assert rejection.reason == "stale-packet" + + +def test_every_rejection_can_be_formatted() -> None: + """A log line with the wrong number of placeholders logs an error instead of + the drop, so each rule's message and args are checked against each other.""" + ctx_deny = _ctx( + stunned=False, + acl_check=_deny(_SRC, bytes_3(214)), + global_rules=AclRules(enabled=True), + system_rules=AclRules(enabled=True), + server_id=2084, + validate_server_ids=True, + known_server_prefixes={"2131"}, + resolve_server_id=lambda _id: False, + ) + rejections = [ + check_network_id(_SYSTEM, _STREAM, expected=bytes_4(1), received=bytes_4(2)), + check_network_id(_SYSTEM, _STREAM, expected=bytes_4(1), received=bytes_4(2), dmre=True), + check_slot(_frame(slot=2)), + check_stun(_frame(), stunned=True), + check_packet_age(_frame(), _envelope(fresh=False), now=_NOW), + check_source_server(_frame(), _envelope(source_server=123), _ctx()), + check_source_server(_frame(), _envelope(source_server=2084), ctx_deny), + check_source_server(_frame(), _envelope(source_server=2130001), ctx_deny), + check_hops(_frame(), _envelope(hops=10)), + check_tg_filter_v1(_frame(9)), + check_tg_filter_v5(_frame(9), _envelope(), ctx_deny), + check_tg_filter_v5(_frame(9990), _envelope(), ctx_deny), + check_tg_filter_v5(_frame(100), _envelope(source_server=21310), ctx_deny), + check_tg_filter_v5(_frame(85), _envelope(source_server=21310), ctx_deny), + check_acl_chain(_frame(), ctx_deny), + check_acl_chain(_frame(), _ctx(acl_check=_deny(_SRC), system_rules=AclRules(enabled=True))), + ] + assert all(r is not None for r in rejections) + seen = set() + for rejection in rejections: + assert rejection is not None + rejection.message % rejection.args # raises if the arity is wrong + seen.add(rejection.reason) + assert len(seen) == len(rejections) - 1 # the two network-id variants share a reason