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.
pull/34/head
Rodrigo Pérez 3 months ago
parent 2e4026f9d7
commit ad9ab2a024

@ -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

@ -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

@ -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 (

@ -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)

@ -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

@ -0,0 +1,60 @@
# ADN DMR Peer Server - UDP receive buffer sizing
# 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, 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"]

@ -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()
Loading…
Cancel
Save

Powered by TurnKey Linux.