From b27c1663c62cc25b3d4671b7114a1cded43da5c5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Mon, 13 Jul 2026 23:14:48 -0400 Subject: [PATCH 1/2] feat: publish OBP keepalive connected status in dashboard_state Expose a boolean connected field on ENHANCED_OBP openbridge entries so the monitor can show mesh leg health from the same 60s BCKA rule as routing. --- schemas/examples/dashboard_state.json | 1 + schemas/report-v2.json | 6 +- .../application/report/dashboard_state.py | 27 +++++++- tests/application/test_dashboard_state.py | 62 +++++++++++++++++++ 4 files changed, 93 insertions(+), 3 deletions(-) diff --git a/schemas/examples/dashboard_state.json b/schemas/examples/dashboard_state.json index 09ff622..42d94a8 100644 --- a/schemas/examples/dashboard_state.json +++ b/schemas/examples/dashboard_state.json @@ -57,6 +57,7 @@ "ip": "44.31.61.68", "port": 62999, "enhanced_obp": true, + "connected": false, "streams": {} } } diff --git a/schemas/report-v2.json b/schemas/report-v2.json index 0f5916f..34e1e4d 100644 --- a/schemas/report-v2.json +++ b/schemas/report-v2.json @@ -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": { diff --git a/src/adn_server/application/report/dashboard_state.py b/src/adn_server/application/report/dashboard_state.py index 237acf3..1764eda 100644 --- a/src/adn_server/application/report/dashboard_state.py +++ b/src/adn_server/application/report/dashboard_state.py @@ -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) diff --git a/tests/application/test_dashboard_state.py b/tests/application/test_dashboard_state.py index fb14460..0b388d4 100644 --- a/tests/application/test_dashboard_state.py +++ b/tests/application/test_dashboard_state.py @@ -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": { From 770523cc3e22a2815009b10077b540f66fac1d73 Mon Sep 17 00:00:00 2001 From: ce5rpy <169016246+ce5rpy@users.noreply.github.com> Date: Tue, 14 Jul 2026 00:16:55 -0400 Subject: [PATCH 2/2] Merge pull request #50 from ce5rpy/feat/obp-proxy-fan-in fix: OBP_PROXY fan-in for centralized OpenBridge ingress --- adn-server.example.yaml | 10 + docs/en/README.md | 3 +- docs/en/server/development/performance.md | 118 ++---- docs/en/server/user-guide/configuration.md | 15 + docs/en/server/user-guide/introduction.md | 3 + docs/en/server/user-guide/obp-proxy.md | 51 +++ docs/es/README.md | 3 +- docs/es/server/development/performance.md | 117 ++---- docs/es/server/user-guide/configuration.md | 15 + docs/es/server/user-guide/introduction.md | 3 + docs/es/server/user-guide/obp-proxy.md | 49 +++ mkdocs.es.yml | 1 + mkdocs.yml | 1 + schemas/examples/dashboard_state.json | 2 +- src/adn_server/application/proxy/__init__.py | 16 +- .../application/proxy/deployment.py | 99 +++++ src/adn_server/application/voice_use_cases.py | 16 +- .../infrastructure/bootstrap/peer_server.py | 50 ++- .../infrastructure/config_reload.py | 3 +- .../infrastructure/config_validator.py | 73 ++++ src/adn_server/infrastructure/doctor.py | 69 +++- .../infrastructure/proxy/__init__.py | 10 + .../infrastructure/proxy/obp_config.py | 41 ++ .../infrastructure/proxy/obp_fanin.py | 225 +++++++++++ .../infrastructure/proxy/obp_runtime.py | 258 +++++++++++++ .../twisted_adapters/udp_hbp.py | 5 +- tests/application/test_dashboard_state.py | 2 +- tests/application/test_proxy_use_cases.py | 6 +- tests/harness/voice_helpers.py | 2 +- tests/infrastructure/test_doctor.py | 62 +++ .../test_echo_special_tg_downlink.py | 6 +- tests/infrastructure/test_obp_proxy.py | 357 ++++++++++++++++++ tests/infrastructure/test_proxy_repeat_e2e.py | 2 +- .../test_self_rf_downlink_filter.py | 2 +- tests/infrastructure/test_udp_fanin.py | 2 +- tests/routing/test_config_reload.py | 2 +- tests/routing/test_private_voice.py | 2 +- tests/schemas/test_report_wire_contract.py | 2 +- tests/voice/test_tts_engine_serial.py | 2 +- 39 files changed, 1496 insertions(+), 209 deletions(-) create mode 100644 docs/en/server/user-guide/obp-proxy.md create mode 100644 docs/es/server/user-guide/obp-proxy.md create mode 100644 src/adn_server/infrastructure/proxy/obp_config.py create mode 100644 src/adn_server/infrastructure/proxy/obp_fanin.py create mode 100644 src/adn_server/infrastructure/proxy/obp_runtime.py create mode 100644 tests/infrastructure/test_obp_proxy.py diff --git a/adn-server.example.yaml b/adn-server.example.yaml index 953ed2c..cb46178 100644 --- a/adn-server.example.yaml +++ b/adn-server.example.yaml @@ -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 diff --git a/docs/en/README.md b/docs/en/README.md index 1baeb9b..0a0bf90 100644 --- a/docs/en/README.md +++ b/docs/en/README.md @@ -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) | diff --git a/docs/en/server/development/performance.md b/docs/en/server/development/performance.md index 5dfdc50..7bad48a 100644 --- a/docs/en/server/development/performance.md +++ b/docs/en/server/development/performance.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 diff --git a/docs/en/server/user-guide/configuration.md b/docs/en/server/user-guide/configuration.md index 728b7c8..4788510 100644 --- a/docs/en/server/user-guide/configuration.md +++ b/docs/en/server/user-guide/configuration.md @@ -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. diff --git a/docs/en/server/user-guide/introduction.md b/docs/en/server/user-guide/introduction.md index 34dc643..9742761 100644 --- a/docs/en/server/user-guide/introduction.md +++ b/docs/en/server/user-guide/introduction.md @@ -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. diff --git a/docs/en/server/user-guide/obp-proxy.md b/docs/en/server/user-guide/obp-proxy.md new file mode 100644 index 0000000..b79b41f --- /dev/null +++ b/docs/en/server/user-guide/obp-proxy.md @@ -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). diff --git a/docs/es/README.md b/docs/es/README.md index 8e1b060..728bd2b 100644 --- a/docs/es/README.md +++ b/docs/es/README.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) | diff --git a/docs/es/server/development/performance.md b/docs/es/server/development/performance.md index 41a2179..bdfc65a 100644 --- a/docs/es/server/development/performance.md +++ b/docs/es/server/development/performance.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 diff --git a/docs/es/server/user-guide/configuration.md b/docs/es/server/user-guide/configuration.md index c1041df..9e87ec5 100644 --- a/docs/es/server/user-guide/configuration.md +++ b/docs/es/server/user-guide/configuration.md @@ -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. diff --git a/docs/es/server/user-guide/introduction.md b/docs/es/server/user-guide/introduction.md index ddb985b..e90e4fb 100644 --- a/docs/es/server/user-guide/introduction.md +++ b/docs/es/server/user-guide/introduction.md @@ -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. diff --git a/docs/es/server/user-guide/obp-proxy.md b/docs/es/server/user-guide/obp-proxy.md new file mode 100644 index 0000000..d3f4337 --- /dev/null +++ b/docs/es/server/user-guide/obp-proxy.md @@ -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). diff --git a/mkdocs.es.yml b/mkdocs.es.yml index d4a444c..7fb3f8e 100644 --- a/mkdocs.es.yml +++ b/mkdocs.es.yml @@ -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: diff --git a/mkdocs.yml b/mkdocs.yml index 8113f63..1d82929 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -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: diff --git a/schemas/examples/dashboard_state.json b/schemas/examples/dashboard_state.json index 42d94a8..5857016 100644 --- a/schemas/examples/dashboard_state.json +++ b/schemas/examples/dashboard_state.json @@ -31,7 +31,7 @@ "MASTER-B": { "mode": "MASTER", "ip": "10.0.0.2", - "port": 62032, + "port": 62040, "peers": { "3120002": { "id": 3120002, diff --git a/src/adn_server/application/proxy/__init__.py b/src/adn_server/application/proxy/__init__.py index c330bb5..137f0f0 100644 --- a/src/adn_server/application/proxy/__init__.py +++ b/src/adn_server/application/proxy/__init__.py @@ -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", ] diff --git a/src/adn_server/application/proxy/deployment.py b/src/adn_server/application/proxy/deployment.py index c6ae37f..ef1d1b1 100644 --- a/src/adn_server/application/proxy/deployment.py +++ b/src/adn_server/application/proxy/deployment.py @@ -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) diff --git a/src/adn_server/application/voice_use_cases.py b/src/adn_server/application/voice_use_cases.py index 2408eb5..ca82b9d 100644 --- a/src/adn_server/application/voice_use_cases.py +++ b/src/adn_server/application/voice_use_cases.py @@ -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]: diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index 8143464..7717d3b 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -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) diff --git a/src/adn_server/infrastructure/config_reload.py b/src/adn_server/infrastructure/config_reload.py index 1ef9112..72ba54f 100644 --- a/src/adn_server/infrastructure/config_reload.py +++ b/src/adn_server/infrastructure/config_reload.py @@ -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 diff --git a/src/adn_server/infrastructure/config_validator.py b/src/adn_server/infrastructure/config_validator.py index 13e4728..f10b021 100644 --- a/src/adn_server/infrastructure/config_validator.py +++ b/src/adn_server/infrastructure/config_validator.py @@ -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) diff --git a/src/adn_server/infrastructure/doctor.py b/src/adn_server/infrastructure/doctor.py index 1618260..78a11c5 100644 --- a/src/adn_server/infrastructure/doctor.py +++ b/src/adn_server/infrastructure/doctor.py @@ -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"), diff --git a/src/adn_server/infrastructure/proxy/__init__.py b/src/adn_server/infrastructure/proxy/__init__.py index 487f147..08a80bc 100644 --- a/src/adn_server/infrastructure/proxy/__init__.py +++ b/src/adn_server/infrastructure/proxy/__init__.py @@ -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", ] diff --git a/src/adn_server/infrastructure/proxy/obp_config.py b/src/adn_server/infrastructure/proxy/obp_config.py new file mode 100644 index 0000000..37c833b --- /dev/null +++ b/src/adn_server/infrastructure/proxy/obp_config.py @@ -0,0 +1,41 @@ +# ADN DMR Peer Server - infrastructure proxy obp config +# +# 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 +############################################################################### + +"""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")), + } diff --git a/src/adn_server/infrastructure/proxy/obp_fanin.py b/src/adn_server/infrastructure/proxy/obp_fanin.py new file mode 100644 index 0000000..c5e7e4d --- /dev/null +++ b/src/adn_server/infrastructure/proxy/obp_fanin.py @@ -0,0 +1,225 @@ +# ADN DMR Peer Server - infrastructure proxy obp fanin +# +# 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 +############################################################################### + +"""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", +] diff --git a/src/adn_server/infrastructure/proxy/obp_runtime.py b/src/adn_server/infrastructure/proxy/obp_runtime.py new file mode 100644 index 0000000..99e61c0 --- /dev/null +++ b/src/adn_server/infrastructure/proxy/obp_runtime.py @@ -0,0 +1,258 @@ +# ADN DMR Peer Server - infrastructure proxy obp runtime +# +# 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 +############################################################################### + +"""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", +] diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 3350c4a..35cae61 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -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, diff --git a/tests/application/test_dashboard_state.py b/tests/application/test_dashboard_state.py index 0b388d4..3c332df 100644 --- a/tests/application/test_dashboard_state.py +++ b/tests/application/test_dashboard_state.py @@ -212,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", diff --git a/tests/application/test_proxy_use_cases.py b/tests/application/test_proxy_use_cases.py index 8d8e6eb..42bb0cc 100644 --- a/tests/application/test_proxy_use_cases.py +++ b/tests/application/test_proxy_use_cases.py @@ -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 diff --git a/tests/harness/voice_helpers.py b/tests/harness/voice_helpers.py index 14382f9..b0d586e 100644 --- a/tests/harness/voice_helpers.py +++ b/tests/harness/voice_helpers.py @@ -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) diff --git a/tests/infrastructure/test_doctor.py b/tests/infrastructure/test_doctor.py index 531ee9e..d0eac06 100644 --- a/tests/infrastructure/test_doctor.py +++ b/tests/infrastructure/test_doctor.py @@ -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) diff --git a/tests/infrastructure/test_echo_special_tg_downlink.py b/tests/infrastructure/test_echo_special_tg_downlink.py index 987cc53..31eb3f0 100644 --- a/tests/infrastructure/test_echo_special_tg_downlink.py +++ b/tests/infrastructure/test_echo_special_tg_downlink.py @@ -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, diff --git a/tests/infrastructure/test_obp_proxy.py b/tests/infrastructure/test_obp_proxy.py new file mode 100644 index 0000000..846e551 --- /dev/null +++ b/tests/infrastructure/test_obp_proxy.py @@ -0,0 +1,357 @@ +# ADN DMR Peer Server - tests infrastructure obp proxy +# +# 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 +############################################################################### + +"""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 diff --git a/tests/infrastructure/test_proxy_repeat_e2e.py b/tests/infrastructure/test_proxy_repeat_e2e.py index 00bc73e..7ed70e6 100644 --- a/tests/infrastructure/test_proxy_repeat_e2e.py +++ b/tests/infrastructure/test_proxy_repeat_e2e.py @@ -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) diff --git a/tests/infrastructure/test_self_rf_downlink_filter.py b/tests/infrastructure/test_self_rf_downlink_filter.py index e874552..165b52f 100644 --- a/tests/infrastructure/test_self_rf_downlink_filter.py +++ b/tests/infrastructure/test_self_rf_downlink_filter.py @@ -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) diff --git a/tests/infrastructure/test_udp_fanin.py b/tests/infrastructure/test_udp_fanin.py index 7e50c4e..23e42cc 100644 --- a/tests/infrastructure/test_udp_fanin.py +++ b/tests/infrastructure/test_udp_fanin.py @@ -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) diff --git a/tests/routing/test_config_reload.py b/tests/routing/test_config_reload.py index 5521bcb..7abf490 100644 --- a/tests/routing/test_config_reload.py +++ b/tests/routing/test_config_reload.py @@ -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, } diff --git a/tests/routing/test_private_voice.py b/tests/routing/test_private_voice.py index c52b6ca..9381e5b 100644 --- a/tests/routing/test_private_voice.py +++ b/tests/routing/test_private_voice.py @@ -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}, diff --git a/tests/schemas/test_report_wire_contract.py b/tests/schemas/test_report_wire_contract.py index 3fda853..a29c130 100644 --- a/tests/schemas/test_report_wire_contract.py +++ b/tests/schemas/test_report_wire_contract.py @@ -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", diff --git a/tests/voice/test_tts_engine_serial.py b/tests/voice/test_tts_engine_serial.py index 742dcc6..9e82582 100644 --- a/tests/voice/test_tts_engine_serial.py +++ b/tests/voice/test_tts_engine_serial.py @@ -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: