From 156731447c62cd9fe5f31562900cfba2d6c5d7ca Mon Sep 17 00:00:00 2001
From: ce5rpy <169016246+ce5rpy@users.noreply.github.com>
Date: Thu, 9 Jul 2026 13:33:55 -0400
Subject: [PATCH] fix: production ingress mitigations (UDP rcvbuf, rate
limiter, TA embed) (#34)
* 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.
* fix: exclude byte-identical duplicates from HBP rate counter (A.2)
Check lastData before incrementing the ingress packet counter so compressed
duplicate bursts do not trigger legitimate RATE DROP on call start.
* fix: duplicate-safe TA embed phase and REPEAT VHEAD DMRA (B)
Ignore byte-identical B-E embed bursts so duplicate uplinks do not desync the
TA phase machine, and re-emit DMRA on every VHEAD on the REPEAT path.
* chore: document echo point-to-point and logged_in reconciliation (EN/ES)
Document multi-hotspot echo/service delivery via RX_PEER, the lst_seen
logged_in reconcile loop, and cross-links between user and dev guides.
---
adn-server.example.yaml | 2 +
docs/en/monitor/self-service.md | 18 ++
.../development/behaviour-and-timers.md | 1 +
.../development/routing-and-contention.md | 20 +-
docs/en/server/user-guide/echo.md | 20 ++
docs/en/server/user-guide/hotspot-proxy.md | 2 +-
docs/es/monitor/self-service.md | 19 ++
.../development/behaviour-and-timers.md | 1 +
.../development/routing-and-contention.md | 20 +-
docs/es/server/user-guide/echo.md | 21 ++
docs/es/server/user-guide/hotspot-proxy.md | 2 +-
.../application/routing/hbp_forward.py | 14 +-
src/adn_server/application/routing/lc_ta.py | 21 +-
.../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/hbp/test_ingress.py | 23 ++
tests/infrastructure/test_udp_rcvbuf.py | 48 +++++
.../test_embed_phase_duplicate.py | 204 ++++++++++++++++++
21 files changed, 494 insertions(+), 15 deletions(-)
create mode 100644 src/adn_server/infrastructure/udp_rcvbuf.py
create mode 100644 tests/infrastructure/test_udp_rcvbuf.py
create mode 100644 tests/talker_alias/test_embed_phase_duplicate.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/docs/en/monitor/self-service.md b/docs/en/monitor/self-service.md
index ddefe72..b88c355 100644
--- a/docs/en/monitor/self-service.md
+++ b/docs/en/monitor/self-service.md
@@ -75,6 +75,24 @@ Important: the proxy sends **RPTO to the master only**, not to the hotspot direc
---
+## `logged_in` reconciliation
+
+The **`logged_in`** flag gates both **password** and **IP-based** login: only
+rows with `logged_in = 1` can authenticate on the dashboard. The **peer
+server** keeps this flag accurate by reconciling it against actually connected
+peers every **120 s** (the `lst_seen` loop):
+
+- Peers currently connected to the inject-only MASTER get `logged_in = 1`.
+- All other rows get `logged_in = 0`.
+- The loop starts with `now=True`, so the **first tick at boot clears every
+ stale flag immediately** — after a server restart, hotspots that did not
+ reconnect cannot authenticate via login-by-IP.
+
+This replaced the legacy hourly `clean_tbl` (24 h idle sweep), which left
+`logged_in = 1` for peers no longer connected after a restart.
+
+---
+
## Password hashing
**`AuthenticateUser`** uses:
diff --git a/docs/en/server/development/behaviour-and-timers.md b/docs/en/server/development/behaviour-and-timers.md
index 1b73668..4c6aacf 100644
--- a/docs/en/server/development/behaviour-and-timers.md
+++ b/docs/en/server/development/behaviour-and-timers.md
@@ -21,6 +21,7 @@ The following intervals are part of the current runtime behavior:
| `bridge_reset` | **6s** | Bridge reset flag cleanup and pending reset completion. |
| OPTIONS refresh | **event-driven** | Static TG / reflector from **RPTO**, **startup/reload** (`apply_startup_bridges`), **dmrd** no-source fallback. No periodic 26s loop (**D-28**). |
| `dynamic_tg_purge_loop` | **60s** | Purge expired **SINGLE=1** rows from `peer_dynamic_tgs` and in-memory `_PEER_UA_SESSIONS`. |
+| `lst_seen` (self-service reconcile) | **120s** | Reconcile `Clients.logged_in` against currently connected peers: connected peers get `logged_in=1`, the rest `0`. Runs with `now=True` on first tick so stale flags clear immediately after a server restart. Prevents the monitor from authenticating disconnected hotspots via login-by-IP. |
| `statTrimmer` | **303s** | Trim stale STAT bridges and transient status entries. |
If you change one of these intervals, document the operational impact for monitoring, loop behavior, and troubleshooting.
diff --git a/docs/en/server/development/routing-and-contention.md b/docs/en/server/development/routing-and-contention.md
index 89eca0c..110afc0 100644
--- a/docs/en/server/development/routing-and-contention.md
+++ b/docs/en/server/development/routing-and-contention.md
@@ -285,7 +285,7 @@ flowchart TD
PKT[DMRD toward peer] --> F1{OPTIONS has TG
or dynamic active?}
F1 -->|No| DROP[Do not deliver]
F1 -->|Yes| F2{Special TG 9990-9999?}
- F2 -->|Yes| PASS1[Bypass slot contention]
+ F2 -->|Yes| PASS1[Point-to-point:
only RX_PEER (or single peer)]
F2 -->|No| F3{Target slot
free?}
F3 -->|Peer TXing ingress| DROP
F3 -->|Bridge hold active
different TG| DROP
@@ -309,6 +309,24 @@ flowchart TD
| **Incompatible SINGLE lock** | SINGLE=1 with a lock on another TG → block, unless it is the lock TG or the peer activated it. |
| **Stale session** | If the per-peer session has had no frames for `_STALE_PEER_SESSION_TIMEOUT`, it is purged and the slot frees. |
+### Special TGs (9990–9999): point-to-point delivery
+
+Echo (9990) and on-demand service TGs (9991–9999) **bypass** the per-peer slot
+contention and OPTIONS filters described above, but they are **not** broadcast
+to all peers. On a multi-peer inject-only MASTER they are delivered
+**point-to-point**:
+
+- The packet goes **only** to `RX_PEER` — the exact peer that originated the
+ call on that slot — when `RX_TGID` matches the special TG.
+- If `RX_TGID` does not yet match (e.g. the very first VHEAD of a 9990
+ transmission arrives before `RX_TGID` is updated) and only one peer is
+ connected, it is delivered to that single peer (legacy fallback).
+- There is **no fuzzy matching** on the source DMR ID: other hotspots of the
+ same user never receive echo or service playback.
+
+This matches legacy `adn-dmr-server`, where each MASTER had a single peer so
+echo naturally returned only to the caller. See [Echo](../user-guide/echo.md).
+
### Mid-call join
When a downlink stream ends (VTERM or timeout), the peer's slot becomes free.
diff --git a/docs/en/server/user-guide/echo.md b/docs/en/server/user-guide/echo.md
index 123f486..e4bde97 100644
--- a/docs/en/server/user-guide/echo.md
+++ b/docs/en/server/user-guide/echo.md
@@ -37,6 +37,26 @@ sudo systemctl enable --now adn-echo
The main server exposes an **ECHO** master on **TG 9990** for the echo bridge. The standalone **echo** service is a separate process with its own config that connects to that master.
+## Multi-hotspot behaviour (inject-only proxy)
+
+When hotspots attach through the **integrated proxy** (`PROXY`), several radios
+may share the same MASTER as peers. In legacy `adn-dmr-server` each MASTER had
+a single peer, so echo naturally returned only to the caller. The multi-peer
+inject-only proxy enforces the same explicitly:
+
+- **Point-to-point delivery.** Echo playback (TG 9990) and on-demand service
+ TGs (**9991–9999**) are delivered **only** to the exact peer that originated
+ the call (`RX_PEER` on the active slot), **never** to other hotspots of the
+ same user. There is no fuzzy matching on the source DMR ID.
+- **Single-peer fallback.** When only one peer is connected, the packet is
+ delivered to it (legacy single-peer behaviour).
+- This applies to both the **data plane** (audio routing) and the **report
+ plane** (`BRDG_EVENT` sent to the monitor): the monitor shows the echo chip
+ on the originating hotspot, not on a sibling.
+
+See [Voice routing and contention — Downlink gate](../development/routing-and-contention.md#downlink-gate-does-a-peer-receive-the-packet)
+and [Hotspot proxy — Multi-hotspot behaviour](hotspot-proxy.md#multi-hotspot-behaviour).
+
## Documentation
This page is the summary shipped with the repository; extend your deployment notes locally as needed.
diff --git a/docs/en/server/user-guide/hotspot-proxy.md b/docs/en/server/user-guide/hotspot-proxy.md
index 1e4361c..63adcfe 100644
--- a/docs/en/server/user-guide/hotspot-proxy.md
+++ b/docs/en/server/user-guide/hotspot-proxy.md
@@ -92,7 +92,7 @@ Details of the dashboard flow: [Self-service](../../monitor/self-service.md).
## Multi-hotspot behaviour
- Each authenticated hotspot is a **peer** on the inject-only MASTER with its own **OPTIONS** (static TGs). **Repeat** and monitor fan-out respect **per-peer OPTIONS** — traffic for a TG is not sent to peers that did not select it.
-- **Parrot / echo** talkgroups **9990–9999** bypass the OPTIONS filter and return to the **calling** hotspot (see [Special numbers](special-numbers.md)).
+- **Parrot / echo** talkgroups **9990–9999** are **point-to-point**: they bypass the OPTIONS filter and return **only** to the exact peer that originated the call (`RX_PEER` on the slot), never to other hotspots of the same user. With a single connected peer, it is delivered to that peer (legacy behaviour). See [Echo — Multi-hotspot behaviour](echo.md#multi-hotspot-behaviour-inject-only-proxy) and [Special numbers](special-numbers.md).
---
diff --git a/docs/es/monitor/self-service.md b/docs/es/monitor/self-service.md
index 76873c2..0185452 100644
--- a/docs/es/monitor/self-service.md
+++ b/docs/es/monitor/self-service.md
@@ -75,6 +75,25 @@ Importante: el proxy envía **RPTO solo al master**, no al hotspot directamente.
---
+## Reconciliación de `logged_in`
+
+El flag **`logged_in`** controla tanto el login por **contraseña** como por
+**IP**: sólo las filas con `logged_in = 1` pueden autenticarse en el dashboard.
+El **peer server** mantiene este flag preciso reconciliándolo contra los peers
+realmente conectados cada **120 s** (el bucle `lst_seen`):
+
+- Los peers actualmente conectados al MASTER inject-only quedan `logged_in = 1`.
+- Las demás filas quedan `logged_in = 0`.
+- El bucle arranca con `now=True`, de modo que el **primer tick al arrancar
+ limpia todos los flags obsoletos inmediatamente** — tras un reinicio del
+ servidor, los hotspots que no se reconectaron no pueden autenticarse vía
+ login-by-IP.
+
+Esto reemplazó el `clean_tbl` horario del legado (barrido de 24 h de inactividad),
+que dejaba `logged_in = 1` en peers ya desconectados tras un reinicio.
+
+---
+
## Hash de contraseñas
**`AuthenticateUser`** usa:
diff --git a/docs/es/server/development/behaviour-and-timers.md b/docs/es/server/development/behaviour-and-timers.md
index c7584bc..86611c9 100644
--- a/docs/es/server/development/behaviour-and-timers.md
+++ b/docs/es/server/development/behaviour-and-timers.md
@@ -21,6 +21,7 @@ Los siguientes intervalos forman parte del comportamiento actual en ejecución:
| `bridge_reset` | **6s** | Limpieza de flags de reset y cierre de resets pendientes. |
| OPTIONS refresh | **por evento** | TG estáticas / reflector vía **RPTO**, **startup/reload** (`apply_startup_bridges`), fallback **dmrd** sin source. Sin loop periódico de 26s (**D-28**). |
| `dynamic_tg_purge_loop` | **60s** | Purga filas **SINGLE=1** expiradas de `peer_dynamic_tgs` y `_PEER_UA_SESSIONS` en memoria. |
+| `lst_seen` (reconcile self-service) | **120s** | Reconcilia `Clients.logged_in` contra los peers conectados actualmente: los conectados quedan `logged_in=1`, el resto `0`. Corre con `now=True` en el primer tick para limpiar flags obsoletos inmediatamente tras un reinicio del servidor. Evita que el monitor autentique hotspots desconectados vía login-by-IP. |
| `statTrimmer` | **303s** | Limpieza de bridges STAT obsoletos y estados transitorios. |
Si cambias uno de estos intervalos, documenta el impacto operativo en monitorización, comportamiento de bucles y troubleshooting.
diff --git a/docs/es/server/development/routing-and-contention.md b/docs/es/server/development/routing-and-contention.md
index 3ab0af0..bb1ec02 100644
--- a/docs/es/server/development/routing-and-contention.md
+++ b/docs/es/server/development/routing-and-contention.md
@@ -285,7 +285,7 @@ flowchart TD
PKT[DMRD hacia peer] --> F1{OPTIONS tiene TG
o dinámico activo?}
F1 -->|No| DROP[No entrega]
F1 -->|Sí| F2{TG especial 9990-9999?}
- F2 -->|Sí| PASS1[Bypass slot contention]
+ F2 -->|Sí| PASS1[Punto-a-punto:
sólo RX_PEER (o peer único)]
F2 -->|No| F3{¿Slot destino
libre?}
F3 -->|Peer TXing ingress| DROP
F3 -->|Bridge hold activo
TG distinto| DROP
@@ -309,6 +309,24 @@ flowchart TD
| **SINGLE lock incompatible** | SINGLE=1 con lock en otro TG → bloquea salvo que sea el TG del lock o el peer lo activara. |
| **Sesión stale** | Si la sesión per-peer lleva `_STALE_PEER_SESSION_TIMEOUT` sin frames, se purga y el slot se libera. |
+### TG especiales (9990–9999): entrega punto-a-punto
+
+El eco (9990) y los TG de servicio bajo demanda (9991–9999) **omiten** la
+contención de slot per-peer y los filtros OPTIONS descritos arriba, pero **no**
+se difunden a todos los peers. En un MASTER inject-only multi-peer se entregan
+**punto-a-punto**:
+
+- El paquete va **sólo** a `RX_PEER` — el peer exacto que originó la llamada
+ en ese slot — cuando `RX_TGID` coincide con el TG especial.
+- Si `RX_TGID` aún no coincide (p. ej. el primer VHEAD de una transmisión a
+ 9990 llega antes de que se actualice `RX_TGID`) y sólo hay un peer conectado,
+ se entrega a ese peer único (fallback legado).
+- **No hay matching difuso** del ID DMR de origen: otros hotspots del mismo
+ usuario nunca reciben el eco ni la reproducción de servicio.
+
+Esto coincide con el legado `adn-dmr-server`, donde cada MASTER tenía un solo
+peer y el eco volvía naturalmente sólo al llamante. Ver [Echo](../user-guide/echo.md).
+
### Mid-call join
Cuando un stream downlink termina (VTERM o timeout), el slot del peer queda
diff --git a/docs/es/server/user-guide/echo.md b/docs/es/server/user-guide/echo.md
index b3fed9a..8e2bfa2 100644
--- a/docs/es/server/user-guide/echo.md
+++ b/docs/es/server/user-guide/echo.md
@@ -37,6 +37,27 @@ sudo systemctl enable --now adn-echo
El servidor principal expone un master **ECHO** en **TG 9990** para el bridge de eco. El servicio **echo** independiente es un proceso aparte con su propia config que se conecta a ese master.
+## Comportamiento con varios hotspots (proxy inject-only)
+
+Cuando los hotspots se conectan a través del **proxy integrado** (`PROXY`),
+varios radios pueden compartir el mismo MASTER como peers. En el legado
+`adn-dmr-server` cada MASTER tenía un solo peer, así que el eco volvía
+naturalmente sólo al llamante. El proxy inject-only multi-peer lo impone
+explícitamente:
+
+- **Entrega punto-a-punto.** La reproducción del eco (TG 9990) y los TG de
+ servicio bajo demanda (**9991–9999**) se entregan **sólo** al peer exacto
+ que originó la llamada (`RX_PEER` en el slot activo), **nunca** a otros
+ hotspots del mismo usuario. No hay matching difuso del ID DMR de origen.
+- **Fallback a peer único.** Cuando sólo hay un peer conectado, el paquete se
+ le entrega (comportamiento legado de peer único).
+- Esto aplica tanto al **plano de datos** (enrutado de audio) como al **plano
+ de reportes** (`BRDG_EVENT` enviado al monitor): el monitor muestra el chip
+ de eco en el hotspot originador, no en un hermano.
+
+Ver [Enrutado de voz y contención — Gate de downlink](../development/routing-and-contention.md#gate-de-downlink-un-peer-recibe-el-paquete)
+y [Proxy hotspot — Comportamiento con varios hotspots](hotspot-proxy.md#comportamiento-con-varios-hotspots).
+
## Documentación
Esta página es el resumen incluido en el repositorio; amplía las notas de despliegue localmente según necesites.
diff --git a/docs/es/server/user-guide/hotspot-proxy.md b/docs/es/server/user-guide/hotspot-proxy.md
index f9a5271..86e63c1 100644
--- a/docs/es/server/user-guide/hotspot-proxy.md
+++ b/docs/es/server/user-guide/hotspot-proxy.md
@@ -92,7 +92,7 @@ Detalle del flujo en el panel: [Self-service](../../monitor/self-service.md).
## Comportamiento con varios hotspots
- Cada hotspot autenticado es un **peer** en el MASTER de inyección con sus **OPTIONS** (TG estáticas). **Repeat** y el fan-out del monitor respetan **OPTIONS por peer** — el tráfico de un TG no se envía a peers que no lo tienen seleccionado.
-- Los talkgroups **eco 9990–9999** omiten el filtro OPTIONS y vuelven al hotspot **llamante** (ver [Números especiales](special-numbers.md)).
+- Los talkgroups **eco 9990–9999** son **punto-a-punto**: omiten el filtro OPTIONS y vuelven **sólo** al peer exacto que originó la llamada (`RX_PEER` en el slot), nunca a otros hotspots del mismo usuario. Con un único peer conectado, se le entrega a ése (comportamiento legado). Ver [Echo — Comportamiento con varios hotspots](echo.md#comportamiento-con-varios-hotspots-proxy-inject-only) y [Números especiales](special-numbers.md).
---
diff --git a/src/adn_server/application/routing/hbp_forward.py b/src/adn_server/application/routing/hbp_forward.py
index b6257cb..4c14230 100644
--- a/src/adn_server/application/routing/hbp_forward.py
+++ b/src/adn_server/application/routing/hbp_forward.py
@@ -216,6 +216,13 @@ class HbpForwardMixin:
_slot_st["lastSeq"] = False
_slot_st["lastData"] = False
_slot_st["RX_START"] = pkt_time
+ if _slot_st.get("lastData") and _slot_st["lastData"] == data and seq > 1:
+ _slot_st["loss"] = _slot_st.get("loss", 0) + 1
+ logger.debug(
+ "(%s) *PacketControl* last packet is a complete duplicate, discarding. Stream ID: %s TGID: %s",
+ system_name, int_id(stream_id), int_id(dst_id),
+ )
+ return False
_slot_st["packets"] = _slot_st.get("packets", 0) + 1
_pkts = _slot_st["packets"]
_rx_start = _slot_st.get("RX_START", pkt_time)
@@ -279,13 +286,6 @@ class HbpForwardMixin:
src_proto._obp_send_bcsq(dst_id, stream_id)
_slot_st["_bcsq"] = True
return False
- if _slot_st.get("lastData") and _slot_st["lastData"] == data and seq > 1:
- _slot_st["loss"] = _slot_st.get("loss", 0) + 1
- logger.debug(
- "(%s) *PacketControl* last packet is a complete duplicate, discarding. Stream ID: %s TGID: %s",
- system_name, int_id(stream_id), int_id(dst_id),
- )
- return False
if seq and seq == _slot_st.get("lastSeq"):
_slot_st["loss"] = _slot_st.get("loss", 0) + 1
return False
diff --git a/src/adn_server/application/routing/lc_ta.py b/src/adn_server/application/routing/lc_ta.py
index efa18e1..df79dbe 100644
--- a/src/adn_server/application/routing/lc_ta.py
+++ b/src/adn_server/application/routing/lc_ta.py
@@ -321,6 +321,7 @@ class LcTaMixin:
tgid: int | None = None,
force: bool = False,
fallback_inject: bool = False,
+ repeat_vhead: bool = False,
) -> None:
"""Emit DMRA to an HBP target on VHEAD (once per target stream)."""
if not self._send_dmra_to_system:
@@ -334,9 +335,9 @@ class LcTaMixin:
if force:
if not self._talker_alias.should_resend_passthrough_dmra(target_system, stream_id):
return
- elif not self._talker_alias.should_send_on_vhead(target_system, stream_id):
+ elif not repeat_vhead and not self._talker_alias.should_send_on_vhead(target_system, stream_id):
return
- elif not self._talker_alias.should_send_on_vhead(target_system, stream_id):
+ elif not repeat_vhead and not self._talker_alias.should_send_on_vhead(target_system, stream_id):
return
if not have_passthrough:
# Legacy resolve_ta (both): inject at VHEAD when the buffer is still empty.
@@ -400,6 +401,7 @@ class LcTaMixin:
self._send_talker_alias_to_target(
system_name, system_name, rf_src, stream_id, source_peer,
wire_slot=wire_slot, tgid=tgid,
+ repeat_vhead=True,
)
def prepare_talker_alias_local_repeat(
@@ -526,6 +528,8 @@ class LcTaMixin:
st.pop("TX_TA_EMB", None)
st.pop("TX_TA_PHASE", None)
st.pop("TX_TA_ON", None)
+ st.pop("_ta_last_embed_burst", None)
+ st.pop("_ta_last_embed_frag", None)
emblcs = self._talker_alias.embedded_emblc_for_stream(
source_system,
rf_src,
@@ -562,6 +566,15 @@ class LcTaMixin:
"""
if dtype_vseq not in (1, 2, 3, 4):
return
+ dmrpkt = dmrbits.tobytes()
+ burst_key = (dtype_vseq, dmrpkt)
+ # Duplicated uplink bursts (same B–E payload) must not advance the TA phase
+ # machine; REPEAT runs before ingress duplicate drops (legacy lastData parity).
+ if burst_key == st.get("_ta_last_embed_burst"):
+ last_frag = st.get("_ta_last_embed_frag")
+ if last_frag is not None:
+ dmrbits[EMB_LC_SLICE] = last_frag
+ return
ta_emb = st.get("TX_TA_EMB")
if ta_emb is not None and st.get("TX_TA_ON"):
phase = st.get("TX_TA_PHASE", 0)
@@ -575,6 +588,8 @@ class LcTaMixin:
if dtype_vseq == 4 and ta_emb is not None:
st["TX_TA_ON"] = True
dmrbits[EMB_LC_SLICE] = frag
+ st["_ta_last_embed_burst"] = burst_key
+ st["_ta_last_embed_frag"] = frag
def _clear_talker_alias_embed(self, st: dict[str, Any]) -> None:
st.pop("TX_TA_EMB", None)
@@ -582,3 +597,5 @@ class LcTaMixin:
st.pop("TX_TA_BLOCK_COUNT", None)
st.pop("TX_TA_ON", None)
st.pop("_ta_embed_kind", None)
+ st.pop("_ta_last_embed_burst", None)
+ st.pop("_ta_last_embed_frag", None)
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/hbp/test_ingress.py b/tests/hbp/test_ingress.py
index 7c098a2..558c946 100644
--- a/tests/hbp/test_ingress.py
+++ b/tests/hbp/test_ingress.py
@@ -53,6 +53,29 @@ def test_hbp_ingress_sets_rx_start_on_new_stream() -> None:
assert slot_st.get("RX_STREAM_ID") == base.data()[16:20]
+def test_hbp_rate_limit_ignores_byte_identical_duplicates() -> None:
+ """Compressed duplicate bursts must not inflate the ingress rate counter."""
+ bridges = active_routing_table(91, (("MASTER-A", 2), ("MASTER-B", 2)))
+ scenario = DeterministicScenario(routing_table=bridges)
+ base = PacketSpec(dst_id=91, stream_id=0x90909090)
+ t0 = scenario.clock.time()
+
+ scenario.inject_hbp(
+ "MASTER-A",
+ DeterministicScenario.voice_head_spec(base),
+ ingress_pkt_time=t0,
+ )
+ for seq in range(1, 25):
+ burst = DeterministicScenario.voice_burst_spec(base, seq=seq, dtype_vseq=min(seq, 4))
+ pkt_time = t0 + seq * 0.06
+ ok = scenario.inject_hbp("MASTER-A", burst, ingress_pkt_time=pkt_time)
+ assert ok is not False
+ scenario.inject_hbp("MASTER-A", burst, ingress_pkt_time=pkt_time + 0.003)
+
+ forwarded = len(scenario.capture.for_system("MASTER-B"))
+ assert forwarded >= 20
+
+
def test_hbp_rate_drop_prevents_bridge_forward() -> None:
"""After ingress RATE DROP, no further packets are bridged."""
bridges = active_routing_table(91, (("MASTER-A", 2), ("MASTER-B", 2)))
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()
diff --git a/tests/talker_alias/test_embed_phase_duplicate.py b/tests/talker_alias/test_embed_phase_duplicate.py
new file mode 100644
index 0000000..9b05772
--- /dev/null
+++ b/tests/talker_alias/test_embed_phase_duplicate.py
@@ -0,0 +1,204 @@
+# ADN DMR Peer Server - tests talker alias embed phase duplicate
+#
+# 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
+###############################################################################
+
+"""TA embed phase machine must ignore byte-identical duplicate voice bursts."""
+
+from __future__ import annotations
+
+from bitarray import bitarray
+from tests.harness.deterministic import DeterministicScenario, PacketSpec
+from tests.harness.scenarios import talker_alias_config
+from tests.support.hbp_repeat_stack import build_hbp_repeat_stack
+
+from adn_server.application.routing_use_cases import RoutingUseCases
+from adn_server.domain import bytes_3, bytes_4
+from adn_server.domain.dmr.bptc import encode_emblc
+from adn_server.domain.dmr.const import LC_OPT
+from adn_server.infrastructure.acl_router import InMemoryAclRouter
+from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
+from adn_server.infrastructure.talker_alias_emblc import default_ta_emblc_encoder
+
+_EMB_SLICE = slice(116, 148)
+_PEER_TX = bytes_4(730039210)
+_PEER_RX = bytes_4(730039101)
+_ADDR_TX = ("10.0.0.1", 62001)
+_ADDR_RX = ("10.0.0.2", 62002)
+
+
+def _embed_bits(dmrpkt: bytes) -> bitarray:
+ bits = bitarray(endian="big")
+ bits.frombytes(dmrpkt)
+ return bits[_EMB_SLICE]
+
+
+def _init_repeat_slot(
+ bridge: RoutingUseCases,
+ *,
+ system_name: str = "MASTER-A",
+ slot: int = 2,
+ stream_id: bytes,
+ rf_src: bytes,
+ dst_id: bytes,
+) -> dict:
+ class _Proto:
+ STATUS = {slot: {}}
+
+ proto = _Proto()
+ protocols = {system_name: proto}
+ bridge._get_protocols = lambda: protocols # type: ignore[method-assign]
+ st = proto.STATUS[slot]
+ st["REP_STREAM_ID"] = stream_id
+ st["REP_EMB_LC"] = encode_emblc(LC_OPT + dst_id + rf_src)
+ bridge._init_talker_alias_embed(st, system_name, system_name, rf_src, stream_id)
+ return st
+
+
+def _run_superframe(
+ bridge: RoutingUseCases,
+ *,
+ system_name: str,
+ slot: int,
+ stream_id: bytes,
+ payload: bytes,
+ duplicate_e: bool = False,
+) -> bytes:
+ for dtype in (1, 2, 3, 4):
+ bridge.rewrite_repeat_voice_burst(
+ system_name, slot, stream_id, dtype, payload,
+ )
+ if duplicate_e:
+ bridge.rewrite_repeat_voice_burst(
+ system_name, slot, stream_id, 4, payload,
+ )
+ return bridge.rewrite_repeat_voice_burst(
+ system_name, slot, stream_id, 1, payload,
+ )
+
+
+def test_duplicate_burst_e_does_not_advance_ta_phase() -> None:
+ """Byte-identical burst E must not double-advance the embed phase machine."""
+ config = talker_alias_config()
+ bridge = RoutingUseCases(
+ InMemoryAclRouter(),
+ config,
+ InMemorySubscriptionStore(),
+ get_protocols=lambda: {},
+ encode_emblc=encode_emblc,
+ ta_emblc_encoder=default_ta_emblc_encoder,
+ )
+ stream_id = bytes_4(0xC0FFEE01)
+ rf_src = bytes_3(3120001)
+ dst_id = bytes_3(7304)
+ payload = b"\x42" + b"\x00" * 32
+ st = _init_repeat_slot(
+ bridge,
+ stream_id=stream_id,
+ rf_src=rf_src,
+ dst_id=dst_id,
+ )
+
+ baseline_next_b = _run_superframe(
+ bridge,
+ system_name="MASTER-A",
+ slot=2,
+ stream_id=stream_id,
+ payload=payload,
+ duplicate_e=False,
+ )
+ baseline_phase = st.get("TX_TA_PHASE", 0)
+ baseline_embed = _embed_bits(baseline_next_b)
+
+ st_dup = _init_repeat_slot(
+ bridge,
+ stream_id=stream_id,
+ rf_src=rf_src,
+ dst_id=dst_id,
+ )
+ dup_next_b = _run_superframe(
+ bridge,
+ system_name="MASTER-A",
+ slot=2,
+ stream_id=stream_id,
+ payload=payload,
+ duplicate_e=True,
+ )
+
+ assert st_dup.get("TX_TA_PHASE", 0) == baseline_phase
+ assert _embed_bits(dup_next_b) == baseline_embed
+
+
+def test_repeat_stack_duplicate_burst_matches_non_duplicate_embed() -> None:
+ """Integration: duplicated REPEAT burst E keeps the next superframe TA embed aligned."""
+
+ def _play_through(duplicate_e: bool) -> bitarray:
+ stack = build_hbp_repeat_stack(talker_alias=True)
+ stack.register_peer(_PEER_TX, _ADDR_TX, options="TS2=7304;")
+ stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;")
+ base = PacketSpec(
+ peer_id=730039210,
+ rf_src=7300392,
+ dst_id=7304,
+ slot=2,
+ stream_id=0xA1B2C3D4,
+ payload=b"\x77" + b"\x00" * 32,
+ )
+ stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
+ for seq, dtype in enumerate((1, 2, 3, 4), start=1):
+ stack.inject_spec(
+ DeterministicScenario.voice_burst_spec(base, seq=seq, dtype_vseq=dtype),
+ _ADDR_TX,
+ )
+ if duplicate_e:
+ stack.inject_spec(
+ DeterministicScenario.voice_burst_spec(base, seq=5, dtype_vseq=4),
+ _ADDR_TX,
+ )
+ stack.transport.clear()
+ next_seq = 6 if duplicate_e else 5
+ stack.inject_spec(
+ DeterministicScenario.voice_burst_spec(base, seq=next_seq, dtype_vseq=1),
+ _ADDR_TX,
+ )
+ downlink = stack.transport.for_addr(_ADDR_RX)
+ assert downlink, "expected downlink after superframe"
+ return _embed_bits(downlink[0][20:53])
+
+ assert _play_through(duplicate_e=False) == _play_through(duplicate_e=True)
+
+
+def test_two_vheads_emit_dmra_on_repeat_path() -> None:
+ """Legacy hblink re-sends DMRA on every VHEAD; REPEAT must not dedupe the second."""
+ stack = build_hbp_repeat_stack(talker_alias=True)
+ stack.register_peer(_PEER_TX, _ADDR_TX, options="TS2=7304;")
+ stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;")
+ base = PacketSpec(
+ peer_id=730039210,
+ rf_src=7300392,
+ dst_id=7304,
+ slot=2,
+ stream_id=0xA1B2C3D4,
+ )
+
+ stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
+ first_dmra = len(stack.dmra_capture)
+ assert first_dmra == 1
+
+ stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
+ assert len(stack.dmra_capture) == first_dmra + 1