Merge pull request #51 from ce5rpy/develop

feat: release 2.4
pull/53/head
ce5rpy 3 months ago committed by GitHub
commit 17f480854c
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -80,6 +80,16 @@ PROXY:
BLACK_LIST: []
IP_BLACK_LIST: {}
# OpenBridge fan-in (standard with OPENBRIDGE). Hotspots: PROXY 62031; OBP mesh: OBP_PROXY 62032.
# OBP_PROXY block is optional: when absent, defaults apply (ENABLED, LISTEN_PORT 62032, BIND_LEGACY_PORTS).
# BIND_LEGACY_PORTS: true keeps each SYSTEMS.*.PORT open during migration; false = LISTEN_PORT only.
OBP_PROXY:
ENABLED: true
LISTEN_PORT: 62032
LISTEN_IP: ""
BIND_LEGACY_PORTS: true
DEBUG: false
# MariaDB (required): dynamic TG persistence + optional self-service (Clients table).
DATABASE:
DB_SERVER: localhost

@ -25,7 +25,7 @@ The **ADN DMR Peer Server** is a [GPL-3.0](https://www.gnu.org/licenses/gpl-3.0.
| Private calls | [Private calls](server/user-guide/private-calls.md) |
| Voice / TTS | [Voice, announcements, and TTS](server/user-guide/voice-and-tts.md) |
| Legacy dashboard + server 2.x | [Report proxy](server/user-guide/report-proxy.md) |
| OpenBridge / DMRE | [OpenBridge](server/protocols/openbridge.md), [DMRE v5](server/protocols/dmre-v5.md) |
| OpenBridge / DMRE | [OpenBridge](server/protocols/openbridge.md), [DMRE v5](server/protocols/dmre-v5.md), [OBP proxy](server/user-guide/obp-proxy.md) |
| HBP | [HBP](server/protocols/hbp.md) |
| Code layout | [Architecture](server/development/architecture.md), [Behaviour and timers](server/development/behaviour-and-timers.md) |
| Credits, license, lineage | [Credits & license](server/user-guide/attribution.md) |
@ -50,6 +50,7 @@ Dashboard, WebSocket live view, FastAPI API, **MySQL** self-service — see [Mon
| I want to… | Start here |
|------------|------------|
| `adn-server.yaml` — integrated `PROXY` / `SELF_SERVICE` | [Hotspot proxy (integrated)](server/user-guide/hotspot-proxy.md) |
| `adn-server.yaml` — `OBP_PROXY` OpenBridge fan-in | [OBP proxy](server/user-guide/obp-proxy.md) |
| `adn-monitor.yaml`, layout | [Monitor configuration](monitor/configuration.md) |
| Integrated hotspot proxy | [Hotspot proxy](server/user-guide/hotspot-proxy.md) |
| Standalone hotspot proxy (removed) | [Hotspot proxy — moved](monitor/hotspot-proxy.md) |

@ -1,99 +1,49 @@
# Performance (2.x)
**adn-server 2.x** and **adn-monitor 2.x** include several changes that reduce CPU work and memory footprint compared with **adn-dmr-server** and the old monitor/proxy stack. This page lists **what** improves and **what causes it**.
**adn-server 2.x** and **adn-monitor 2.x** use less CPU and RAM than **adn-dmr-server**
and the old monitor/proxy stack. You do not need to tune anything — the gains come from
how routing, reporting, and the integrated proxy are built.
## At a glance
Pair **server 2.x with monitor 2.x** to get the full benefit on the dashboard side.
| Area | Typical effect | Main cause |
|------|----------------|------------|
| **Voice downlink (inject proxy)** | Lower CPU under busy group traffic | **`PeerDownlinkIndex`** — fan-out to peers that match `(slot, TG)` instead of scanning every connected hotspot per packet |
| **Bridge source lookup** | Faster “am I the ACTIVE source?” | **`SubscriptionStore`** indexes (`relay_tables_with_active_source`) — O(1) by `(system, slot, tgid)` vs scanning table rows |
| **Background CPU** | Fewer wakeups | **Event-driven OPTIONS / static TG** — removed legacy **26 s** `options_config_loop` ([Behaviour and timers](behaviour-and-timers.md)) |
| **Mass peer login** | Less redundant CONFIG traffic | **`ConfigPushThrottle`** — adaptive debounce on CONFIG push to the monitor |
| **Reporting vs voice** | Voice path less blocked by reports | **`BoundedReportQueue`** — coalesced snapshots, bounded drain per tick |
| **Server → monitor wire** | Less serialize/send work | **Report v2** JSON (`routing_table`, `topology`, `voice_event`) instead of periodic full pickle of `CONFIG`/`BRIDGES` ([Report protocol v2](../protocols/report-v2.md)) |
| **Process count (RAM)** | One Python process instead of two | **Integrated `PROXY`** in `adn-server.py` — no separate **adn-proxy** process ([Hotspot proxy](../user-guide/hotspot-proxy.md)) |
| **Monitor RAM / WS load** | Smaller in-memory dashboard state | **Slim `dashboard_state` wire**, `clean_sys_dict`, lighter WebSocket fingerprints ([Monitor architecture](../../monitor/architecture.md)) |
## What improved
## Server: inject-only downlink index
| Improvement | What you get |
|-------------|--------------|
| **Smarter voice routing** | Under busy group traffic the server does less work per packet — especially on proxies with many hotspots and on OpenBridge-heavy nodes. |
| **Integrated hotspot proxy** | One `adn-server` process instead of server + standalone **adn-proxy** — less RAM and simpler ops. See [Hotspot proxy](../user-guide/hotspot-proxy.md). |
| **No legacy 26 s timer** | Static talkgroups refresh on events (startup, config reload, peer OPTIONS), not a background loop every 26 seconds. See [Behaviour and timers](behaviour-and-timers.md). |
| **Lighter monitor link** | **Report v2** sends compact JSON instead of heavy periodic pickle dumps. See [Report protocol v2](../protocols/report-v2.md). |
| **Monitor 2.x** | Slimmer dashboard state, less memory growth on panels left open for days. See [Monitor architecture](../../monitor/architecture.md). |
The largest **CPU** win on many ADN networks is on the **MASTER inject-only** path (`PROXY` with inject-only mode).
## OpenBridge-heavy servers (2.3.3+)
**Legacy:** `send_peers` walks **every registered peer** for each downlink packet → cost grows as **O(peers × packets/s)**.
If you run **many OpenBridges** and upgrade from an earlier **2.x** build to **2.3.3 or
newer**, a production node under comparable OBP load showed roughly:
**2.x:** `PeerDownlinkIndex` precomputes candidates from each peer’s **OPTIONS** (static TGs) and **UA session** state. For each group voice frame, only peers that **might** want that `(slot, TG)` are considered; each candidate still passes `peer_should_receive_group_voice`.
| | Before (2.x) | After (2.3.3+) |
|--|--------------|----------------|
| **CPU** | baseline | **~25% of before** |
| **RAM** | baseline | **~75% of before** |
```text
Legacy: every DMRD → try all N peers
2.x: every DMRD → index lookup → try k peers (k ≪ N on busy proxies)
```
Exact numbers depend on traffic and hardware; treat these as a reference, not a guarantee.
OPTIONS parsing is **cached per peer** (`_CACHED_OPTIONS_STATIC`): if the OPTIONS blob is unchanged, already-parsed static TGs are reused instead of re-parsing on every packet.
## When you will notice it most
| Code | Role |
|------|------|
| `application/routing/peer_downlink_index.py` | Index build and `(slot, tgid) → candidates` |
| `infrastructure/twisted_adapters/udp_hbp.py` | `_iter_downlink_peers`, `send_peers` |
| `tests/infrastructure/test_peer_downlink_fanout.py` | Inject-only fan-out tests |
| Your network | Effect |
|--------------|--------|
| Small install, few peers, light traffic | Modest — the stack is simply lighter overall. |
| **Inject proxy with many hotspots** | **Clear CPU win** when group voice is busy. |
| **Many OpenBridges, steady mesh traffic** | **Clear CPU and RAM win** after **2.3.3+** (see table above). |
| **Many hotspots logging in at once** | Less load on server and monitor during login bursts. |
| **Monitor open 24/7** | Lower and more stable RAM on **adn-monitor 2.x**. |
**When it matters:** proxy with **tens to hundreds** of hotspots and steady group voice. On a small conference with few peers, the difference is minor.
## Server: routing indexes
On every group voice frame the server must find relay tables where **this system is the ACTIVE source**.
**Legacy:** scan rows inside `BRIDGES[table_key]` (and related tables).
**2.x:** `InMemorySubscriptionStore.relay_tables_with_active_source()` uses a maintained **`_source_tables`** index — lookup by `(system, slot, dst_tgid)` without walking all legs.
This lives in the subscription store implementation; it is an **algorithmic index**, not a separate feature you configure.
| Code | Role |
|------|------|
| `infrastructure/subscription_store.py` | `_source_tables`, `_by_table`, `_active_target_counts` |
| `application/subscription/router.py` | `SubscriptionRouter.resolve()` |
## Server: less periodic and login-storm work
| Change | What it avoids |
|--------|----------------|
| **No 26 s OPTIONS loop** | Timer firing every 26 s across all systems to refresh static bridges when RPTO/startup/reload already handle it |
| **`ConfigPushThrottle`** | Flooding the monitor with CONFIG snapshots when many peers connect within a few seconds (debounce widens from ~0.3 s to ~2 s during bursts) |
| **`BoundedReportQueue`** | Doing pickle/JSON encode and TCP send synchronously on the voice hot path; coalesces duplicate config/bridge snapshots |
## Server: reporting and deployment
- **Report v2** — structured JSON replaces opaque pickle snapshots for bridge/config state on the **2.x monitor** wire. See [Monitoring and reports](../user-guide/monitoring.md) and [Report protocol v2](../protocols/report-v2.md).
- **Integrated proxy** — `PROXY` runs **in-process**; dropping the standalone **adn-proxy** saves baseline **RAM** (one interpreter, shared config) and simplifies ops.
## Monitor (adn-monitor 2.x)
Pair **adn-server 2.x** with **adn-monitor 2.x** to get the reporting-side gains:
| Change | Effect |
|--------|--------|
| **Slim wire / `dashboard_state`** | Monitor ingests compact JSON state instead of holding full duplicated pickle trees from v1 |
| **`clean_sys_dict`** | Periodic eviction of stale in-memory entries (caps runaway growth on long-lived panels) |
| **Last-heard row cache, lighter WS fingerprints** | Less work per dashboard refresh |
| **Unified FastAPI stack** | Removed separate PHP API and standalone monitor **proxy** process |
Details: [Monitor architecture](../../monitor/architecture.md).
## When you will notice a difference
| Deployment | CPU | RAM |
|------------|-----|-----|
| Few masters, no inject proxy, light traffic | Small | Small |
| **Inject-only proxy, many hotspots, busy TG** | **Clear** (downlink index) | Moderate (single server process vs server+proxy) |
| Long-lived monitor + report v2 | Moderate (less serialize on wire) | **Clearer** on monitor (slim state, `clean_sys_dict`) |
Crypto, AMBE, and OpenBridge MAC work still dominate on OpenBridge-heavy paths — routing-table optimizations do not remove that cost.
Crypto, AMBE, and OpenBridge wire work still cost CPU on busy OBP paths — routing
optimizations remove redundant bridge work, not voice codec or encryption overhead.
## Related reading
- [Architecture](architecture.md) — layers and entrypoint
- [BRIDGES vs Subscriptions](bridges-vs-subscriptions.md) — routing model (not a performance feature)
- [Behaviour and timers](behaviour-and-timers.md) — event-driven OPTIONS vs legacy 26 s loop
- [Hotspot proxy](../user-guide/hotspot-proxy.md) — integrated `PROXY` / inject-only
- [Report protocol v2](../protocols/report-v2.md) — JSON wire to monitor
- Release notes: `CHANGELOG.md` at the repository root (`Performance` under **2.0.0-rc.1**).
- [Hotspot proxy](../user-guide/hotspot-proxy.md) — integrated `PROXY`
- [Monitoring and reports](../user-guide/monitoring.md) — report v2 pairing
- [Behaviour and timers](behaviour-and-timers.md) — event-driven OPTIONS
- Release notes: `CHANGELOG.md` at the repository root

@ -155,6 +155,20 @@ The **echo** example (`adn-echo.example.yaml`) is a PEER that attaches to the **
Ingress filters and loop control: [OpenBridge protocol](../protocols/openbridge.md) (including [DMRE vs OpenBridge v5](../protocols/openbridge.md#dmre-and-openbridge-v5)) and [Special numbers — OpenBridge ingress](special-numbers.md#openbridge-ingress--group-tg-filters).
Inbound UDP for all OPENBRIDGE systems can use the shared **`OBP_PROXY`** fan-in (standard port **62032**). See [OBP proxy](obp-proxy.md).
---
## `OBP_PROXY` (integrated OpenBridge fan-in)
Active when **`OBP_PROXY`** is present or when any **OPENBRIDGE** system exists (defaults apply). All OpenBridge inbound UDP is demuxed on **`LISTEN_PORT`** (default **62032**); per-bridge legacy listeners are optional. Full guide: [OBP proxy](obp-proxy.md).
| Key | Meaning |
|-----|---------|
| **LISTEN_PORT** / **LISTEN_IP** | UDP bind for OpenBridge fan-in (pair to hotspot **`PROXY`** **62031**). |
| **BIND_LEGACY_PORTS** | When `true`, also listen each **`SYSTEMS.*.PORT`** that differs from **`LISTEN_PORT`** (per-bridge migration). |
| **ENABLED** | Set `false` to restore per-bridge UDP bind (legacy mode). |
---
## `BRIDGES` (runtime)
@ -295,3 +309,4 @@ Use the project interpreter (see workspace rules), e.g. `python3.11` from pyenv,
- [Special numbers](special-numbers.md) — reserved TGs and server IDs.
- [Echo](echo.md) — PEER example (echo process).
- [Hotspot proxy](hotspot-proxy.md) — integrated **`PROXY`** / **`SELF_SERVICE`**.
- [OBP proxy](obp-proxy.md) — integrated **`OBP_PROXY`** OpenBridge fan-in.

@ -23,11 +23,13 @@ Routing, timers, OpenBridge loop control, and protocol handling are implemented
| **Voice** | AMBE files, scheduled announcements, TTS pipeline, on-demand playback (TG 9991–9999). |
| **Reporting** | TCP netstring channel to **adn-monitor** (and compatible dashboards): config, bridge state, call events (report v2 JSON). |
| **Hotspot proxy** | Optional integrated UDP fan-in (`PROXY` in `adn-server.yaml`) plus MySQL **self-service** (`SELF_SERVICE`) for dashboard-driven hotspot options. |
| **OBP proxy** | Optional integrated OpenBridge UDP fan-in (`OBP_PROXY` in `adn-server.yaml`; default when OPENBRIDGE systems exist). |
## Related programs
- **Echo / playback** — `adn-server.py --echo` with minimal `adn-echo.yaml`; see [Echo](echo.md).
- **Integrated hotspot proxy** — `PROXY` in **`adn-server.yaml`**; see [Hotspot proxy](hotspot-proxy.md).
- **Integrated OBP proxy** — `OBP_PROXY` fan-in for OpenBridge; see [OBP proxy](obp-proxy.md).
- **Report proxy (legacy dashboards)** — optional **[ADN-report-proxy](https://github.com/ce5rpy/ADN-report-proxy)** so **adn-server 2.x** can feed old HBMonitor / FDMR-style monitors (v1 wire); see [Report proxy](report-proxy.md). Not used with **adn-monitor 2.x**.
## Next steps
@ -37,6 +39,7 @@ Routing, timers, OpenBridge loop control, and protocol handling are implemented
- [Voice routing and contention](../development/routing-and-contention.md) — the full packet flow, contention rules, SINGLE, slot mapping, and divergences.
- [Special numbers](special-numbers.md) — TG 4000, information services, echo.
- [Hotspot proxy](hotspot-proxy.md) — integrated **`PROXY`** / **`SELF_SERVICE`** in `adn-server.yaml`.
- [OBP proxy](obp-proxy.md) — integrated **`OBP_PROXY`** OpenBridge fan-in.
- [ADN Monitor](../../monitor/index.md) — dashboard, `adn-monitor.yaml`, self-service UI (separate repo, deployed with the server).
- [Performance (2.x)](../development/performance.md) — CPU/RAM improvements in this release and what causes them.
- [Credits & license](attribution.md) — ADN → FreeDMR → hblink3, license.

@ -0,0 +1,51 @@
# OBP proxy (single inbound port)
Optional `OBP_PROXY` stanza configures the fan-in listener for all `MODE: OPENBRIDGE` systems. When active, OpenBridge instances are **inject-only** (no per-bridge `listenUDP` in `HBPProtocol`); the proxy owns every inbound OBP socket.
## Activation
| YAML | Behaviour |
|------|-----------|
| No `OBP_PROXY` block, no OPENBRIDGE | N/A (proxy not started). |
| No `OBP_PROXY` block, OPENBRIDGE present | **Default proxy on** (`LISTEN_PORT` 62032, `BIND_LEGACY_PORTS` true). |
| `OBP_PROXY.ENABLED: false` | Legacy mode: each OPENBRIDGE binds its own `PORT`. |
| `OBP_PROXY.ENABLED: true` | Proxy manages all OBP inbound UDP (same as absent block). |
## Configuration
```yaml
OBP_PROXY:
ENABLED: true
LISTEN_PORT: 62032 # ADN standard OBP fan-in (pair to PROXY 62031)
LISTEN_IP: "" # optional bind address
BIND_LEGACY_PORTS: true # default: also listen each SYSTEMS.*.PORT
DEBUG: false
```
OPENBRIDGE sections stay unchanged (`PORT`, `NETWORK_ID`, `PASSPHRASE`, `TARGET_*`, ACL, etc.). With proxy enabled, `PORT` is kept as metadata (`_REPORT_PORT` internally) for monitor/report and optional legacy listeners.
## Per-bridge migration (`BIND_LEGACY_PORTS: true`)
When the global flag is true, each OPENBRIDGE can migrate individually:
| `SYSTEMS.*.PORT` | Behaviour |
|------------------|-----------|
| Same as `OBP_PROXY.LISTEN_PORT` (e.g. 62032) | Fan-in only for this bridge (no extra legacy listener). |
| Omitted, `0`, or empty | Same as `LISTEN_PORT` — fan-in only (migrated bridge). |
| Any other port (e.g. 62999) | Legacy listener stays open for that bridge. |
Example: migrate `OBP-CL2` to the shared fan-in while `OBP-EU` keeps `PORT: 62999`.
## Migration
1. Existing configs with OPENBRIDGE but no `OBP_PROXY` stanza already use defaults (`BIND_LEGACY_PORTS: true`) — no remote changes required.
2. Optionally add an explicit `OBP_PROXY` block to tune `LISTEN_PORT` / `BIND_LEGACY_PORTS`.
3. Set `BIND_LEGACY_PORTS: false` and close legacy ports when all remotes use `LISTEN_PORT`.
## Requirements
- `NETWORK_ID` must be unique among enabled OPENBRIDGE systems.
- `LISTEN_PORT` must not collide with any OPENBRIDGE `PORT` when `BIND_LEGACY_PORTS` is true.
- `RELAX_CHECKS: true` is recommended so `TARGET_SOCK` is learned from the first valid packet.
See also: [OpenBridge protocol](../protocols/openbridge.md).

@ -24,7 +24,7 @@ El **ADN DMR Peer Server** es un puente de conferencia [GPL-3.0](https://www.gnu
| Llamadas privadas | [Llamadas privadas](server/user-guide/private-calls.md) |
| Voz / TTS | [Voz, anuncios y TTS](server/user-guide/voice-and-tts.md) |
| Panel legacy + servidor 2.x | [Proxy de informes](server/user-guide/report-proxy.md) |
| OpenBridge / DMRE | [OpenBridge](server/protocols/openbridge.md), [DMRE v5](server/protocols/dmre-v5.md) |
| OpenBridge / DMRE | [OpenBridge](server/protocols/openbridge.md), [DMRE v5](server/protocols/dmre-v5.md), [Proxy OBP](server/user-guide/obp-proxy.md) |
| HBP | [HBP](server/protocols/hbp.md) |
| Código | [Arquitectura](server/development/architecture.md), [Comportamiento y temporizadores](server/development/behaviour-and-timers.md) |
| Créditos, licencia, linaje | [Créditos y licencia](server/user-guide/attribution.md) |
@ -49,6 +49,7 @@ Panel, **WebSocket** en vivo, API FastAPI, **MySQL** self-service — ver [Descr
| Quiero… | Empieza aquí |
|---------|----------------|
| `adn-server.yaml` — `PROXY` / `SELF_SERVICE` integrados | [Proxy hotspot (integrado)](server/user-guide/hotspot-proxy.md) |
| `adn-server.yaml` — fan-in OpenBridge `OBP_PROXY` | [Proxy OBP](server/user-guide/obp-proxy.md) |
| `adn-monitor.yaml`, despliegue | [Configuración del monitor](monitor/configuration.md) |
| Proxy hotspot integrado | [Proxy hotspot](server/user-guide/hotspot-proxy.md) |
| Proxy hotspot (UDP, `PROXY`, rango de puertos) | [Proxy hotspot](monitor/hotspot-proxy.md) |

@ -1,99 +1,50 @@
# Rendimiento (2.x)
**adn-server 2.x** y **adn-monitor 2.x** incluyen varios cambios que reducen trabajo de CPU y huella de memoria frente a **adn-dmr-server** y al stack antiguo de monitor/proxy. Esta página resume **qué** mejora y **qué lo provoca**.
**adn-server 2.x** y **adn-monitor 2.x** consumen menos CPU y RAM que **adn-dmr-server**
y el stack antiguo de monitor/proxy. No hace falta ajustar nada — las mejoras vienen de
cómo están hechos el enrutado, los informes y el proxy integrado.
## Resumen
Empareja **servidor 2.x con monitor 2.x** para aprovechar también el lado del panel.
| Área | Efecto típico | Causa principal |
|------|---------------|-----------------|
| **Downlink de voz (proxy inject)** | Menos CPU con tráfico de grupo intenso | **`PeerDownlinkIndex`** — fan-out solo a peers que encajan `(slot, TG)` en lugar de escanear todos los hotspots por paquete |
| **Origen ACTIVE en bridge** | Lookup más rápido | **Índices del `SubscriptionStore`** (`relay_tables_with_active_source`) — O(1) por `(system, slot, tgid)` frente a recorrer filas |
| **CPU de fondo** | Menos despertares | **OPTIONS / TG estática por eventos** — eliminado el bucle legacy cada **26 s** `options_config_loop` ([Comportamiento y temporizadores](behaviour-and-timers.md)) |
| **Ráfaga de logins** | Menos CONFIG redundante | **`ConfigPushThrottle`** — debounce adaptativo al empujar CONFIG al monitor |
| **Informes vs voz** | La voz se bloquea menos por informes | **`BoundedReportQueue`** — snapshots coalescidos, drenado acotado por tick |
| **Cable servidor → monitor** | Menos serializar/enviar | **Informe v2** JSON (`routing_table`, `topology`, `voice_event`) en lugar de pickle periódico de `CONFIG`/`BRIDGES` ([Protocolo de informes v2](../protocols/report-v2.md)) |
| **Procesos (RAM)** | Un proceso Python en lugar de dos | **`PROXY` integrado** en `adn-server.py` — sin proceso **adn-proxy** aparte ([Proxy hotspot](../user-guide/hotspot-proxy.md)) |
| **RAM / WS del monitor** | Estado de panel más compacto | **Wire slim `dashboard_state`**, `clean_sys_dict`, fingerprints WS más ligeros ([Arquitectura del monitor](../../monitor/architecture.md)) |
## Qué mejoró
## Servidor: índice de downlink inject-only
La mayor ganancia de **CPU** en muchas redes ADN está en el camino **MASTER inject-only** (`PROXY` en modo inject-only).
**Legacy:** `send_peers` recorre **todos los peers registrados** por cada paquete de downlink → coste **O(peers × paquetes/s)**.
**2.x:** `PeerDownlinkIndex` precalcula candidatos desde **OPTIONS** (TG estáticas) y estado **UA** de cada peer. Por cada frame de voz de grupo solo se consideran peers que **podrían** querer ese `(slot, TG)`; cada candidato sigue pasando `peer_should_receive_group_voice`.
```text
Legacy: cada DMRD → probar los N peers
2.x: cada DMRD → lookup en índice → probar k peers (k ≪ N en proxies cargados)
```
El parse de OPTIONS se **guarda en caché por peer** (`_CACHED_OPTIONS_STATIC`): si el blob OPTIONS no cambió, se reutilizan las TG estáticas ya parseadas en lugar de volver a interpretarlo en cada paquete.
| Código | Rol |
|--------|-----|
| `application/routing/peer_downlink_index.py` | Construcción del índice y `(slot, tgid) → candidatos` |
| `infrastructure/twisted_adapters/udp_hbp.py` | `_iter_downlink_peers`, `send_peers` |
| `tests/infrastructure/test_peer_downlink_fanout.py` | Tests de fan-out inject-only |
**Cuándo se nota:** proxy con **decenas o cientos** de hotspots y voz de grupo continua. En una conferencia pequeña con pocos peers, la diferencia es pequeña.
## Servidor: índices de enrutado
En cada frame de voz de grupo el servidor debe encontrar tablas donde **este system es origen ACTIVE**.
**Legacy:** recorrer filas dentro de `BRIDGES[clave]`.
**2.x:** `InMemorySubscriptionStore.relay_tables_with_active_source()` usa el índice **`_source_tables`** — lookup por `(system, slot, dst_tgid)` sin recorrer todas las patas.
Está en la implementación del store; es un **índice algorítmico**, no una opción de configuración aparte.
| Código | Rol |
|--------|-----|
| `infrastructure/subscription_store.py` | `_source_tables`, `_by_table`, `_active_target_counts` |
| `application/subscription/router.py` | `SubscriptionRouter.resolve()` |
## Servidor: menos trabajo periódico y en tormenta de logins
| Cambio | Qué evita |
| Mejora | Qué ganas |
|--------|-----------|
| **Sin bucle OPTIONS 26 s** | Timer cada 26 s en todos los systems cuando RPTO/arranque/reload ya refrescan bridges estáticos |
| **`ConfigPushThrottle`** | Inundar al monitor con snapshots CONFIG cuando muchos peers conectan en pocos segundos (debounce ~0,3 s → ~2 s en ráfaga) |
| **`BoundedReportQueue`** | Encode pickle/JSON y envío TCP en el hot path de voz; coalesce de snapshots config/bridge duplicados |
| **Enrutado de voz más eficiente** | Con tráfico de grupo intenso el servidor hace menos trabajo por paquete — sobre todo en proxies con muchos hotspots y en nodos con muchos OpenBridge. |
| **Proxy hotspot integrado** | Un solo proceso `adn-server` en lugar de servidor + **adn-proxy** aparte — menos RAM y operación más simple. Ver [Proxy hotspot](../user-guide/hotspot-proxy.md). |
| **Sin timer legacy de 26 s** | Las TG estáticas se refrescan por eventos (arranque, recarga de config, OPTIONS del peer), no con un bucle de fondo cada 26 segundos. Ver [Comportamiento y temporizadores](behaviour-and-timers.md). |
| **Cable al monitor más liviano** | **Informe v2** envía JSON compacto en lugar de volcados pickle pesados. Ver [Protocolo de informes v2](../protocols/report-v2.md). |
| **Monitor 2.x** | Estado de panel más compacto, menos crecimiento de memoria con el panel abierto días. Ver [Arquitectura del monitor](../../monitor/architecture.md). |
## Servidor: informes y despliegue
## Servidores con muchos OpenBridge (2.3.3+)
- **Informe v2** — JSON estructurado sustituye snapshots pickle opacos de bridge/config en el cable hacia **monitor 2.x**. Ver [Monitor e informes](../user-guide/monitoring.md) y [Protocolo de informes v2](../protocols/report-v2.md).
- **Proxy integrado** — `PROXY` **in-process**; quitar **adn-proxy** standalone ahorra **RAM** base (un intérprete, config compartida) y simplifica operación.
Si tienes **muchos OpenBridge** y actualizas desde un **2.x** anterior a **2.3.3 o
superior**, un nodo en producción con carga OBP comparable mostró aproximadamente:
## Monitor (adn-monitor 2.x)
| | Antes (2.x) | Después (2.3.3+) |
|--|-------------|------------------|
| **CPU** | línea base | **~25% de la línea base** |
| **RAM** | línea base | **~75% de la línea base** |
Empareja **adn-server 2.x** con **adn-monitor 2.x** para las mejoras del lado informes:
Las cifras exactas dependen del tráfico y del hardware; tómalo como referencia, no como garantía.
| Cambio | Efecto |
|--------|--------|
| **Wire slim / `dashboard_state`** | El monitor ingiere JSON compacto en lugar de duplicar árboles pickle v1 |
| **`clean_sys_dict`** | Expulsión periódica de entradas obsoletas en memoria (tope de crecimiento en paneles largos) |
| **Caché lastheard, fingerprints WS ligeros** | Menos trabajo por refresco del dashboard |
| **Stack FastAPI unificado** | Eliminados API PHP y proceso **proxy** standalone del monitor |
Detalle: [Arquitectura del monitor](../../monitor/architecture.md).
## Cuándo se nota más
## Cuándo se nota la diferencia
| Despliegue | CPU | RAM |
|------------|-----|-----|
| Pocos masters, sin proxy inject, tráfico bajo | Poca | Poca |
| **Proxy inject-only, muchos hotspots, TG activa** | **Clara** (índice downlink) | Moderada (un proceso servidor vs servidor+proxy) |
| Monitor largo + informe v2 | Moderada (menos serializar en cable) | **Más clara** en monitor (estado slim, `clean_sys_dict`) |
| Tu red | Efecto |
|--------|--------|
| Instalación pequeña, pocos peers, poco tráfico | Moderado — el stack es más liviano en general. |
| **Proxy inject con muchos hotspots** | **Ganancia clara de CPU** con voz de grupo activa. |
| **Muchos OpenBridge, tráfico mesh continuo** | **Ganancia clara de CPU y RAM** tras **2.3.3+** (ver tabla anterior). |
| **Muchos hotspots conectando a la vez** | Menos carga en servidor y monitor en ráfagas de login. |
| **Monitor abierto 24/7** | RAM más baja y estable en **adn-monitor 2.x**. |
Crypto, AMBE y MAC OpenBridge siguen dominando en tramos OBP cargados — optimizar la tabla de bridge no elimina ese coste.
Crypto, AMBE y el trabajo de cable OpenBridge siguen costando CPU en tramos OBP
cargados — las optimizaciones de routing quitan trabajo redundante de bridges, no el
codec de voz ni el cifrado.
## Lecturas relacionadas
- [Arquitectura](architecture.md) — capas y entrypoint
- [BRIDGES vs Subscriptions](bridges-vs-subscriptions.md) — modelo de enrutado (no es feature de rendimiento)
- [Comportamiento y temporizadores](behaviour-and-timers.md) — OPTIONS por eventos vs bucle 26 s legacy
- [Proxy hotspot](../user-guide/hotspot-proxy.md) — `PROXY` integrado / inject-only
- [Protocolo de informes v2](../protocols/report-v2.md) — cable JSON al monitor
- Notas de versión: `CHANGELOG.md` en la raíz del repositorio (`Performance` en **2.0.0-rc.1**).
- [Proxy hotspot](../user-guide/hotspot-proxy.md) — `PROXY` integrado
- [Monitor e informes](../user-guide/monitoring.md) — emparejar informe v2
- [Comportamiento y temporizadores](behaviour-and-timers.md) — OPTIONS por eventos
- Notas de versión: `CHANGELOG.md` en la raíz del repositorio

@ -155,6 +155,20 @@ El ejemplo **echo** (`adn-echo.example.yaml`) es un PEER que se une al MASTER **
Filtros de ingreso y control de bucle: [Protocolo OpenBridge](../protocols/openbridge.md) (incl. [DMRE frente a OpenBridge v5](../protocols/openbridge.md#dmre-and-openbridge-v5)) y [Números especiales — ingreso OpenBridge](special-numbers.md#openbridge-ingress--group-tg-filters).
El UDP entrante de todos los OPENBRIDGE puede usar el fan-in compartido **`OBP_PROXY`** (puerto estándar **62032**). Ver [Proxy OBP](obp-proxy.md).
---
## `OBP_PROXY` (fan-in OpenBridge integrado)
Activo cuando existe **`OBP_PROXY`** o algún sistema **OPENBRIDGE** (se aplican defaults). Todo el UDP OBP entrante se demultiplexa en **`LISTEN_PORT`** (default **62032**); los listeners legacy por bridge son opcionales. Guía completa: [Proxy OBP](obp-proxy.md).
| Clave | Significado |
|-------|-------------|
| **LISTEN_PORT** / **LISTEN_IP** | Bind UDP del fan-in OpenBridge (pareja del **`PROXY`** hotspot **62031**). |
| **BIND_LEGACY_PORTS** | Con `true`, también escucha cada **`SYSTEMS.*.PORT`** distinto de **`LISTEN_PORT`** (migración bridge a bridge). |
| **ENABLED** | `false` restaura bind UDP por bridge (modo legacy). |
---
## `BRIDGES` (tiempo de ejecución)
@ -295,3 +309,4 @@ Usa el intérprete del proyecto (ver reglas del workspace), p. ej. `python3.11`
- [Números especiales](special-numbers.md) — TG e IDs reservados.
- [Echo](echo.md) — ejemplo PEER (proceso echo).
- [Proxy hotspot](hotspot-proxy.md) — **`PROXY`** / **`SELF_SERVICE`** integrados.
- [Proxy OBP](obp-proxy.md) — fan-in OpenBridge **`OBP_PROXY`** integrado.

@ -23,11 +23,13 @@ Enrutado, temporizadores, control de bucle OpenBridge y manejo de protocolo est
| **Voz** | Ficheros AMBE, anuncios programados, tubería TTS, reproducción bajo demanda (TG 9991–9999). |
| **Informes** | Canal TCP netstring hacia **adn-monitor** (y paneles compatibles): config, estado de bridges, eventos de llamada (informe v2 JSON). |
| **Proxy hotspot** | Fan-in UDP integrado opcional (`PROXY` en `adn-server.yaml`) y **self-service** MySQL (`SELF_SERVICE`) para opciones de hotspot desde el panel. |
| **Proxy OBP** | Fan-in UDP OpenBridge integrado opcional (`OBP_PROXY` en `adn-server.yaml`; por defecto si hay sistemas OPENBRIDGE). |
## Programas relacionados
- **Echo / playback** — `adn-server.py --echo` con `adn-echo.yaml` mínimo; ver [Echo](echo.md).
- **Proxy hotspot integrado** — `PROXY` en **`adn-server.yaml`**; ver [Proxy hotspot](hotspot-proxy.md).
- **Proxy OBP integrado** — fan-in `OBP_PROXY` para OpenBridge; ver [Proxy OBP](obp-proxy.md).
- **Proxy de informes (paneles legacy)** — **[ADN-report-proxy](https://github.com/ce5rpy/ADN-report-proxy)** opcional para que **adn-server 2.x** alimente monitores antiguos estilo HBMonitor / FDMR (wire v1); ver [Proxy de informes](report-proxy.md). No se usa con **adn-monitor 2.x**.
## Siguientes pasos
@ -37,6 +39,7 @@ Enrutado, temporizadores, control de bucle OpenBridge y manejo de protocolo est
- [Enrutado de voz y contención](../development/routing-and-contention.md) — el flujo completo de paquetes, reglas de contención, SINGLE, mapeo de slot y divergencias.
- [Números especiales](special-numbers.md) — TG 4000, servicios de información, eco.
- [Proxy hotspot](hotspot-proxy.md) — **`PROXY`** / **`SELF_SERVICE`** integrados en `adn-server.yaml`.
- [Proxy OBP](obp-proxy.md) — fan-in OpenBridge **`OBP_PROXY`** integrado.
- [ADN Monitor](../../monitor/index.md) — panel, `adn-monitor.yaml`, UI self-service (repo aparte, desplegado con el servidor).
- [Rendimiento (2.x)](../development/performance.md) — mejoras de CPU/RAM en esta versión y qué las provoca.
- [Créditos y licencia](attribution.md) — ADN → FreeDMR → hblink3, licencia.

@ -0,0 +1,49 @@
# Proxy OBP (puerto de entrada único)
La stanza opcional `OBP_PROXY` configura el listener fan-in para todos los sistemas `MODE: OPENBRIDGE`. Con el proxy activo, las instancias OpenBridge son **inject-only** (sin `listenUDP` por bridge en `HBPProtocol`); el proxy gestiona todo el UDP OBP entrante.
## Activación
| YAML | Comportamiento |
|------|----------------|
| Sin `OBP_PROXY`, sin OPENBRIDGE | N/A (no se arranca proxy). |
| Sin `OBP_PROXY`, con OPENBRIDGE | **Proxy por defecto** (`LISTEN_PORT` 62032, `BIND_LEGACY_PORTS` true). |
| `OBP_PROXY.ENABLED: false` | Modo legacy: cada OPENBRIDGE hace bind en su `PORT`. |
| `OBP_PROXY.ENABLED: true` | El proxy gestiona toda la entrada OBP (igual que bloque ausente). |
## Configuración
```yaml
OBP_PROXY:
ENABLED: true
LISTEN_PORT: 62032 # puerto estándar OBP fan-in (pareja de PROXY 62031)
LISTEN_IP: "" # dirección de bind opcional
BIND_LEGACY_PORTS: true # default: también escucha cada SYSTEMS.*.PORT
DEBUG: false
```
Las secciones OPENBRIDGE no cambian (`PORT`, `NETWORK_ID`, `PASSPHRASE`, `TARGET_*`, ACL, etc.). Con proxy activo, `PORT` se conserva como metadato (`_REPORT_PORT` internamente) para monitor/report y listeners legacy opcionales.
## Migración por bridge (`BIND_LEGACY_PORTS: true`)
Con el flag global activo, cada OPENBRIDGE migra de forma individual:
| `SYSTEMS.*.PORT` | Comportamiento |
|------------------|----------------|
| Igual a `OBP_PROXY.LISTEN_PORT` (p. ej. 62032) | Solo fan-in para ese bridge (sin listener legacy extra). |
| Omitido, `0` o vacío | Igual que `LISTEN_PORT` — solo fan-in (bridge migrado). |
| Otro puerto (p. ej. 62999) | Se mantiene el listener legacy de ese bridge. |
Ejemplo: migrar `OBP-CL2` al fan-in compartido mientras `OBP-EU` conserva `PORT: 62999`.
## Migración
1. Configs existentes con OPENBRIDGE sin stanza `OBP_PROXY` ya usan defaults (`BIND_LEGACY_PORTS: true`) — sin cambios en remotos.
2. Opcionalmente añadir bloque `OBP_PROXY` explícito para ajustar `LISTEN_PORT` / `BIND_LEGACY_PORTS`.
3. `BIND_LEGACY_PORTS: false` y cerrar puertos legacy cuando todos usen `LISTEN_PORT`.
- `NETWORK_ID` único entre OPENBRIDGE habilitados.
- `LISTEN_PORT` sin colisión con ningún `PORT` de sección si `BIND_LEGACY_PORTS` es true.
- `RELAX_CHECKS: true` recomendado para aprender `TARGET_SOCK` del primer paquete válido.
Ver también: [protocolo OpenBridge](../protocols/openbridge.md).

@ -63,6 +63,7 @@ nav:
- Monitor e informes: server/user-guide/monitoring.md
- Proxy de informes (paneles legacy): server/user-guide/report-proxy.md
- Proxy hotspot (integrado): server/user-guide/hotspot-proxy.md
- Proxy OBP (fan-in OpenBridge): server/user-guide/obp-proxy.md
- Echo (reproducción): server/user-guide/echo.md
- Créditos y licencia: server/user-guide/attribution.md
- Protocolos:

@ -63,6 +63,7 @@ nav:
- Monitoring and reports: server/user-guide/monitoring.md
- Report proxy (legacy dashboards): server/user-guide/report-proxy.md
- Hotspot proxy (integrated): server/user-guide/hotspot-proxy.md
- OBP proxy (OpenBridge fan-in): server/user-guide/obp-proxy.md
- Echo (playback): server/user-guide/echo.md
- Credits & license: server/user-guide/attribution.md
- Protocols:

@ -31,7 +31,7 @@
"MASTER-B": {
"mode": "MASTER",
"ip": "10.0.0.2",
"port": 62032,
"port": 62040,
"peers": {
"3120002": {
"id": 3120002,
@ -57,6 +57,7 @@
"ip": "44.31.61.68",
"port": 62999,
"enhanced_obp": true,
"connected": false,
"streams": {}
}
}

@ -155,7 +155,11 @@
"network_id": { "$ref": "#/$defs/dmr_id" },
"ip": { "type": "string" },
"port": { "type": "integer", "minimum": 0, "maximum": 65535 },
"enhanced_obp": { "type": "boolean" }
"enhanced_obp": { "type": "boolean" },
"connected": {
"type": "boolean",
"description": "ENHANCED_OBP BCKA keepalive OK (peer answered within 60s). Omitted when ENHANCED_OBP is false."
}
}
},
"dashboard_ctable": {

@ -20,15 +20,29 @@
"""Proxy application layer (Phase 3)."""
from .deployment import is_proxy_inject_only, normalize_proxy_target, proxy_target_system
from .deployment import (
is_obp_proxy_managed,
is_proxy_inject_only,
normalize_obp_proxy_targets,
normalize_proxy_target,
obp_bridge_legacy_listen_port,
obp_proxy_bind_legacy_ports,
obp_proxy_enabled,
proxy_target_system,
)
from .packet_helpers import peer_id_from_packet
from .use_cases import ProxySlotError, ProxyUseCases
__all__ = [
"ProxySlotError",
"ProxyUseCases",
"is_obp_proxy_managed",
"is_proxy_inject_only",
"normalize_obp_proxy_targets",
"normalize_proxy_target",
"obp_bridge_legacy_listen_port",
"obp_proxy_bind_legacy_ports",
"obp_proxy_enabled",
"peer_id_from_packet",
"proxy_target_system",
]

@ -61,3 +61,102 @@ def normalize_proxy_target(config: dict[str, Any]) -> None:
sys_cfg["_REPORT_BASE_PORT"] = int(port)
else:
sys_cfg.setdefault("_REPORT_BASE_PORT", 56400)
def _obp_proxy_block(config: dict[str, Any]) -> dict[str, Any] | None:
block = config.get("OBP_PROXY")
return block if isinstance(block, dict) else None
def config_has_enabled_openbridge(config: dict[str, Any]) -> bool:
"""True when config defines at least one enabled OPENBRIDGE system."""
systems = config.get("SYSTEMS", {})
if not isinstance(systems, dict):
return False
return any(
isinstance(cfg, dict) and cfg.get("ENABLED", True) and cfg.get("MODE") == "OPENBRIDGE"
for cfg in systems.values()
)
def obp_proxy_enabled(config: dict[str, Any]) -> bool:
"""True when OBP proxy manages inbound UDP (default on for OPENBRIDGE configs)."""
block = _obp_proxy_block(config)
if block is not None:
return bool(block.get("ENABLED", True))
return config_has_enabled_openbridge(config)
def obp_proxy_bind_legacy_ports(config: dict[str, Any]) -> bool:
"""When OBP proxy is enabled, also listen on each OPENBRIDGE section PORT (default true)."""
if not obp_proxy_enabled(config):
return False
block = _obp_proxy_block(config)
if block is None:
return True
return bool(block.get("BIND_LEGACY_PORTS", True))
def obp_bridge_legacy_listen_port(
sys_cfg: dict[str, Any],
*,
listen_port: int,
bind_legacy_ports: bool,
) -> int | None:
"""Per-bridge legacy UDP port, or None when inbound uses OBP_PROXY fan-in only.
When BIND_LEGACY_PORTS is true globally, a bridge with PORT equal to LISTEN_PORT is
treated as migrated to the shared fan-in (no extra legacy listener).
"""
if not bind_legacy_ports:
return None
report_port = sys_cfg.get("_REPORT_PORT")
if report_port is not None:
port = int(report_port)
elif "PORT" in sys_cfg:
port = int(sys_cfg.get("PORT", 0) or 0)
else:
return None
if port <= 0 or port == listen_port:
return None
return port
def is_obp_proxy_managed(config: dict[str, Any], system_name: str) -> bool:
"""OPENBRIDGE under active OBP proxy — no direct HBPProtocol UDP bind."""
if not obp_proxy_enabled(config):
return False
sys_cfg = config.get("SYSTEMS", {}).get(system_name)
if not isinstance(sys_cfg, dict):
return False
return sys_cfg.get("MODE") == "OPENBRIDGE"
def normalize_obp_proxy_targets(config: dict[str, Any]) -> None:
"""Strip bind fields from OPENBRIDGE systems when OBP proxy manages inbound UDP."""
if not obp_proxy_enabled(config):
return
block = _obp_proxy_block(config) or {}
listen_port = int(block.get("LISTEN_PORT", 62032))
systems = config.get("SYSTEMS", {})
if not isinstance(systems, dict):
return
for sys_cfg in systems.values():
if not isinstance(sys_cfg, dict):
continue
if sys_cfg.get("MODE") != "OPENBRIDGE":
continue
port = sys_cfg.pop("PORT", None)
bind_ip = sys_cfg.pop("IP", None)
parsed_port = 0
if port is not None and port != "":
try:
parsed_port = int(port)
except (TypeError, ValueError):
parsed_port = 0
if parsed_port > 0:
sys_cfg["_REPORT_PORT"] = parsed_port
else:
sys_cfg["_REPORT_PORT"] = listen_port
if bind_ip is not None and str(bind_ip).strip():
sys_cfg["_REPORT_BIND_IP"] = str(bind_ip)

@ -73,8 +73,28 @@ def _upstream_peer_block(name: str, cfg: dict[str, Any]) -> dict[str, Any]:
return block
def _openbridge_block(name: str, cfg: dict[str, Any], topology_row: dict[str, Any] | None) -> dict[str, Any]:
def _obp_ka_connected(cfg: dict[str, Any], now: float) -> bool | None:
"""BCKA keepalive status for ENHANCED OBP legs; ``None`` when KA gating does not apply."""
if not cfg.get("ENHANCED_OBP"):
return None
bcka = cfg.get("_bcka")
if bcka is None:
return False
try:
return float(bcka) >= now - 60
except (TypeError, ValueError):
return False
def _openbridge_block(
name: str,
cfg: dict[str, Any],
topology_row: dict[str, Any] | None,
*,
now: float,
) -> dict[str, Any]:
"""Enabled OPENBRIDGE legs (``CTABLE.OPENBRIDGES``); STREAMS stay empty here (live chips = monitor/voice)."""
del name
block: dict[str, Any] = {"mode": "OPENBRIDGE", "streams": {}}
network_id = cfg.get("NETWORK_ID")
if network_id is not None:
@ -86,6 +106,9 @@ def _openbridge_block(name: str, cfg: dict[str, Any], topology_row: dict[str, An
block["port"] = int(row["port"])
if row.get("enhanced_obp") or cfg.get("ENHANCED_OBP"):
block["enhanced_obp"] = True
connected = _obp_ka_connected(cfg, now)
if connected is not None:
block["connected"] = connected
return block
@ -137,7 +160,7 @@ def build_dashboard_state(
block["port"] = int(topo["port"])
masters[name] = block
elif mode == "OPENBRIDGE":
openbridges[name] = _openbridge_block(name, cfg, topo)
openbridges[name] = _openbridge_block(name, cfg, topo, now=epoch)
elif mode in ("PEER", "XLXPEER") and _upstream_peer_connected(cfg):
peers[name] = _upstream_peer_block(name, cfg)

@ -35,7 +35,6 @@ from ..domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3, bytes_4, int_id
from .ports import VoiceProvider
from .server_voice import (
announcement_item_source_bytes,
server_voice_id,
server_voice_rf_src_bytes,
)
@ -128,6 +127,15 @@ class VoiceUseCases:
from .routing.helpers import master_dynamic_tg_slots, master_slot_holds_server_broadcast
from .server_voice import all_server_voice_ids
if master_slot_holds_server_broadcast(
slot,
time.time(),
server_voice_rf_srcs=all_server_voice_ids(self._config),
):
return False
# External QSO on RX while slot is in hangtime — defer announcement inject.
if slot.get("RX_TYPE") != HBPF_SLT_VTERM and slot.get("TX_TYPE") == HBPF_SLT_VTERM:
return True
if wire_ts in master_dynamic_tg_slots(sys_cfg, int(tg)):
return False
ptt_system = self._announcement_ptt_system or ""
@ -135,12 +143,6 @@ class VoiceUseCases:
return False
if slot.get("RX_TYPE") == HBPF_SLT_VTERM and slot.get("TX_TYPE") == HBPF_SLT_VTERM:
return False
if master_slot_holds_server_broadcast(
slot,
time.time(),
server_voice_rf_srcs=all_server_voice_ids(self._config),
):
return False
return True
def _inject_ptt_slot_order(self, tg: int, sys_cfg: dict[str, Any]) -> list[int]:

@ -39,8 +39,11 @@ from adn_server.application import (
)
from adn_server.application.dynamic_tg_use_cases import DynamicTgUseCases
from adn_server.application.proxy.deployment import (
is_obp_proxy_managed,
is_proxy_inject_only,
normalize_obp_proxy_targets,
normalize_proxy_target,
obp_proxy_enabled,
proxy_target_system,
)
from adn_server.application.report.queue import BoundedReportQueue, QueuedReportSender
@ -77,6 +80,10 @@ from adn_server.infrastructure.persistence.dynamic_tg_repository import MysqlDyn
from adn_server.infrastructure.persistence.keys_store import JsonKeysStore
from adn_server.infrastructure.persistence.mysql_pool import create_mysql_pool, ensure_database_sync
from adn_server.infrastructure.proxy import apply_proxy_config_reload, start_proxy_service
from adn_server.infrastructure.proxy.obp_runtime import (
apply_obp_proxy_config_reload,
start_obp_proxy_service,
)
from adn_server.infrastructure.security.password_download import (
DefaultSecurityDownloader,
StubSecurityDownloader,
@ -219,6 +226,7 @@ def run_peer_server(
_ensure_system_runtime_config(config)
_normalize_peer_config(config)
_normalize_obp_config(config)
normalize_obp_proxy_targets(config)
runtime_holder = RuntimeContextHolder(RuntimeContext(config=config, config_path=config_path))
config = ConfigProxy(runtime_holder)
@ -592,10 +600,16 @@ def run_peer_server(
# UDP / proxy listeners (proxy_state declared before reload handler uses nonlocal)
proxy_state = None
obp_proxy_state = None
proxy_enabled = not no_proxy and proxy_target_system(config) is not None
obp_proxy_on = obp_proxy_enabled(config)
def _should_bind_udp(system_name: str, sys_cfg: dict[str, Any]) -> bool:
return not is_proxy_inject_only(config, system_name)
if is_proxy_inject_only(config, system_name):
return False
if is_obp_proxy_managed(config, system_name):
return False
return True
def _on_config_systems_changed() -> None:
routing_use_cases.apply_startup_subscriptions()
@ -604,9 +618,10 @@ def run_peer_server(
user_passwords_loader.load(config)
def _apply_reload_success(result: Any, *, new_config: dict[str, Any], mqtt_before: Any) -> None:
nonlocal report_mqtt, proxy_state
nonlocal report_mqtt, proxy_state, obp_proxy_state
swap_runtime_config(runtime_holder, new_config, config_path=config_path)
normalize_proxy_target(config)
normalize_obp_proxy_targets(config)
report_factory.set_config(config)
mqtt_after = mqtt_settings_from_config(config)
report_mqtt = reconcile_mqtt_publisher(
@ -631,11 +646,20 @@ def run_peer_server(
_wire_proxy_report_slots(report_factory, proxy_state)
else:
_wire_proxy_report_slots(report_factory, None)
if obp_proxy_state is not None:
apply_obp_proxy_config_reload(
obp_proxy_state,
config,
protocols,
logger=logger,
)
elif obp_proxy_enabled(config):
obp_proxy_state = start_obp_proxy_service(config, protocols, logger=logger)
if result.added or result.removed or result.updated or result.rebound:
_on_config_systems_changed()
def _do_config_reload() -> None:
nonlocal report_mqtt, proxy_state
nonlocal report_mqtt, proxy_state, obp_proxy_state
mqtt_before = mqtt_settings_from_config(config)
new_config = prepare_reload_config(runtime_holder)
@ -732,7 +756,10 @@ def run_peer_server(
protocol = _create_hbp_protocol(system_name)
protocols[system_name] = protocol
if not _should_bind_udp(system_name, sys_cfg):
logger.info("(PROXY) %s inject-only (no UDP bind)", system_name)
if is_obp_proxy_managed(config, system_name):
logger.info("(OBP_PROXY) %s inject-only (no UDP bind)", system_name)
else:
logger.info("(PROXY) %s inject-only (no UDP bind)", system_name)
continue
bind = BindSpec(ip=str(sys_cfg.get("IP") or "0.0.0.0"), port=int(sys_cfg.get("PORT", 56400)))
udp_ports[system_name] = _listen_system(system_name, bind, protocol)
@ -784,6 +811,21 @@ def run_peer_server(
reactor.addSystemEventTrigger("before", "shutdown", _stop_proxy)
_wire_proxy_report_slots(report_factory, proxy_state)
if obp_proxy_on:
try:
obp_proxy_state = start_obp_proxy_service(config, protocols, logger=logger)
except Exception as exc:
logger.error("(OBP_PROXY) failed to start OBP proxy: %s", exc)
raise
def _stop_obp_proxy(_: Any = None) -> None:
if obp_proxy_state is not None:
for deferred in obp_proxy_state.stop():
if deferred is not None:
pass
reactor.addSystemEventTrigger("before", "shutdown", _stop_obp_proxy)
if proxy_state is None or proxy_state.self_service is None:
def dynamic_reload_loop() -> None:
dynamic_tg_uc.process_reload_queue(try_purge=_try_purge_dynamic_reload)

@ -29,7 +29,7 @@ from typing import Any, Callable
from twisted.internet import defer
from adn_server.application.proxy.deployment import normalize_proxy_target
from adn_server.application.proxy.deployment import normalize_obp_proxy_targets, normalize_proxy_target
from ..domain.errors import ConfigError
from .config_loader import YamlConfigLoader
@ -88,6 +88,7 @@ def prepare_incoming_config(
ensure_system_runtime_config(incoming)
normalize_peer_config(incoming)
normalize_obp_config(incoming)
normalize_obp_proxy_targets(incoming)
return incoming

@ -23,6 +23,8 @@ from __future__ import annotations
from typing import Any
from adn_server.application.proxy.deployment import config_has_enabled_openbridge
from ..domain.errors import ConfigError
ACL_KEYS = frozenset(
@ -262,6 +264,71 @@ def _validate_proxy(proxy_cfg: dict[str, Any] | None, systems: dict[str, Any], e
)
def _validate_obp_proxy(
obp_cfg: dict[str, Any] | None,
systems: dict[str, Any],
errors: list[str],
) -> None:
if isinstance(obp_cfg, dict) and not obp_cfg.get("ENABLED", True):
return
if not isinstance(obp_cfg, dict) and not config_has_enabled_openbridge({"SYSTEMS": systems}):
return
effective = obp_cfg if isinstance(obp_cfg, dict) else {}
for key in ("DEBUG", "BIND_LEGACY_PORTS"):
if key in effective:
_expect_bool(f"OBP_PROXY.{key}", effective[key], errors)
if "LISTEN_IP" in effective:
_expect_str("OBP_PROXY.LISTEN_IP", effective["LISTEN_IP"], errors)
listen_port = effective.get("LISTEN_PORT", 62032)
if isinstance(listen_port, bool) or not isinstance(listen_port, int) or listen_port < 1:
errors.append("OBP_PROXY.LISTEN_PORT: required >= 1 when ENABLED.")
return
bind_legacy = bool(effective.get("BIND_LEGACY_PORTS", True))
network_ids: dict[Any, str] = {}
legacy_ports: set[int] = set()
if not isinstance(systems, dict):
return
for name, sys_cfg in systems.items():
if not isinstance(sys_cfg, dict) or not sys_cfg.get("ENABLED", True):
continue
if sys_cfg.get("MODE") != "OPENBRIDGE":
continue
raw_nid = sys_cfg.get("NETWORK_ID")
if raw_nid is not None and not _is_empty(raw_nid):
if isinstance(raw_nid, bytes):
nid_key = raw_nid
else:
try:
nid_key = int(raw_nid) & 0xFFFFFFFF
except (TypeError, ValueError):
errors.append(f"SYSTEMS.{name}.NETWORK_ID: invalid value {raw_nid!r}.")
continue
prev = network_ids.get(nid_key)
if prev is not None:
errors.append(
f"SYSTEMS.{name}.NETWORK_ID: duplicate OPENBRIDGE identity "
f"(same as SYSTEMS.{prev})."
)
else:
network_ids[nid_key] = name
if bind_legacy:
port = sys_cfg.get("PORT", 0)
if not _is_empty(port):
try:
legacy_port = int(port)
except (TypeError, ValueError):
errors.append(f"SYSTEMS.{name}.PORT: expected integer.")
continue
if legacy_port > 0:
if legacy_port == listen_port:
continue
if legacy_port in legacy_ports:
errors.append(
f"SYSTEMS.{name}.PORT: duplicate OPENBRIDGE listen port {legacy_port}."
)
legacy_ports.add(legacy_port)
def _config_requires_database(config: dict[str, Any]) -> bool:
"""True for ``run_peer_server`` configs (proxy/master); not echo-only PEER fleets."""
proxy = config.get("PROXY")
@ -352,6 +419,12 @@ def validate_config(config: dict[str, Any], *, config_path: str | None = None) -
proxy_cfg = config.get("PROXY")
_validate_proxy(proxy_cfg if isinstance(proxy_cfg, dict) else None, systems if isinstance(systems, dict) else {}, errors)
obp_proxy_cfg = config.get("OBP_PROXY")
_validate_obp_proxy(
obp_proxy_cfg if isinstance(obp_proxy_cfg, dict) else None,
systems if isinstance(systems, dict) else {},
errors,
)
if _config_requires_database(config):
_validate_database(config.get("DATABASE"), errors)

@ -31,7 +31,10 @@ from typing import TextIO
from adn_server.application.proxy.deployment import (
is_proxy_inject_only,
normalize_obp_proxy_targets,
normalize_proxy_target,
obp_bridge_legacy_listen_port,
obp_proxy_enabled,
proxy_target_system,
)
from adn_server.domain.errors import ConfigError
@ -43,6 +46,7 @@ from adn_server.infrastructure.config_normalizer import (
normalize_obp_config,
normalize_peer_config,
)
from adn_server.infrastructure.proxy.obp_config import obp_proxy_settings
@dataclass(frozen=True)
@ -121,6 +125,7 @@ def collect_findings(
ensure_system_runtime_config(config)
normalize_peer_config(config)
normalize_obp_config(config)
normalize_obp_proxy_targets(config)
findings.append(Finding("ok", "config", f"loaded {config_path}"))
@ -170,19 +175,49 @@ def collect_findings(
)
)
elif mode == "OPENBRIDGE":
ip = str(sys_cfg.get("IP") or "0.0.0.0")
port = int(sys_cfg.get("PORT", 62044))
proto_ver = sys_cfg.get("VER", sys_cfg.get("PROTO_VER", 5))
enhanced = sys_cfg.get("ENHANCED_OBP", True)
ok, detail = _check_udp_bind(ip, port)
level = "ok" if ok else "error"
findings.append(
Finding(
level,
"ports",
f"{name}: OPENBRIDGE UDP {ip}:{port} — {detail}; PROTO_VER={proto_ver}, ENHANCED_OBP={enhanced}",
if obp_proxy_enabled(config):
obp = obp_proxy_settings(config)
legacy_port = obp_bridge_legacy_listen_port(
sys_cfg,
listen_port=int(obp["listen_port"]),
bind_legacy_ports=bool(obp["bind_legacy_ports"]),
)
if legacy_port is not None:
ip = str(sys_cfg.get("_REPORT_BIND_IP") or "0.0.0.0")
ok, detail = _check_udp_bind(ip, legacy_port)
level = "ok" if ok else "error"
findings.append(
Finding(
level,
"ports",
f"{name}: OPENBRIDGE legacy UDP {ip}:{legacy_port} — {detail}; "
f"PROTO_VER={proto_ver}, ENHANCED_OBP={enhanced}",
)
)
else:
findings.append(
Finding(
"ok",
"systems",
f"{name}: OPENBRIDGE fan-in only (OBP_PROXY {obp['listen_port']}); "
f"PROTO_VER={proto_ver}, ENHANCED_OBP={enhanced}",
)
)
else:
ip = str(sys_cfg.get("IP") or "0.0.0.0")
port = int(sys_cfg.get("PORT", 62044))
ok, detail = _check_udp_bind(ip, port)
level = "ok" if ok else "error"
findings.append(
Finding(
level,
"ports",
f"{name}: OPENBRIDGE UDP {ip}:{port} — {detail}; "
f"PROTO_VER={proto_ver}, ENHANCED_OBP={enhanced}",
)
)
)
target_ip = sys_cfg.get("TARGET_IP", "")
if isinstance(target_ip, bytes):
target_ip = target_ip.decode("utf-8", errors="replace")
@ -217,6 +252,20 @@ def collect_findings(
elif proxy_target and no_proxy:
findings.append(Finding("warn", "proxy", f"PROXY configured ({proxy_target}) but --no-proxy set"))
if obp_proxy_enabled(config):
obp = obp_proxy_settings(config)
ip = str(obp["listen_ip"] or "0.0.0.0")
port = int(obp["listen_port"])
ok, detail = _check_udp_bind(ip, port)
level = "ok" if ok else "error"
findings.append(
Finding(
level,
"proxy",
f"OBP_PROXY UDP {ip}:{port} — {detail}; BIND_LEGACY_PORTS={obp['bind_legacy_ports']}",
)
)
aliases = config.get("ALIASES", {})
for key, default in (
("PEER_FILE", "peer_ids.json"),

@ -23,6 +23,9 @@
from .config import apply_proxy_env_overrides, proxy_settings
from .hbp_adapters import FanInClientSender, HbpMasterPeerRegistry, InProcessHbpSink
from .ip_blacklist import InMemoryProxyIpBlacklist
from .obp_config import obp_proxy_settings
from .obp_fanin import ObpFanInDemux, ObpFanInProtocol, listen_obp_fanin
from .obp_runtime import ObpProxyServiceState, apply_obp_proxy_config_reload, start_obp_proxy_service
from .reply_transport import ProxyReplyTransport
from .rpto_queue import InMemoryPendingRptoQueue
from .runtime import ProxyServiceState, apply_proxy_config_reload, start_proxy_service
@ -38,14 +41,21 @@ __all__ = [
"InMemoryProxyIpBlacklist",
"InMemoryProxySlotStore",
"InProcessHbpSink",
"ObpFanInDemux",
"ObpFanInProtocol",
"ObpProxyServiceState",
"ProxyFanInProtocol",
"ProxyReplyTransport",
"ProxyServiceState",
"apply_obp_proxy_config_reload",
"apply_proxy_config_reload",
"apply_proxy_env_overrides",
"apply_session_teardown",
"listen_obp_fanin",
"listen_proxy_fanin",
"obp_proxy_settings",
"proxy_settings",
"self_service_settings",
"start_obp_proxy_service",
"start_proxy_service",
]

@ -0,0 +1,41 @@
# ADN DMR Peer Server - infrastructure proxy obp config
#
# 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
###############################################################################
"""OBP_PROXY runtime settings from config dict (infrastructure; no business rules)."""
from __future__ import annotations
from typing import Any
from adn_server.application.proxy.deployment import obp_proxy_bind_legacy_ports, obp_proxy_enabled
def obp_proxy_settings(config: dict[str, Any]) -> dict[str, Any]:
"""Resolved OBP_PROXY runtime settings with defaults (block may be absent)."""
block = config.get("OBP_PROXY", {})
if not isinstance(block, dict):
block = {}
return {
"enabled": obp_proxy_enabled(config),
"listen_port": int(block.get("LISTEN_PORT", 62032)),
"listen_ip": str(block.get("LISTEN_IP") or ""),
"bind_legacy_ports": obp_proxy_bind_legacy_ports(config),
"debug": bool(block.get("DEBUG")),
}

@ -0,0 +1,225 @@
# ADN DMR Peer Server - infrastructure proxy obp fanin
#
# 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
###############################################################################
"""UDP fan-in for OPENBRIDGE: demux by NETWORK_ID and optional legacy PORT."""
from __future__ import annotations
import logging
from dataclasses import dataclass, field
from typing import Any, Protocol
from twisted.internet.protocol import DatagramProtocol
from adn_server.infrastructure.hbp_constants import BCKA, BCSQ, BCST, BCVE, DMRD, DMRE
from adn_server.infrastructure.mesh.obp_v1 import verify_bcka, verify_bcsq, verify_bcst, verify_bcve
from adn_server.infrastructure.udp_rcvbuf import apply_udp_rcvbuf, udp_rcvbuf_bytes
_logger = logging.getLogger(__name__)
class _DatagramWriter(Protocol):
def write(self, data: bytes, addr: tuple[str, int]) -> None:
...
class _ObpReceiver(Protocol):
def _obp_datagram_received(self, data: bytes, sockaddr: tuple[str, int]) -> None:
...
class ObpIngressReplyTransport:
"""Route OBP egress through the fan-in socket that last received for this bridge."""
def __init__(self, fallback: _DatagramWriter) -> None:
self._fallback = fallback
self._active: _DatagramWriter | None = None
def note_ingress(self, transport: _DatagramWriter) -> None:
self._active = transport
def write(self, data: bytes, addr: tuple[str, int]) -> None:
transport = self._active or self._fallback
transport.write(data, addr)
class InProcessObpSink:
"""Deliver datagrams to an OPENBRIDGE HBPProtocol without a UDP hop."""
def __init__(self, hbp: _ObpReceiver) -> None:
self._hbp = hbp
def inject(self, data: bytes, client_addr: tuple[str, int]) -> None:
self._hbp._obp_datagram_received(data, client_addr)
@dataclass
class ObpBridgeEntry:
system_name: str
network_id: bytes
passphrase: bytes
sink: InProcessObpSink
reply_transport: ObpIngressReplyTransport
legacy_port: int | None = None
@dataclass
class ObpBridgeRegistry:
"""NETWORK_ID and legacy PORT lookup for OBP fan-in demux."""
by_network_id: dict[bytes, str] = field(default_factory=dict)
by_legacy_port: dict[int, str] = field(default_factory=dict)
bridges: dict[str, ObpBridgeEntry] = field(default_factory=dict)
def register(self, entry: ObpBridgeEntry) -> None:
self.bridges[entry.system_name] = entry
self.by_network_id[entry.network_id] = entry.system_name
if entry.legacy_port is not None:
self.by_legacy_port[entry.legacy_port] = entry.system_name
def clear(self) -> None:
self.by_network_id.clear()
self.by_legacy_port.clear()
self.bridges.clear()
class ObpFanInDemux:
"""Shared demux handler for one or more OBP proxy UDP listeners."""
def __init__(
self,
registry: ObpBridgeRegistry,
*,
debug: bool = False,
logger: logging.Logger | None = None,
) -> None:
self._registry = registry
self.debug = debug
self._log = logger or _logger
def deliver(
self,
data: bytes,
addr: tuple[str, int],
*,
local_port: int,
transport: _DatagramWriter,
) -> None:
host, port = addr
system_name = self._registry.by_legacy_port.get(local_port)
if system_name is None:
system_name = self._lookup_by_network_id(data)
if system_name is None:
system_name = self._lookup_control(data)
if system_name is None:
if self.debug:
self._log.debug(
"(OBP_PROXY) dropped packet from %s:%s len=%d local_port=%s",
host,
port,
len(data),
local_port,
)
return
entry = self._registry.bridges.get(system_name)
if entry is None:
return
if self.debug:
self._log.debug(
"(OBP_PROXY) RX %s from %s:%s len=%d -> %s",
data[:4],
host,
port,
len(data),
system_name,
)
entry.reply_transport.note_ingress(transport)
entry.sink.inject(data, addr)
def _lookup_by_network_id(self, data: bytes) -> str | None:
if len(data) < 15:
return None
opcode = data[:4]
if opcode not in (DMRD, DMRE):
return None
network_id = data[11:15]
return self._registry.by_network_id.get(network_id)
def _lookup_control(self, data: bytes) -> str | None:
if len(data) < 4:
return None
opcode = data[:4]
for name, entry in self._registry.bridges.items():
passphrase = entry.passphrase
if opcode == BCKA and verify_bcka(data, passphrase):
return name
if opcode == BCSQ and verify_bcsq(data, passphrase) is not None:
return name
if opcode == BCST and verify_bcst(data, passphrase):
return name
if opcode == BCVE:
ok, _ver = verify_bcve(data, passphrase)
if ok:
return name
return None
class ObpFanInProtocol(DatagramProtocol):
"""Thin UDP listener delegating to shared OBP demux."""
def __init__(self, demux: ObpFanInDemux) -> None:
self._demux = demux
def datagramReceived(self, data: bytes, addr: tuple[str, int]) -> None:
transport = self.transport
if transport is None:
return
local_port = int(transport.getHost().port)
self._demux.deliver(data, addr, local_port=local_port, transport=transport)
def listen_obp_fanin(
reactor: Any,
listen_ip: str,
listen_port: int,
demux: ObpFanInDemux,
*,
config: dict[str, Any] | None = None,
udp_rcvbuf: int | None = None,
logger: logging.Logger | None = None,
) -> tuple[ObpFanInProtocol, Any]:
"""Bind one OBP proxy UDP port and return ``(protocol, udp_port)``."""
log = logger or _logger
proto = ObpFanInProtocol(demux)
udp_port = reactor.listenUDP(listen_port, proto, 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="OBP_PROXY", logger=log)
return proto, udp_port
__all__ = [
"InProcessObpSink",
"ObpBridgeEntry",
"ObpBridgeRegistry",
"ObpFanInDemux",
"ObpFanInProtocol",
"ObpIngressReplyTransport",
"listen_obp_fanin",
]

@ -0,0 +1,258 @@
# ADN DMR Peer Server - infrastructure proxy obp runtime
#
# 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
###############################################################################
"""Wire integrated OBP proxy at startup (composition root wiring)."""
from __future__ import annotations
import logging
from dataclasses import dataclass, field
from typing import Any
from twisted.internet import reactor
from adn_server.application.proxy.deployment import obp_bridge_legacy_listen_port, obp_proxy_enabled
from adn_server.infrastructure.proxy.obp_config import obp_proxy_settings
from adn_server.infrastructure.proxy.obp_fanin import (
InProcessObpSink,
ObpBridgeEntry,
ObpBridgeRegistry,
ObpFanInDemux,
ObpIngressReplyTransport,
listen_obp_fanin,
)
def build_obp_bridge_registry(
config: dict[str, Any],
protocols: dict[str, Any],
*,
bind_legacy_ports: bool,
listen_port: int,
primary_transport: Any,
) -> ObpBridgeRegistry:
"""Register enabled OPENBRIDGE systems for fan-in demux."""
registry = ObpBridgeRegistry()
systems = config.get("SYSTEMS", {})
if not isinstance(systems, dict):
return registry
for name, sys_cfg in systems.items():
if not isinstance(sys_cfg, dict):
continue
if not sys_cfg.get("ENABLED", True):
continue
if sys_cfg.get("MODE") != "OPENBRIDGE":
continue
proto = protocols.get(name)
if proto is None:
continue
network_id = sys_cfg.get("NETWORK_ID")
if not isinstance(network_id, bytes) or len(network_id) != 4:
continue
legacy_port = obp_bridge_legacy_listen_port(
sys_cfg,
listen_port=listen_port,
bind_legacy_ports=bind_legacy_ports,
)
reply = ObpIngressReplyTransport(primary_transport)
passphrase = sys_cfg.get("PASSPHRASE") or b""
if isinstance(passphrase, str):
passphrase = (passphrase.strip().encode("utf-8") + b"\x00" * 20)[:20]
entry = ObpBridgeEntry(
system_name=name,
network_id=network_id,
passphrase=passphrase,
sink=InProcessObpSink(proto),
reply_transport=reply,
legacy_port=legacy_port if legacy_port and legacy_port > 0 else None,
)
registry.register(entry)
proto.transport = reply # type: ignore[assignment]
start = getattr(proto, "startProtocol", None)
if callable(start):
start()
return registry
@dataclass
class ObpProxyServiceState:
"""Live OBP proxy handles (for shutdown / reload)."""
demux: ObpFanInDemux
registry: ObpBridgeRegistry
udp_ports: list[Any] = field(default_factory=list)
listen_port: int = 62032
listen_ip: str = ""
bind_legacy_ports: bool = True
_runtime: dict[str, Any] = field(default_factory=dict)
def stop(self) -> list[Any]:
"""Stop all UDP listeners. Returns Deferred list when ports were bound."""
deferreds: list[Any] = []
for port in self.udp_ports:
if port is not None:
deferreds.append(port.stopListening())
self.udp_ports.clear()
return deferreds
def _start_listeners(
state: ObpProxyServiceState,
config: dict[str, Any],
protocols: dict[str, Any],
runtime: dict[str, Any],
*,
logger: logging.Logger,
) -> None:
primary_proto, primary_port = listen_obp_fanin(
reactor,
runtime["listen_ip"],
runtime["listen_port"],
state.demux,
config=config,
logger=logger,
)
state.udp_ports.append(primary_port)
primary_transport = primary_proto.transport
if primary_transport is None:
raise RuntimeError("(OBP_PROXY) fan-in transport missing after bind")
state.registry = build_obp_bridge_registry(
config,
protocols,
bind_legacy_ports=runtime["bind_legacy_ports"],
listen_port=runtime["listen_port"],
primary_transport=primary_transport,
)
state.demux._registry = state.registry # noqa: SLF001
if runtime["bind_legacy_ports"]:
seen_ports: set[int] = {runtime["listen_port"]}
for entry in state.registry.bridges.values():
if entry.legacy_port is None or entry.legacy_port in seen_ports:
continue
seen_ports.add(entry.legacy_port)
sys_cfg = config.get("SYSTEMS", {}).get(entry.system_name, {})
bind_ip = str(sys_cfg.get("_REPORT_BIND_IP") or runtime["listen_ip"] or "")
_, legacy_port = listen_obp_fanin(
reactor,
bind_ip,
entry.legacy_port,
state.demux,
config=config,
logger=logger,
)
state.udp_ports.append(legacy_port)
logger.info(
"(OBP_PROXY) Legacy port %s:%s -> %s",
bind_ip or "*",
entry.legacy_port,
entry.system_name,
)
def start_obp_proxy_service(
config: dict[str, Any],
protocols: dict[str, Any],
*,
logger: logging.Logger,
) -> ObpProxyServiceState:
"""Start OBP fan-in and wire inject-only OPENBRIDGE instances."""
if not obp_proxy_enabled(config):
raise RuntimeError("(OBP_PROXY) start_obp_proxy_service called but OBP_PROXY is disabled")
runtime = obp_proxy_settings(config)
registry = ObpBridgeRegistry()
demux = ObpFanInDemux(registry, debug=runtime["debug"], logger=logger)
state = ObpProxyServiceState(
demux=demux,
registry=registry,
listen_port=runtime["listen_port"],
listen_ip=runtime["listen_ip"],
bind_legacy_ports=runtime["bind_legacy_ports"],
_runtime=dict(runtime),
)
_start_listeners(state, config, protocols, runtime, logger=logger)
bridge_count = len(state.registry.bridges)
logger.info(
"(OBP_PROXY) Fan-in on %s:%s (%s bridge(s), BIND_LEGACY_PORTS=%s)",
runtime["listen_ip"] or "*",
runtime["listen_port"],
bridge_count,
runtime["bind_legacy_ports"],
)
if bridge_count == 0:
logger.warning("(OBP_PROXY) No enabled OPENBRIDGE systems registered")
return state
def apply_obp_proxy_config_reload(
state: ObpProxyServiceState,
config: dict[str, Any],
protocols: dict[str, Any],
*,
logger: logging.Logger,
) -> None:
"""Hot-apply OBP_PROXY settings on SIGHUP without rebinding listeners."""
incoming = obp_proxy_settings(config)
bind_changed = (
state.listen_port != incoming["listen_port"]
or state.listen_ip != incoming["listen_ip"]
or state.bind_legacy_ports != incoming["bind_legacy_ports"]
)
state._runtime.update(incoming)
state.demux.debug = bool(incoming["debug"])
if state.udp_ports:
primary_proto = getattr(state.udp_ports[0], "protocol", None)
transport = getattr(primary_proto, "transport", None) if primary_proto is not None else None
if transport is not None:
state.registry = build_obp_bridge_registry(
config,
protocols,
bind_legacy_ports=incoming["bind_legacy_ports"],
listen_port=incoming["listen_port"],
primary_transport=transport,
)
state.demux._registry = state.registry # noqa: SLF001
if bind_changed:
logger.warning(
"(CONFIG-RELOAD) OBP_PROXY bind change ignored at runtime "
"(still listening on %s:%s BIND_LEGACY_PORTS=%s); restart adn-server to apply "
"%s:%s BIND_LEGACY_PORTS=%s",
state.listen_ip or "*",
state.listen_port,
state.bind_legacy_ports,
incoming["listen_ip"] or "*",
incoming["listen_port"],
incoming["bind_legacy_ports"],
)
logger.debug(
"(CONFIG-RELOAD) OBP proxy settings applied (%s bridge(s))",
len(state.registry.bridges),
)
__all__ = [
"ObpProxyServiceState",
"apply_obp_proxy_config_reload",
"build_obp_bridge_registry",
"start_obp_proxy_service",
]

@ -39,7 +39,6 @@ from twisted.internet import reactor, task
from twisted.internet.protocol import DatagramProtocol
from ...application.proxy.deployment import is_proxy_inject_only
from ...application.server_voice import all_server_voice_ids
from ...application.routing.downlink import (
DownlinkContext,
iter_downlink_voice_slots,
@ -74,6 +73,7 @@ from ...application.routing.peer_downlink_index import (
count_connected_peers,
invalidate_peer_options_cache,
)
from ...application.server_voice import all_server_voice_ids
from ...domain import bytes_3, bytes_4, int_id
from ...domain.dmr import decode
from ...domain.dmr.const import LC_OPT
@ -275,6 +275,9 @@ class HBPProtocol(DatagramProtocol):
def startProtocol(self) -> None:
if self._config.get("MODE") == "OPENBRIDGE":
if getattr(self, "_obp_protocol_started", False):
return
self._obp_protocol_started = True
logger.info(
"(%s) Starting OBP. TARGET_IP: %s, TARGET_PORT: %s",
self._system,

@ -103,9 +103,71 @@ def test_dashboard_state_includes_enabled_openbridge():
assert obp["mode"] == "OPENBRIDGE"
assert obp["network_id"] == 73010
assert obp["enhanced_obp"] is True
assert obp["connected"] is False
assert obp["streams"] == {}
def test_openbridge_connected_true_when_bcka_fresh() -> None:
now = 1000.0
systems = {
"OBP-CL": {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"NETWORK_ID": 73010,
"ENHANCED_OBP": True,
"_bcka": now - 10,
"PEERS": {},
},
}
state = build_dashboard_state(systems, ts=now)
assert state["ctable"]["OPENBRIDGES"]["OBP-CL"]["connected"] is True
def test_openbridge_connected_false_when_bcka_stale() -> None:
now = 1000.0
systems = {
"OBP-CL": {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"NETWORK_ID": 73010,
"ENHANCED_OBP": True,
"_bcka": now - 61,
"PEERS": {},
},
}
state = build_dashboard_state(systems, ts=now)
assert state["ctable"]["OPENBRIDGES"]["OBP-CL"]["connected"] is False
def test_openbridge_connected_false_when_bcka_missing() -> None:
now = 1000.0
systems = {
"OBP-CL": {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"NETWORK_ID": 73010,
"ENHANCED_OBP": True,
"PEERS": {},
},
}
state = build_dashboard_state(systems, ts=now)
assert state["ctable"]["OPENBRIDGES"]["OBP-CL"]["connected"] is False
def test_openbridge_connected_omitted_without_enhanced_obp() -> None:
systems = {
"OBP-CL": {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"NETWORK_ID": 73010,
"PEERS": {},
},
}
state = build_dashboard_state(systems)
obp = state["ctable"]["OPENBRIDGES"]["OBP-CL"]
assert "connected" not in obp
def test_dashboard_state_includes_connected_upstream_peer():
systems = {
"XLX-730": {
@ -150,7 +212,7 @@ def test_build_dashboard_state_matches_example(validator: jsonschema.Draft202012
"MODE": "MASTER",
"ENABLED": True,
"IP": "10.0.0.2",
"PORT": 62032,
"PORT": 62040,
"PEERS": {
bytes_4(3120002): {
"CONNECTION": "YES",

@ -53,9 +53,9 @@ def test_attach_creates_session_and_refreshes_client() -> None:
slot = first.value
assert slot.client == ClientEndpoint(host="10.0.0.1", port=62031)
again = svc.attach_client(_PEER_A, "10.0.0.2", 62032)
again = svc.attach_client(_PEER_A, "10.0.0.2", 62040)
assert isinstance(again, Success)
assert again.value.client == ClientEndpoint(host="10.0.0.2", port=62032)
assert again.value.client == ClientEndpoint(host="10.0.0.2", port=62040)
assert len(svc.list_slots()) == 1
@ -117,7 +117,7 @@ def test_ip_blacklist_blocks_attach() -> None:
def test_attach_client_assigns_report_slot() -> None:
svc = _service(max_peers=4)
first = svc.attach_client(_PEER_A, "10.0.0.1", 62031)
second = svc.attach_client(_PEER_B, "10.0.0.2", 62032)
second = svc.attach_client(_PEER_B, "10.0.0.2", 62040)
assert isinstance(first, Success)
assert isinstance(second, Success)
assert first.value.report_slot == 0

@ -75,7 +75,7 @@ def voice_master_scenario(tg: int = 91) -> tuple[DeterministicScenario, FakeMast
master = FakeMasterForVoice("MASTER-A")
config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"}
config["SYSTEMS"]["MASTER-A"]["PEERS"] = {
"1001": {"CALLSIGN": "TEST", "IP": "127.0.0.1", "PORT": 62032},
"1001": {"CALLSIGN": "TEST", "IP": "127.0.0.1", "PORT": 62040},
}
bridges = active_routing_table(tg, (("MASTER-A", 2), ("MASTER-B", 2)))
scenario = DeterministicScenario(config=config, routing_table=bridges)

@ -89,3 +89,65 @@ def test_collect_findings_peer_mesh_protocol() -> None:
findings = collect_findings(config, project_root=".", config_path="cfg.yaml")
peer_msgs = [f.message for f in findings if f.section == "peer"]
assert any("MESH_PROTOCOL=dmre_v5" in m for m in peer_msgs)
def test_collect_findings_obp_per_bridge_migration() -> None:
"""7301-style: OBP-CL2 on fan-in 62032; another bridge keeps legacy 62999."""
config = {
"GLOBAL": {"SERVER_ID": 7301},
"REPORTS": {"REPORT": False},
"PROXY": {"LISTEN_PORT": 62031, "TARGET_SYSTEM": "HOTSPOT"},
"OBP_PROXY": {"ENABLED": True, "LISTEN_PORT": 62032, "BIND_LEGACY_PORTS": True},
"SYSTEMS": {
"HOTSPOT": {"MODE": "MASTER", "ENABLED": True, "MAX_PEERS": 1},
"OBP-CL2": {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"PORT": 62032,
"NETWORK_ID": 7302,
"PASSPHRASE": "dev",
"TARGET_IP": "44.31.61.68",
"TARGET_PORT": 62032,
},
"OBP-EU": {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"PORT": 62999,
"NETWORK_ID": 73045,
"PASSPHRASE": "dev",
"TARGET_IP": "10.0.0.1",
"TARGET_PORT": 62044,
},
},
}
findings = collect_findings(config, project_root=".", config_path="cfg.yaml")
messages = [f.message for f in findings]
assert any("OBP-CL2" in m and "fan-in only" in m for m in messages)
assert any("OBP-EU" in m and "legacy UDP" in m and ":62999" in m for m in messages)
def test_collect_findings_obp_fanin_only_with_remote_target() -> None:
"""7302-style doctor: inject-only local; TARGET_PORT points at peer OBP_PROXY fan-in."""
config = {
"GLOBAL": {"SERVER_ID": 7302},
"REPORTS": {"REPORT": False},
"PROXY": {"LISTEN_PORT": 62031, "TARGET_SYSTEM": "HOTSPOT"},
"OBP_PROXY": {"ENABLED": True, "LISTEN_PORT": 62032, "BIND_LEGACY_PORTS": False},
"SYSTEMS": {
"HOTSPOT": {"MODE": "MASTER", "ENABLED": True, "MAX_PEERS": 1},
"OBP-CL": {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"PORT": 62032,
"NETWORK_ID": 73010,
"PASSPHRASE": "dev",
"TARGET_IP": "44.31.61.66",
"TARGET_PORT": 62032,
},
},
}
findings = collect_findings(config, project_root=".", config_path="cfg.yaml")
messages = [f.message for f in findings]
assert any("OBP-CL" in m and "fan-in only" in m for m in messages)
assert any("target 44.31.61.66:62032" in m for m in messages)
assert any("OBP_PROXY UDP" in m and ":62032" in m for m in messages)

@ -44,7 +44,7 @@ def test_echo_9990_only_to_originating_peer() -> None:
peer_b = bytes_4(730039210)
proto._peers = {
peer_a: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62031)},
peer_b: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62032)},
peer_b: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62040)},
}
# Simulate: peer_a transmitted to TG 9990 on slot 2 (RX_PEER set by dmrd_received).
proto.STATUS[2] = {
@ -70,7 +70,7 @@ def test_echo_9990_no_fuzzy_match_when_no_rx_peer() -> None:
peer_b = bytes_4(730039210)
proto._peers = {
peer_a: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62031)},
peer_b: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62032)},
peer_b: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62040)},
}
# RX_PEER not set (stale / different TG on slot)
proto.STATUS[2] = {
@ -110,7 +110,7 @@ def test_special_tg_9991_same_isolation() -> None:
peer_b = bytes_4(730039210)
proto._peers = {
peer_a: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62031)},
peer_b: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62032)},
peer_b: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62040)},
}
proto.STATUS[2] = {
"RX_PEER": peer_b,

@ -0,0 +1,357 @@
# ADN DMR Peer Server - tests infrastructure obp proxy
#
# 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
###############################################################################
"""OBP_PROXY configuration and fan-in demux tests."""
from __future__ import annotations
import pytest
from tests.conftest import minimal_valid_config
from adn_server.application.proxy.deployment import (
config_has_enabled_openbridge,
is_obp_proxy_managed,
normalize_obp_proxy_targets,
obp_bridge_legacy_listen_port,
obp_proxy_bind_legacy_ports,
obp_proxy_enabled,
)
from adn_server.domain import bytes_4
from adn_server.domain.errors import ConfigError
from adn_server.infrastructure.config_validator import validate_config
from adn_server.infrastructure.hbp_constants import DMRD
from adn_server.infrastructure.mesh.obp_v1 import build_bcka, build_dmrd_v1
from adn_server.infrastructure.proxy.obp_config import obp_proxy_settings
from adn_server.infrastructure.proxy.obp_fanin import (
InProcessObpSink,
ObpBridgeEntry,
ObpBridgeRegistry,
ObpFanInDemux,
ObpIngressReplyTransport,
)
from adn_server.infrastructure.proxy.obp_runtime import build_obp_bridge_registry
_PASS = b"test-passphrase\x00\x00\x00\x00\x00\x00"
_NETWORK = bytes_4(73044)
_ADDR = ("10.0.0.9", 62044)
class _RecordingTransport:
def __init__(self) -> None:
self.sent: list[tuple[bytes, tuple[str, int]]] = []
self.port = 62032
def write(self, data: bytes, addr: tuple[str, int]) -> None:
self.sent.append((data, addr))
def getHost(self) -> _RecordingTransport:
return self
class _RecordingObp:
def __init__(self) -> None:
self.packets: list[tuple[bytes, tuple[str, int]]] = []
def _obp_datagram_received(self, data: bytes, sockaddr: tuple[str, int]) -> None:
self.packets.append((data, sockaddr))
def _sample_dmr_voice() -> bytes:
return b"".join(
[
DMRD,
bytes([1]),
bytes_4(1001)[1:4],
bytes_4(52090)[1:4],
bytes_4(1),
bytes([0x10]),
bytes_4(0xAABBCCDD),
b"\x00" * 33,
]
)
def _obp_config(*, bind_legacy: bool = True) -> dict:
return {
"GLOBAL": {"SERVER_ID": 73010},
"DATABASE": {
"DB_SERVER": "localhost",
"DB_USERNAME": "hbmon",
"DB_PASSWORD": "secret",
"DB_NAME": "hbmon",
"DB_PORT": 3306,
},
"PROXY": {"LISTEN_PORT": 62031, "TARGET_SYSTEM": "HOTSPOT"},
"OBP_PROXY": {
"ENABLED": True,
"LISTEN_PORT": 62032,
"BIND_LEGACY_PORTS": bind_legacy,
},
"SYSTEMS": {
"HOTSPOT": {"MODE": "MASTER", "ENABLED": True, "MAX_PEERS": 1},
"OBP-CL": {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"PORT": 62044,
"NETWORK_ID": 73044,
"PASSPHRASE": "test-passphrase",
"TARGET_IP": "127.0.0.1",
"TARGET_PORT": 62030,
},
},
}
def test_obp_proxy_disabled_when_block_absent_and_no_openbridge() -> None:
config = minimal_valid_config()
assert not config_has_enabled_openbridge(config)
assert not obp_proxy_enabled(config)
def test_obp_proxy_defaults_when_block_absent_with_openbridge() -> None:
config = _obp_config()
del config["OBP_PROXY"]
assert obp_proxy_enabled(config)
assert obp_proxy_bind_legacy_ports(config)
settings = obp_proxy_settings(config)
assert settings == {
"enabled": True,
"listen_port": 62032,
"listen_ip": "",
"bind_legacy_ports": True,
"debug": False,
}
normalize_obp_proxy_targets(config)
assert "PORT" not in config["SYSTEMS"]["OBP-CL"]
assert config["SYSTEMS"]["OBP-CL"]["_REPORT_PORT"] == 62044
def test_obp_proxy_explicit_disable_with_openbridge() -> None:
config = _obp_config()
config["OBP_PROXY"] = {"ENABLED": False}
assert not obp_proxy_enabled(config)
assert not is_obp_proxy_managed(config, "OBP-CL")
def test_validate_obp_proxy_duplicate_legacy_port() -> None:
config = _obp_config()
config["SYSTEMS"]["OBP-EU"] = {
**config["SYSTEMS"]["OBP-CL"],
"NETWORK_ID": 73045,
}
with pytest.raises(ConfigError) as exc:
validate_config(config)
assert "duplicate" in str(exc.value).lower()
def test_validate_obp_per_bridge_fanin_migration() -> None:
"""7301-style: BIND_LEGACY_PORTS true; one bridge on fan-in, another on legacy."""
config = _obp_config(bind_legacy=True)
config["SYSTEMS"]["OBP-CL"]["PORT"] = 62032
config["SYSTEMS"]["OBP-EU"] = {
**config["SYSTEMS"]["OBP-CL"],
"PORT": 62999,
"NETWORK_ID": 73045,
}
validate_config(config)
def test_validate_obp_fanin_only_port_equals_listen_with_remote_target() -> None:
"""7302-style: no legacy bind; PORT=62032 is metadata; TARGET_PORT is remote fan-in."""
config = _obp_config(bind_legacy=False)
config["OBP_PROXY"]["BIND_LEGACY_PORTS"] = False
config["SYSTEMS"]["OBP-CL"]["PORT"] = 62032
config["SYSTEMS"]["OBP-CL"]["TARGET_IP"] = "44.31.61.66"
config["SYSTEMS"]["OBP-CL"]["TARGET_PORT"] = 62032
validate_config(config)
def test_validate_obp_legacy_local_port_with_remote_fanin_target() -> None:
"""7301-style: legacy 62999 locally; mesh to peer fan-in on 62032."""
config = _obp_config(bind_legacy=True)
config["SYSTEMS"]["OBP-CL"]["PORT"] = 62999
config["SYSTEMS"]["OBP-CL"]["TARGET_IP"] = "44.31.61.68"
config["SYSTEMS"]["OBP-CL"]["TARGET_PORT"] = 62032
validate_config(config)
def test_obp_bridge_legacy_listen_port_per_bridge_migration() -> None:
migrated = {"_REPORT_PORT": 62032}
legacy = {"_REPORT_PORT": 62999}
assert obp_bridge_legacy_listen_port(migrated, listen_port=62032, bind_legacy_ports=True) is None
assert obp_bridge_legacy_listen_port(legacy, listen_port=62032, bind_legacy_ports=True) == 62999
def test_obp_proxy_enabled_defaults() -> None:
config = _obp_config()
assert obp_proxy_enabled(config)
assert obp_proxy_bind_legacy_ports(config)
def test_normalize_obp_proxy_targets_omitted_port_uses_fanin() -> None:
config = _obp_config(bind_legacy=True)
del config["SYSTEMS"]["OBP-CL"]["PORT"]
normalize_obp_proxy_targets(config)
assert config["SYSTEMS"]["OBP-CL"]["_REPORT_PORT"] == 62032
assert obp_bridge_legacy_listen_port(
config["SYSTEMS"]["OBP-CL"],
listen_port=62032,
bind_legacy_ports=True,
) is None
def test_is_obp_proxy_managed_only_for_openbridge() -> None:
config = _obp_config()
assert is_obp_proxy_managed(config, "OBP-CL")
assert not is_obp_proxy_managed(config, "HOTSPOT")
def test_normalize_obp_proxy_targets_strips_bind_fields() -> None:
config = _obp_config()
normalize_obp_proxy_targets(config)
obp = config["SYSTEMS"]["OBP-CL"]
assert "PORT" not in obp
assert obp["_REPORT_PORT"] == 62044
def test_validate_obp_proxy_duplicate_network_id() -> None:
config = _obp_config()
config["SYSTEMS"]["OBP-EU"] = {
**config["SYSTEMS"]["OBP-CL"],
"PORT": 62045,
}
with pytest.raises(ConfigError) as exc:
validate_config(config)
assert "NETWORK_ID" in str(exc.value)
def test_validate_obp_proxy_migrated_bridge_port_matches_listen() -> None:
config = _obp_config()
config["OBP_PROXY"]["LISTEN_PORT"] = 62044
config["SYSTEMS"]["OBP-CL"]["PORT"] = 62044
validate_config(config)
def test_obp_proxy_settings_resolved() -> None:
settings = obp_proxy_settings(_obp_config())
assert settings["listen_port"] == 62032
assert settings["bind_legacy_ports"] is True
def test_obp_proxy_settings_default_listen_port() -> None:
config = _obp_config()
del config["OBP_PROXY"]["LISTEN_PORT"]
settings = obp_proxy_settings(config)
assert settings["listen_port"] == 62032
assert settings["enabled"] is True
def test_obp_fanin_demux_by_network_id() -> None:
receiver = _RecordingObp()
transport = _RecordingTransport()
reply = ObpIngressReplyTransport(transport)
registry = ObpBridgeRegistry()
registry.register(
ObpBridgeEntry(
system_name="OBP-CL",
network_id=_NETWORK,
passphrase=_PASS,
sink=InProcessObpSink(receiver),
reply_transport=reply,
)
)
demux = ObpFanInDemux(registry)
wire = build_dmrd_v1(_sample_dmr_voice(), _NETWORK, _PASS)
demux.deliver(wire, _ADDR, local_port=62032, transport=transport)
assert len(receiver.packets) == 1
assert receiver.packets[0][0] == wire
def test_obp_fanin_demux_legacy_port_routes_without_network_id() -> None:
receiver = _RecordingObp()
transport = _RecordingTransport()
transport.port = 62044
reply = ObpIngressReplyTransport(transport)
registry = ObpBridgeRegistry()
registry.register(
ObpBridgeEntry(
system_name="OBP-CL",
network_id=_NETWORK,
passphrase=_PASS,
sink=InProcessObpSink(receiver),
reply_transport=reply,
legacy_port=62044,
)
)
demux = ObpFanInDemux(registry)
wire = build_bcka(_PASS)
demux.deliver(wire, _ADDR, local_port=62044, transport=transport)
assert len(receiver.packets) == 1
def test_obp_fanin_demux_control_on_listen_port() -> None:
receiver = _RecordingObp()
transport = _RecordingTransport()
reply = ObpIngressReplyTransport(transport)
registry = ObpBridgeRegistry()
registry.register(
ObpBridgeEntry(
system_name="OBP-CL",
network_id=_NETWORK,
passphrase=_PASS,
sink=InProcessObpSink(receiver),
reply_transport=reply,
)
)
demux = ObpFanInDemux(registry)
wire = build_bcka(_PASS)
demux.deliver(wire, _ADDR, local_port=62032, transport=transport)
assert len(receiver.packets) == 1
def test_build_obp_bridge_registry_starts_inject_protocol() -> None:
class _InjectProto:
def __init__(self) -> None:
self.transport = None
self.started = False
def startProtocol(self) -> None:
self.started = True
proto = _InjectProto()
config = {
"SYSTEMS": {
"OBP-CL": {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"NETWORK_ID": _NETWORK,
"PASSPHRASE": _PASS,
}
}
}
build_obp_bridge_registry(
config,
{"OBP-CL": proto},
bind_legacy_ports=False,
listen_port=62032,
primary_transport=_RecordingTransport(),
)
assert proto.started is True
assert proto.transport is not None

@ -44,7 +44,7 @@ pytestmark = pytest.mark.integration
_PEER_TX = bytes_4(730039210)
_PEER_RX = bytes_4(730039101)
_ADDR_TX = ("192.168.50.10", 62031)
_ADDR_RX = ("192.168.50.20", 62032)
_ADDR_RX = ("192.168.50.20", 62040)
_EMB_SLICE = slice(116, 148)

@ -70,7 +70,7 @@ def test_group_downlink_allowed_for_shared_base_rf_when_not_tx() -> None:
other: {
"CONNECTION": "YES",
"OPTIONS": b"TS2=730502;",
"SOCKADDR": ("127.0.0.1", 62032),
"SOCKADDR": ("127.0.0.1", 62040),
},
}
foreign_rf = bytes_3(730039265 // 100)

@ -105,7 +105,7 @@ def test_existing_client_refreshes_endpoint_before_inject() -> None:
packet = RPTL + _PEER
protocol.datagramReceived(packet, _CLIENT_ADDR)
inject_spy.reset_mock()
new_addr = ("192.168.1.51", 62032)
new_addr = ("192.168.1.51", 62040)
protocol.datagramReceived(packet, new_addr)
inject_spy.assert_called_once_with(packet, new_addr)
client = proxy.resolve_client(_PEER)

@ -64,7 +64,7 @@ def test_merge_new_master_keeps_yaml_static_tg() -> None:
"MODE": "MASTER",
"ENABLED": True,
"IP": "127.0.0.1",
"PORT": 62032,
"PORT": 62040,
"TS2_STATIC": "52090",
"DEFAULT_UA_TIMER": 10,
}

@ -31,7 +31,7 @@ from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD
def _private_call_scenario(dst_subscriber: int = 7123456) -> tuple[DeterministicScenario, int]:
config = minimal_config(("MASTER-A", "MASTER-B"))
config["SYSTEMS"]["MASTER-A"]["PEERS"] = {
"1001": {"CALLSIGN": "SRC", "IP": "127.0.0.1", "PORT": 62032},
"1001": {"CALLSIGN": "SRC", "IP": "127.0.0.1", "PORT": 62040},
}
config["SYSTEMS"]["MASTER-B"]["PEERS"] = {
"1002": {"CALLSIGN": "DST", "IP": "127.0.0.1", "PORT": 62033},

@ -69,7 +69,7 @@ def _dashboard_systems() -> dict:
"MODE": "MASTER",
"ENABLED": True,
"IP": "10.0.0.2",
"PORT": 62032,
"PORT": 62040,
"PEERS": {
bytes_4(3120002): {
"CONNECTION": "YES",

@ -28,7 +28,7 @@ import time
import wave
from adn_server.infrastructure.voice import tts_engine
from adn_server.infrastructure.voice.tts_engine import DV3K_SAMPLES_PER_FRAME
from adn_server.infrastructure.voice.tts_engine import DV3K_PRODID_REQ, DV3K_SAMPLES_PER_FRAME
def test_text_to_ambe_serializes_parallel_conversions(tmp_path, monkeypatch) -> None:

Loading…
Cancel
Save

Powered by TurnKey Linux.