From ad9ab2a024f6792241a866be24ab8df0200d50b8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Thu, 9 Jul 2026 02:40:48 -0400 Subject: [PATCH] fix: raise UDP SO_RCVBUF on voice listeners (C-LOCAL) Apply a 4 MB receive buffer on system and proxy UDP sockets to reduce kernel RcvbufErrors under OBP load; size is configurable via GLOBAL.UDP_RCVBUF. --- adn-server.example.yaml | 2 + .../infrastructure/bootstrap/peer_server.py | 2 + .../infrastructure/config_validator.py | 2 +- .../infrastructure/proxy/runtime.py | 1 + .../infrastructure/proxy/udp_fanin.py | 8 ++- src/adn_server/infrastructure/udp_rcvbuf.py | 60 +++++++++++++++++++ tests/infrastructure/test_udp_rcvbuf.py | 48 +++++++++++++++ 7 files changed, 121 insertions(+), 2 deletions(-) create mode 100644 src/adn_server/infrastructure/udp_rcvbuf.py create mode 100644 tests/infrastructure/test_udp_rcvbuf.py diff --git a/adn-server.example.yaml b/adn-server.example.yaml index 559072e..953ed2c 100644 --- a/adn-server.example.yaml +++ b/adn-server.example.yaml @@ -23,6 +23,8 @@ GLOBAL: TALKER_ALIAS_FORMAT: "{callsign} {fname}" # utf8 (Motorola), iso8 (Hytera), 7bit; both vendors by default TALKER_ALIAS_TEXT_FORMAT: "utf8,iso8" + # UDP SO_RCVBUF for voice listeners (bytes); default 4 MB avoids kernel RcvbufErrors under load + # UDP_RCVBUF: 4194304 REPORTS: REPORT: true diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index 994b6fe..d4436d5 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -44,6 +44,7 @@ from adn_server.application.proxy.deployment import ( proxy_target_system, ) from adn_server.application.report.queue import BoundedReportQueue, QueuedReportSender +from adn_server.infrastructure.udp_rcvbuf import apply_udp_rcvbuf, udp_rcvbuf_bytes from adn_server.application.runtime_context import ( ConfigProxy, RuntimeContext, @@ -583,6 +584,7 @@ def run_peer_server( def _listen_system(_name: str, bind: BindSpec, protocol: Any) -> Any: port = reactor.listenUDP(bind.port, protocol, interface=bind.ip or "0.0.0.0") + apply_udp_rcvbuf(port.socket, udp_rcvbuf_bytes(config), label=_name, logger=logger) logger.info("(GLOBAL) UDP %s listening on %s:%s", _name, bind.ip or "*", bind.port) return port diff --git a/src/adn_server/infrastructure/config_validator.py b/src/adn_server/infrastructure/config_validator.py index 4c2fb7d..13e4728 100644 --- a/src/adn_server/infrastructure/config_validator.py +++ b/src/adn_server/infrastructure/config_validator.py @@ -131,7 +131,7 @@ def _section_string_keys(section_name: str, section: dict[str, Any], keys: froze def _validate_global(global_cfg: dict[str, Any], errors: list[str]) -> None: _section_string_keys("GLOBAL", global_cfg, GLOBAL_STRING_KEYS, errors) - for key in ("PING_TIME", "MAX_MISSED", "SERVER_ID"): + for key in ("PING_TIME", "MAX_MISSED", "SERVER_ID", "UDP_RCVBUF"): if key in global_cfg: _expect_int(f"GLOBAL.{key}", global_cfg[key], errors) for key in ( diff --git a/src/adn_server/infrastructure/proxy/runtime.py b/src/adn_server/infrastructure/proxy/runtime.py index 3effd4c..8d4c713 100644 --- a/src/adn_server/infrastructure/proxy/runtime.py +++ b/src/adn_server/infrastructure/proxy/runtime.py @@ -282,6 +282,7 @@ def start_proxy_service( debug=runtime["debug"], logger=logger, protocol=fanin, + config=config, ) state.udp_port = udp_port state.client_sender = FanInClientSender(fanin_proto.transport) diff --git a/src/adn_server/infrastructure/proxy/udp_fanin.py b/src/adn_server/infrastructure/proxy/udp_fanin.py index 90622ae..d5ef44b 100644 --- a/src/adn_server/infrastructure/proxy/udp_fanin.py +++ b/src/adn_server/infrastructure/proxy/udp_fanin.py @@ -32,6 +32,7 @@ from adn_server.application.ports import ProxyMasterSink from adn_server.application.proxy import ProxyUseCases, peer_id_from_packet from adn_server.domain.result import is_fail from adn_server.infrastructure.hbp_constants import RPTC, RPTO +from adn_server.infrastructure.udp_rcvbuf import apply_udp_rcvbuf, udp_rcvbuf_bytes if TYPE_CHECKING: from .self_service_bridge import ProxySelfServiceBridge @@ -118,10 +119,15 @@ def listen_proxy_fanin( debug: bool = False, logger: logging.Logger | None = None, protocol: ProxyFanInProtocol | None = None, + config: dict[str, Any] | None = None, + udp_rcvbuf: int | None = None, ) -> tuple[ProxyFanInProtocol, Any]: """Bind LISTEN_PORT and return ``(protocol, udp_port)``.""" - fanin = protocol or ProxyFanInProtocol(proxy, master_sink, debug=debug, logger=logger) + log = logger or _logger + fanin = protocol or ProxyFanInProtocol(proxy, master_sink, debug=debug, logger=log) udp_port = reactor.listenUDP(listen_port, fanin, interface=listen_ip or "0.0.0.0") + buf_size = udp_rcvbuf if udp_rcvbuf is not None else udp_rcvbuf_bytes(config) + apply_udp_rcvbuf(udp_port.socket, buf_size, label="PROXY", logger=log) return fanin, udp_port diff --git a/src/adn_server/infrastructure/udp_rcvbuf.py b/src/adn_server/infrastructure/udp_rcvbuf.py new file mode 100644 index 0000000..016730a --- /dev/null +++ b/src/adn_server/infrastructure/udp_rcvbuf.py @@ -0,0 +1,60 @@ +# ADN DMR Peer Server - UDP receive buffer sizing +# 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 +############################################################################### + +"""Raise SO_RCVBUF on voice UDP listeners to avoid kernel RcvbufErrors under load.""" + +from __future__ import annotations + +import logging +import socket +from typing import Any + +DEFAULT_UDP_RCVBUF = 4 * 1024 * 1024 + + +def udp_rcvbuf_bytes(config: dict[str, Any] | None) -> int: + if not config: + return DEFAULT_UDP_RCVBUF + raw = config.get("GLOBAL", {}).get("UDP_RCVBUF", DEFAULT_UDP_RCVBUF) + if isinstance(raw, bool) or not isinstance(raw, int) or raw <= 0: + return DEFAULT_UDP_RCVBUF + return raw + + +def apply_udp_rcvbuf( + sock: socket.socket, + requested: int, + *, + label: str, + logger: logging.Logger, +) -> None: + try: + sock.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, requested) + except OSError as exc: + logger.warning("(%s) UDP RX buffer not raised (requested %s): %s", label, requested, exc) + return + try: + effective = sock.getsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF) + except OSError as exc: + logger.warning("(%s) UDP RX buffer set but getsockopt failed: %s", label, exc) + return + logger.info("(%s) UDP RX buffer raised to %s bytes (requested %s)", label, effective, requested) + + +__all__ = ["DEFAULT_UDP_RCVBUF", "apply_udp_rcvbuf", "udp_rcvbuf_bytes"] diff --git a/tests/infrastructure/test_udp_rcvbuf.py b/tests/infrastructure/test_udp_rcvbuf.py new file mode 100644 index 0000000..c28b8d6 --- /dev/null +++ b/tests/infrastructure/test_udp_rcvbuf.py @@ -0,0 +1,48 @@ +"""Tests for UDP receive buffer helper.""" + +from __future__ import annotations + +import logging +import socket +from unittest.mock import MagicMock + +from adn_server.infrastructure.udp_rcvbuf import ( + DEFAULT_UDP_RCVBUF, + apply_udp_rcvbuf, + udp_rcvbuf_bytes, +) + + +def test_udp_rcvbuf_bytes_default() -> None: + assert udp_rcvbuf_bytes(None) == DEFAULT_UDP_RCVBUF + assert udp_rcvbuf_bytes({}) == DEFAULT_UDP_RCVBUF + assert udp_rcvbuf_bytes({"GLOBAL": {}}) == DEFAULT_UDP_RCVBUF + + +def test_udp_rcvbuf_bytes_from_config() -> None: + assert udp_rcvbuf_bytes({"GLOBAL": {"UDP_RCVBUF": 2097152}}) == 2097152 + + +def test_udp_rcvbuf_bytes_rejects_invalid() -> None: + assert udp_rcvbuf_bytes({"GLOBAL": {"UDP_RCVBUF": 0}}) == DEFAULT_UDP_RCVBUF + assert udp_rcvbuf_bytes({"GLOBAL": {"UDP_RCVBUF": -1}}) == DEFAULT_UDP_RCVBUF + assert udp_rcvbuf_bytes({"GLOBAL": {"UDP_RCVBUF": "big"}}) == DEFAULT_UDP_RCVBUF + assert udp_rcvbuf_bytes({"GLOBAL": {"UDP_RCVBUF": True}}) == DEFAULT_UDP_RCVBUF + + +def test_apply_udp_rcvbuf_sets_socket_buffer() -> None: + sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) + try: + requested = 2 * 1024 * 1024 + log = MagicMock(spec=logging.Logger) + apply_udp_rcvbuf(sock, requested, label="TEST", logger=log) + effective = sock.getsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF) + assert effective >= requested + log.info.assert_called_once() + args = log.info.call_args[0] + assert args[0] == "(%s) UDP RX buffer raised to %s bytes (requested %s)" + assert args[1] == "TEST" + assert args[2] == effective + assert args[3] == requested + finally: + sock.close()