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