Merge pull request #50 from ce5rpy/feat/obp-proxy-fan-in
fix: OBP_PROXY fan-in for centralized OpenBridge ingresspull/51/head
parent
13e621b2da
commit
770523cc3e
@ -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).
|
||||
@ -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).
|
||||
@ -0,0 +1,41 @@
|
||||
# ADN DMR Peer Server - infrastructure proxy obp config
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""OBP_PROXY runtime settings from config dict (infrastructure; no business rules)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
from adn_server.application.proxy.deployment import obp_proxy_bind_legacy_ports, obp_proxy_enabled
|
||||
|
||||
|
||||
def obp_proxy_settings(config: dict[str, Any]) -> dict[str, Any]:
|
||||
"""Resolved OBP_PROXY runtime settings with defaults (block may be absent)."""
|
||||
block = config.get("OBP_PROXY", {})
|
||||
if not isinstance(block, dict):
|
||||
block = {}
|
||||
return {
|
||||
"enabled": obp_proxy_enabled(config),
|
||||
"listen_port": int(block.get("LISTEN_PORT", 62032)),
|
||||
"listen_ip": str(block.get("LISTEN_IP") or ""),
|
||||
"bind_legacy_ports": obp_proxy_bind_legacy_ports(config),
|
||||
"debug": bool(block.get("DEBUG")),
|
||||
}
|
||||
@ -0,0 +1,225 @@
|
||||
# ADN DMR Peer Server - infrastructure proxy obp fanin
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""UDP fan-in for OPENBRIDGE: demux by NETWORK_ID and optional legacy PORT."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Any, Protocol
|
||||
|
||||
from twisted.internet.protocol import DatagramProtocol
|
||||
|
||||
from adn_server.infrastructure.hbp_constants import BCKA, BCSQ, BCST, BCVE, DMRD, DMRE
|
||||
from adn_server.infrastructure.mesh.obp_v1 import verify_bcka, verify_bcsq, verify_bcst, verify_bcve
|
||||
from adn_server.infrastructure.udp_rcvbuf import apply_udp_rcvbuf, udp_rcvbuf_bytes
|
||||
|
||||
_logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class _DatagramWriter(Protocol):
|
||||
def write(self, data: bytes, addr: tuple[str, int]) -> None:
|
||||
...
|
||||
|
||||
|
||||
class _ObpReceiver(Protocol):
|
||||
def _obp_datagram_received(self, data: bytes, sockaddr: tuple[str, int]) -> None:
|
||||
...
|
||||
|
||||
|
||||
class ObpIngressReplyTransport:
|
||||
"""Route OBP egress through the fan-in socket that last received for this bridge."""
|
||||
|
||||
def __init__(self, fallback: _DatagramWriter) -> None:
|
||||
self._fallback = fallback
|
||||
self._active: _DatagramWriter | None = None
|
||||
|
||||
def note_ingress(self, transport: _DatagramWriter) -> None:
|
||||
self._active = transport
|
||||
|
||||
def write(self, data: bytes, addr: tuple[str, int]) -> None:
|
||||
transport = self._active or self._fallback
|
||||
transport.write(data, addr)
|
||||
|
||||
|
||||
class InProcessObpSink:
|
||||
"""Deliver datagrams to an OPENBRIDGE HBPProtocol without a UDP hop."""
|
||||
|
||||
def __init__(self, hbp: _ObpReceiver) -> None:
|
||||
self._hbp = hbp
|
||||
|
||||
def inject(self, data: bytes, client_addr: tuple[str, int]) -> None:
|
||||
self._hbp._obp_datagram_received(data, client_addr)
|
||||
|
||||
|
||||
@dataclass
|
||||
class ObpBridgeEntry:
|
||||
system_name: str
|
||||
network_id: bytes
|
||||
passphrase: bytes
|
||||
sink: InProcessObpSink
|
||||
reply_transport: ObpIngressReplyTransport
|
||||
legacy_port: int | None = None
|
||||
|
||||
|
||||
@dataclass
|
||||
class ObpBridgeRegistry:
|
||||
"""NETWORK_ID and legacy PORT lookup for OBP fan-in demux."""
|
||||
|
||||
by_network_id: dict[bytes, str] = field(default_factory=dict)
|
||||
by_legacy_port: dict[int, str] = field(default_factory=dict)
|
||||
bridges: dict[str, ObpBridgeEntry] = field(default_factory=dict)
|
||||
|
||||
def register(self, entry: ObpBridgeEntry) -> None:
|
||||
self.bridges[entry.system_name] = entry
|
||||
self.by_network_id[entry.network_id] = entry.system_name
|
||||
if entry.legacy_port is not None:
|
||||
self.by_legacy_port[entry.legacy_port] = entry.system_name
|
||||
|
||||
def clear(self) -> None:
|
||||
self.by_network_id.clear()
|
||||
self.by_legacy_port.clear()
|
||||
self.bridges.clear()
|
||||
|
||||
|
||||
class ObpFanInDemux:
|
||||
"""Shared demux handler for one or more OBP proxy UDP listeners."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
registry: ObpBridgeRegistry,
|
||||
*,
|
||||
debug: bool = False,
|
||||
logger: logging.Logger | None = None,
|
||||
) -> None:
|
||||
self._registry = registry
|
||||
self.debug = debug
|
||||
self._log = logger or _logger
|
||||
|
||||
def deliver(
|
||||
self,
|
||||
data: bytes,
|
||||
addr: tuple[str, int],
|
||||
*,
|
||||
local_port: int,
|
||||
transport: _DatagramWriter,
|
||||
) -> None:
|
||||
host, port = addr
|
||||
system_name = self._registry.by_legacy_port.get(local_port)
|
||||
if system_name is None:
|
||||
system_name = self._lookup_by_network_id(data)
|
||||
if system_name is None:
|
||||
system_name = self._lookup_control(data)
|
||||
if system_name is None:
|
||||
if self.debug:
|
||||
self._log.debug(
|
||||
"(OBP_PROXY) dropped packet from %s:%s len=%d local_port=%s",
|
||||
host,
|
||||
port,
|
||||
len(data),
|
||||
local_port,
|
||||
)
|
||||
return
|
||||
entry = self._registry.bridges.get(system_name)
|
||||
if entry is None:
|
||||
return
|
||||
if self.debug:
|
||||
self._log.debug(
|
||||
"(OBP_PROXY) RX %s from %s:%s len=%d -> %s",
|
||||
data[:4],
|
||||
host,
|
||||
port,
|
||||
len(data),
|
||||
system_name,
|
||||
)
|
||||
entry.reply_transport.note_ingress(transport)
|
||||
entry.sink.inject(data, addr)
|
||||
|
||||
def _lookup_by_network_id(self, data: bytes) -> str | None:
|
||||
if len(data) < 15:
|
||||
return None
|
||||
opcode = data[:4]
|
||||
if opcode not in (DMRD, DMRE):
|
||||
return None
|
||||
network_id = data[11:15]
|
||||
return self._registry.by_network_id.get(network_id)
|
||||
|
||||
def _lookup_control(self, data: bytes) -> str | None:
|
||||
if len(data) < 4:
|
||||
return None
|
||||
opcode = data[:4]
|
||||
for name, entry in self._registry.bridges.items():
|
||||
passphrase = entry.passphrase
|
||||
if opcode == BCKA and verify_bcka(data, passphrase):
|
||||
return name
|
||||
if opcode == BCSQ and verify_bcsq(data, passphrase) is not None:
|
||||
return name
|
||||
if opcode == BCST and verify_bcst(data, passphrase):
|
||||
return name
|
||||
if opcode == BCVE:
|
||||
ok, _ver = verify_bcve(data, passphrase)
|
||||
if ok:
|
||||
return name
|
||||
return None
|
||||
|
||||
|
||||
class ObpFanInProtocol(DatagramProtocol):
|
||||
"""Thin UDP listener delegating to shared OBP demux."""
|
||||
|
||||
def __init__(self, demux: ObpFanInDemux) -> None:
|
||||
self._demux = demux
|
||||
|
||||
def datagramReceived(self, data: bytes, addr: tuple[str, int]) -> None:
|
||||
transport = self.transport
|
||||
if transport is None:
|
||||
return
|
||||
local_port = int(transport.getHost().port)
|
||||
self._demux.deliver(data, addr, local_port=local_port, transport=transport)
|
||||
|
||||
|
||||
def listen_obp_fanin(
|
||||
reactor: Any,
|
||||
listen_ip: str,
|
||||
listen_port: int,
|
||||
demux: ObpFanInDemux,
|
||||
*,
|
||||
config: dict[str, Any] | None = None,
|
||||
udp_rcvbuf: int | None = None,
|
||||
logger: logging.Logger | None = None,
|
||||
) -> tuple[ObpFanInProtocol, Any]:
|
||||
"""Bind one OBP proxy UDP port and return ``(protocol, udp_port)``."""
|
||||
log = logger or _logger
|
||||
proto = ObpFanInProtocol(demux)
|
||||
udp_port = reactor.listenUDP(listen_port, proto, interface=listen_ip or "0.0.0.0")
|
||||
buf_size = udp_rcvbuf if udp_rcvbuf is not None else udp_rcvbuf_bytes(config)
|
||||
apply_udp_rcvbuf(udp_port.socket, buf_size, label="OBP_PROXY", logger=log)
|
||||
return proto, udp_port
|
||||
|
||||
|
||||
__all__ = [
|
||||
"InProcessObpSink",
|
||||
"ObpBridgeEntry",
|
||||
"ObpBridgeRegistry",
|
||||
"ObpFanInDemux",
|
||||
"ObpFanInProtocol",
|
||||
"ObpIngressReplyTransport",
|
||||
"listen_obp_fanin",
|
||||
]
|
||||
@ -0,0 +1,258 @@
|
||||
# ADN DMR Peer Server - infrastructure proxy obp runtime
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""Wire integrated OBP proxy at startup (composition root wiring)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Any
|
||||
|
||||
from twisted.internet import reactor
|
||||
|
||||
from adn_server.application.proxy.deployment import obp_bridge_legacy_listen_port, obp_proxy_enabled
|
||||
from adn_server.infrastructure.proxy.obp_config import obp_proxy_settings
|
||||
from adn_server.infrastructure.proxy.obp_fanin import (
|
||||
InProcessObpSink,
|
||||
ObpBridgeEntry,
|
||||
ObpBridgeRegistry,
|
||||
ObpFanInDemux,
|
||||
ObpIngressReplyTransport,
|
||||
listen_obp_fanin,
|
||||
)
|
||||
|
||||
|
||||
def build_obp_bridge_registry(
|
||||
config: dict[str, Any],
|
||||
protocols: dict[str, Any],
|
||||
*,
|
||||
bind_legacy_ports: bool,
|
||||
listen_port: int,
|
||||
primary_transport: Any,
|
||||
) -> ObpBridgeRegistry:
|
||||
"""Register enabled OPENBRIDGE systems for fan-in demux."""
|
||||
registry = ObpBridgeRegistry()
|
||||
systems = config.get("SYSTEMS", {})
|
||||
if not isinstance(systems, dict):
|
||||
return registry
|
||||
for name, sys_cfg in systems.items():
|
||||
if not isinstance(sys_cfg, dict):
|
||||
continue
|
||||
if not sys_cfg.get("ENABLED", True):
|
||||
continue
|
||||
if sys_cfg.get("MODE") != "OPENBRIDGE":
|
||||
continue
|
||||
proto = protocols.get(name)
|
||||
if proto is None:
|
||||
continue
|
||||
network_id = sys_cfg.get("NETWORK_ID")
|
||||
if not isinstance(network_id, bytes) or len(network_id) != 4:
|
||||
continue
|
||||
legacy_port = obp_bridge_legacy_listen_port(
|
||||
sys_cfg,
|
||||
listen_port=listen_port,
|
||||
bind_legacy_ports=bind_legacy_ports,
|
||||
)
|
||||
reply = ObpIngressReplyTransport(primary_transport)
|
||||
passphrase = sys_cfg.get("PASSPHRASE") or b""
|
||||
if isinstance(passphrase, str):
|
||||
passphrase = (passphrase.strip().encode("utf-8") + b"\x00" * 20)[:20]
|
||||
entry = ObpBridgeEntry(
|
||||
system_name=name,
|
||||
network_id=network_id,
|
||||
passphrase=passphrase,
|
||||
sink=InProcessObpSink(proto),
|
||||
reply_transport=reply,
|
||||
legacy_port=legacy_port if legacy_port and legacy_port > 0 else None,
|
||||
)
|
||||
registry.register(entry)
|
||||
proto.transport = reply # type: ignore[assignment]
|
||||
start = getattr(proto, "startProtocol", None)
|
||||
if callable(start):
|
||||
start()
|
||||
return registry
|
||||
|
||||
|
||||
@dataclass
|
||||
class ObpProxyServiceState:
|
||||
"""Live OBP proxy handles (for shutdown / reload)."""
|
||||
|
||||
demux: ObpFanInDemux
|
||||
registry: ObpBridgeRegistry
|
||||
udp_ports: list[Any] = field(default_factory=list)
|
||||
listen_port: int = 62032
|
||||
listen_ip: str = ""
|
||||
bind_legacy_ports: bool = True
|
||||
_runtime: dict[str, Any] = field(default_factory=dict)
|
||||
|
||||
def stop(self) -> list[Any]:
|
||||
"""Stop all UDP listeners. Returns Deferred list when ports were bound."""
|
||||
deferreds: list[Any] = []
|
||||
for port in self.udp_ports:
|
||||
if port is not None:
|
||||
deferreds.append(port.stopListening())
|
||||
self.udp_ports.clear()
|
||||
return deferreds
|
||||
|
||||
|
||||
def _start_listeners(
|
||||
state: ObpProxyServiceState,
|
||||
config: dict[str, Any],
|
||||
protocols: dict[str, Any],
|
||||
runtime: dict[str, Any],
|
||||
*,
|
||||
logger: logging.Logger,
|
||||
) -> None:
|
||||
primary_proto, primary_port = listen_obp_fanin(
|
||||
reactor,
|
||||
runtime["listen_ip"],
|
||||
runtime["listen_port"],
|
||||
state.demux,
|
||||
config=config,
|
||||
logger=logger,
|
||||
)
|
||||
state.udp_ports.append(primary_port)
|
||||
primary_transport = primary_proto.transport
|
||||
if primary_transport is None:
|
||||
raise RuntimeError("(OBP_PROXY) fan-in transport missing after bind")
|
||||
|
||||
state.registry = build_obp_bridge_registry(
|
||||
config,
|
||||
protocols,
|
||||
bind_legacy_ports=runtime["bind_legacy_ports"],
|
||||
listen_port=runtime["listen_port"],
|
||||
primary_transport=primary_transport,
|
||||
)
|
||||
state.demux._registry = state.registry # noqa: SLF001
|
||||
|
||||
if runtime["bind_legacy_ports"]:
|
||||
seen_ports: set[int] = {runtime["listen_port"]}
|
||||
for entry in state.registry.bridges.values():
|
||||
if entry.legacy_port is None or entry.legacy_port in seen_ports:
|
||||
continue
|
||||
seen_ports.add(entry.legacy_port)
|
||||
sys_cfg = config.get("SYSTEMS", {}).get(entry.system_name, {})
|
||||
bind_ip = str(sys_cfg.get("_REPORT_BIND_IP") or runtime["listen_ip"] or "")
|
||||
_, legacy_port = listen_obp_fanin(
|
||||
reactor,
|
||||
bind_ip,
|
||||
entry.legacy_port,
|
||||
state.demux,
|
||||
config=config,
|
||||
logger=logger,
|
||||
)
|
||||
state.udp_ports.append(legacy_port)
|
||||
logger.info(
|
||||
"(OBP_PROXY) Legacy port %s:%s -> %s",
|
||||
bind_ip or "*",
|
||||
entry.legacy_port,
|
||||
entry.system_name,
|
||||
)
|
||||
|
||||
|
||||
def start_obp_proxy_service(
|
||||
config: dict[str, Any],
|
||||
protocols: dict[str, Any],
|
||||
*,
|
||||
logger: logging.Logger,
|
||||
) -> ObpProxyServiceState:
|
||||
"""Start OBP fan-in and wire inject-only OPENBRIDGE instances."""
|
||||
if not obp_proxy_enabled(config):
|
||||
raise RuntimeError("(OBP_PROXY) start_obp_proxy_service called but OBP_PROXY is disabled")
|
||||
|
||||
runtime = obp_proxy_settings(config)
|
||||
registry = ObpBridgeRegistry()
|
||||
demux = ObpFanInDemux(registry, debug=runtime["debug"], logger=logger)
|
||||
state = ObpProxyServiceState(
|
||||
demux=demux,
|
||||
registry=registry,
|
||||
listen_port=runtime["listen_port"],
|
||||
listen_ip=runtime["listen_ip"],
|
||||
bind_legacy_ports=runtime["bind_legacy_ports"],
|
||||
_runtime=dict(runtime),
|
||||
)
|
||||
_start_listeners(state, config, protocols, runtime, logger=logger)
|
||||
|
||||
bridge_count = len(state.registry.bridges)
|
||||
logger.info(
|
||||
"(OBP_PROXY) Fan-in on %s:%s (%s bridge(s), BIND_LEGACY_PORTS=%s)",
|
||||
runtime["listen_ip"] or "*",
|
||||
runtime["listen_port"],
|
||||
bridge_count,
|
||||
runtime["bind_legacy_ports"],
|
||||
)
|
||||
if bridge_count == 0:
|
||||
logger.warning("(OBP_PROXY) No enabled OPENBRIDGE systems registered")
|
||||
return state
|
||||
|
||||
|
||||
def apply_obp_proxy_config_reload(
|
||||
state: ObpProxyServiceState,
|
||||
config: dict[str, Any],
|
||||
protocols: dict[str, Any],
|
||||
*,
|
||||
logger: logging.Logger,
|
||||
) -> None:
|
||||
"""Hot-apply OBP_PROXY settings on SIGHUP without rebinding listeners."""
|
||||
incoming = obp_proxy_settings(config)
|
||||
bind_changed = (
|
||||
state.listen_port != incoming["listen_port"]
|
||||
or state.listen_ip != incoming["listen_ip"]
|
||||
or state.bind_legacy_ports != incoming["bind_legacy_ports"]
|
||||
)
|
||||
state._runtime.update(incoming)
|
||||
state.demux.debug = bool(incoming["debug"])
|
||||
if state.udp_ports:
|
||||
primary_proto = getattr(state.udp_ports[0], "protocol", None)
|
||||
transport = getattr(primary_proto, "transport", None) if primary_proto is not None else None
|
||||
if transport is not None:
|
||||
state.registry = build_obp_bridge_registry(
|
||||
config,
|
||||
protocols,
|
||||
bind_legacy_ports=incoming["bind_legacy_ports"],
|
||||
listen_port=incoming["listen_port"],
|
||||
primary_transport=transport,
|
||||
)
|
||||
state.demux._registry = state.registry # noqa: SLF001
|
||||
if bind_changed:
|
||||
logger.warning(
|
||||
"(CONFIG-RELOAD) OBP_PROXY bind change ignored at runtime "
|
||||
"(still listening on %s:%s BIND_LEGACY_PORTS=%s); restart adn-server to apply "
|
||||
"%s:%s BIND_LEGACY_PORTS=%s",
|
||||
state.listen_ip or "*",
|
||||
state.listen_port,
|
||||
state.bind_legacy_ports,
|
||||
incoming["listen_ip"] or "*",
|
||||
incoming["listen_port"],
|
||||
incoming["bind_legacy_ports"],
|
||||
)
|
||||
logger.debug(
|
||||
"(CONFIG-RELOAD) OBP proxy settings applied (%s bridge(s))",
|
||||
len(state.registry.bridges),
|
||||
)
|
||||
|
||||
|
||||
__all__ = [
|
||||
"ObpProxyServiceState",
|
||||
"apply_obp_proxy_config_reload",
|
||||
"build_obp_bridge_registry",
|
||||
"start_obp_proxy_service",
|
||||
]
|
||||
@ -0,0 +1,357 @@
|
||||
# ADN DMR Peer Server - tests infrastructure obp proxy
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""OBP_PROXY configuration and fan-in demux tests."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import pytest
|
||||
from tests.conftest import minimal_valid_config
|
||||
|
||||
from adn_server.application.proxy.deployment import (
|
||||
config_has_enabled_openbridge,
|
||||
is_obp_proxy_managed,
|
||||
normalize_obp_proxy_targets,
|
||||
obp_bridge_legacy_listen_port,
|
||||
obp_proxy_bind_legacy_ports,
|
||||
obp_proxy_enabled,
|
||||
)
|
||||
from adn_server.domain import bytes_4
|
||||
from adn_server.domain.errors import ConfigError
|
||||
from adn_server.infrastructure.config_validator import validate_config
|
||||
from adn_server.infrastructure.hbp_constants import DMRD
|
||||
from adn_server.infrastructure.mesh.obp_v1 import build_bcka, build_dmrd_v1
|
||||
from adn_server.infrastructure.proxy.obp_config import obp_proxy_settings
|
||||
from adn_server.infrastructure.proxy.obp_fanin import (
|
||||
InProcessObpSink,
|
||||
ObpBridgeEntry,
|
||||
ObpBridgeRegistry,
|
||||
ObpFanInDemux,
|
||||
ObpIngressReplyTransport,
|
||||
)
|
||||
from adn_server.infrastructure.proxy.obp_runtime import build_obp_bridge_registry
|
||||
|
||||
_PASS = b"test-passphrase\x00\x00\x00\x00\x00\x00"
|
||||
_NETWORK = bytes_4(73044)
|
||||
_ADDR = ("10.0.0.9", 62044)
|
||||
|
||||
|
||||
class _RecordingTransport:
|
||||
def __init__(self) -> None:
|
||||
self.sent: list[tuple[bytes, tuple[str, int]]] = []
|
||||
self.port = 62032
|
||||
|
||||
def write(self, data: bytes, addr: tuple[str, int]) -> None:
|
||||
self.sent.append((data, addr))
|
||||
|
||||
def getHost(self) -> _RecordingTransport:
|
||||
return self
|
||||
|
||||
|
||||
class _RecordingObp:
|
||||
def __init__(self) -> None:
|
||||
self.packets: list[tuple[bytes, tuple[str, int]]] = []
|
||||
|
||||
def _obp_datagram_received(self, data: bytes, sockaddr: tuple[str, int]) -> None:
|
||||
self.packets.append((data, sockaddr))
|
||||
|
||||
|
||||
def _sample_dmr_voice() -> bytes:
|
||||
return b"".join(
|
||||
[
|
||||
DMRD,
|
||||
bytes([1]),
|
||||
bytes_4(1001)[1:4],
|
||||
bytes_4(52090)[1:4],
|
||||
bytes_4(1),
|
||||
bytes([0x10]),
|
||||
bytes_4(0xAABBCCDD),
|
||||
b"\x00" * 33,
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
def _obp_config(*, bind_legacy: bool = True) -> dict:
|
||||
return {
|
||||
"GLOBAL": {"SERVER_ID": 73010},
|
||||
"DATABASE": {
|
||||
"DB_SERVER": "localhost",
|
||||
"DB_USERNAME": "hbmon",
|
||||
"DB_PASSWORD": "secret",
|
||||
"DB_NAME": "hbmon",
|
||||
"DB_PORT": 3306,
|
||||
},
|
||||
"PROXY": {"LISTEN_PORT": 62031, "TARGET_SYSTEM": "HOTSPOT"},
|
||||
"OBP_PROXY": {
|
||||
"ENABLED": True,
|
||||
"LISTEN_PORT": 62032,
|
||||
"BIND_LEGACY_PORTS": bind_legacy,
|
||||
},
|
||||
"SYSTEMS": {
|
||||
"HOTSPOT": {"MODE": "MASTER", "ENABLED": True, "MAX_PEERS": 1},
|
||||
"OBP-CL": {
|
||||
"MODE": "OPENBRIDGE",
|
||||
"ENABLED": True,
|
||||
"PORT": 62044,
|
||||
"NETWORK_ID": 73044,
|
||||
"PASSPHRASE": "test-passphrase",
|
||||
"TARGET_IP": "127.0.0.1",
|
||||
"TARGET_PORT": 62030,
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def test_obp_proxy_disabled_when_block_absent_and_no_openbridge() -> None:
|
||||
config = minimal_valid_config()
|
||||
assert not config_has_enabled_openbridge(config)
|
||||
assert not obp_proxy_enabled(config)
|
||||
|
||||
|
||||
def test_obp_proxy_defaults_when_block_absent_with_openbridge() -> None:
|
||||
config = _obp_config()
|
||||
del config["OBP_PROXY"]
|
||||
assert obp_proxy_enabled(config)
|
||||
assert obp_proxy_bind_legacy_ports(config)
|
||||
settings = obp_proxy_settings(config)
|
||||
assert settings == {
|
||||
"enabled": True,
|
||||
"listen_port": 62032,
|
||||
"listen_ip": "",
|
||||
"bind_legacy_ports": True,
|
||||
"debug": False,
|
||||
}
|
||||
normalize_obp_proxy_targets(config)
|
||||
assert "PORT" not in config["SYSTEMS"]["OBP-CL"]
|
||||
assert config["SYSTEMS"]["OBP-CL"]["_REPORT_PORT"] == 62044
|
||||
|
||||
|
||||
def test_obp_proxy_explicit_disable_with_openbridge() -> None:
|
||||
config = _obp_config()
|
||||
config["OBP_PROXY"] = {"ENABLED": False}
|
||||
assert not obp_proxy_enabled(config)
|
||||
assert not is_obp_proxy_managed(config, "OBP-CL")
|
||||
|
||||
|
||||
def test_validate_obp_proxy_duplicate_legacy_port() -> None:
|
||||
config = _obp_config()
|
||||
config["SYSTEMS"]["OBP-EU"] = {
|
||||
**config["SYSTEMS"]["OBP-CL"],
|
||||
"NETWORK_ID": 73045,
|
||||
}
|
||||
with pytest.raises(ConfigError) as exc:
|
||||
validate_config(config)
|
||||
assert "duplicate" in str(exc.value).lower()
|
||||
|
||||
|
||||
def test_validate_obp_per_bridge_fanin_migration() -> None:
|
||||
"""7301-style: BIND_LEGACY_PORTS true; one bridge on fan-in, another on legacy."""
|
||||
config = _obp_config(bind_legacy=True)
|
||||
config["SYSTEMS"]["OBP-CL"]["PORT"] = 62032
|
||||
config["SYSTEMS"]["OBP-EU"] = {
|
||||
**config["SYSTEMS"]["OBP-CL"],
|
||||
"PORT": 62999,
|
||||
"NETWORK_ID": 73045,
|
||||
}
|
||||
validate_config(config)
|
||||
|
||||
|
||||
def test_validate_obp_fanin_only_port_equals_listen_with_remote_target() -> None:
|
||||
"""7302-style: no legacy bind; PORT=62032 is metadata; TARGET_PORT is remote fan-in."""
|
||||
config = _obp_config(bind_legacy=False)
|
||||
config["OBP_PROXY"]["BIND_LEGACY_PORTS"] = False
|
||||
config["SYSTEMS"]["OBP-CL"]["PORT"] = 62032
|
||||
config["SYSTEMS"]["OBP-CL"]["TARGET_IP"] = "44.31.61.66"
|
||||
config["SYSTEMS"]["OBP-CL"]["TARGET_PORT"] = 62032
|
||||
validate_config(config)
|
||||
|
||||
|
||||
def test_validate_obp_legacy_local_port_with_remote_fanin_target() -> None:
|
||||
"""7301-style: legacy 62999 locally; mesh to peer fan-in on 62032."""
|
||||
config = _obp_config(bind_legacy=True)
|
||||
config["SYSTEMS"]["OBP-CL"]["PORT"] = 62999
|
||||
config["SYSTEMS"]["OBP-CL"]["TARGET_IP"] = "44.31.61.68"
|
||||
config["SYSTEMS"]["OBP-CL"]["TARGET_PORT"] = 62032
|
||||
validate_config(config)
|
||||
|
||||
|
||||
def test_obp_bridge_legacy_listen_port_per_bridge_migration() -> None:
|
||||
migrated = {"_REPORT_PORT": 62032}
|
||||
legacy = {"_REPORT_PORT": 62999}
|
||||
assert obp_bridge_legacy_listen_port(migrated, listen_port=62032, bind_legacy_ports=True) is None
|
||||
assert obp_bridge_legacy_listen_port(legacy, listen_port=62032, bind_legacy_ports=True) == 62999
|
||||
|
||||
|
||||
def test_obp_proxy_enabled_defaults() -> None:
|
||||
config = _obp_config()
|
||||
assert obp_proxy_enabled(config)
|
||||
assert obp_proxy_bind_legacy_ports(config)
|
||||
|
||||
|
||||
def test_normalize_obp_proxy_targets_omitted_port_uses_fanin() -> None:
|
||||
config = _obp_config(bind_legacy=True)
|
||||
del config["SYSTEMS"]["OBP-CL"]["PORT"]
|
||||
normalize_obp_proxy_targets(config)
|
||||
assert config["SYSTEMS"]["OBP-CL"]["_REPORT_PORT"] == 62032
|
||||
assert obp_bridge_legacy_listen_port(
|
||||
config["SYSTEMS"]["OBP-CL"],
|
||||
listen_port=62032,
|
||||
bind_legacy_ports=True,
|
||||
) is None
|
||||
|
||||
|
||||
def test_is_obp_proxy_managed_only_for_openbridge() -> None:
|
||||
config = _obp_config()
|
||||
assert is_obp_proxy_managed(config, "OBP-CL")
|
||||
assert not is_obp_proxy_managed(config, "HOTSPOT")
|
||||
|
||||
|
||||
def test_normalize_obp_proxy_targets_strips_bind_fields() -> None:
|
||||
config = _obp_config()
|
||||
normalize_obp_proxy_targets(config)
|
||||
obp = config["SYSTEMS"]["OBP-CL"]
|
||||
assert "PORT" not in obp
|
||||
assert obp["_REPORT_PORT"] == 62044
|
||||
|
||||
|
||||
def test_validate_obp_proxy_duplicate_network_id() -> None:
|
||||
config = _obp_config()
|
||||
config["SYSTEMS"]["OBP-EU"] = {
|
||||
**config["SYSTEMS"]["OBP-CL"],
|
||||
"PORT": 62045,
|
||||
}
|
||||
with pytest.raises(ConfigError) as exc:
|
||||
validate_config(config)
|
||||
assert "NETWORK_ID" in str(exc.value)
|
||||
|
||||
|
||||
def test_validate_obp_proxy_migrated_bridge_port_matches_listen() -> None:
|
||||
config = _obp_config()
|
||||
config["OBP_PROXY"]["LISTEN_PORT"] = 62044
|
||||
config["SYSTEMS"]["OBP-CL"]["PORT"] = 62044
|
||||
validate_config(config)
|
||||
|
||||
|
||||
def test_obp_proxy_settings_resolved() -> None:
|
||||
settings = obp_proxy_settings(_obp_config())
|
||||
assert settings["listen_port"] == 62032
|
||||
assert settings["bind_legacy_ports"] is True
|
||||
|
||||
|
||||
def test_obp_proxy_settings_default_listen_port() -> None:
|
||||
config = _obp_config()
|
||||
del config["OBP_PROXY"]["LISTEN_PORT"]
|
||||
settings = obp_proxy_settings(config)
|
||||
assert settings["listen_port"] == 62032
|
||||
assert settings["enabled"] is True
|
||||
|
||||
|
||||
def test_obp_fanin_demux_by_network_id() -> None:
|
||||
receiver = _RecordingObp()
|
||||
transport = _RecordingTransport()
|
||||
reply = ObpIngressReplyTransport(transport)
|
||||
registry = ObpBridgeRegistry()
|
||||
registry.register(
|
||||
ObpBridgeEntry(
|
||||
system_name="OBP-CL",
|
||||
network_id=_NETWORK,
|
||||
passphrase=_PASS,
|
||||
sink=InProcessObpSink(receiver),
|
||||
reply_transport=reply,
|
||||
)
|
||||
)
|
||||
demux = ObpFanInDemux(registry)
|
||||
wire = build_dmrd_v1(_sample_dmr_voice(), _NETWORK, _PASS)
|
||||
demux.deliver(wire, _ADDR, local_port=62032, transport=transport)
|
||||
assert len(receiver.packets) == 1
|
||||
assert receiver.packets[0][0] == wire
|
||||
|
||||
|
||||
def test_obp_fanin_demux_legacy_port_routes_without_network_id() -> None:
|
||||
receiver = _RecordingObp()
|
||||
transport = _RecordingTransport()
|
||||
transport.port = 62044
|
||||
reply = ObpIngressReplyTransport(transport)
|
||||
registry = ObpBridgeRegistry()
|
||||
registry.register(
|
||||
ObpBridgeEntry(
|
||||
system_name="OBP-CL",
|
||||
network_id=_NETWORK,
|
||||
passphrase=_PASS,
|
||||
sink=InProcessObpSink(receiver),
|
||||
reply_transport=reply,
|
||||
legacy_port=62044,
|
||||
)
|
||||
)
|
||||
demux = ObpFanInDemux(registry)
|
||||
wire = build_bcka(_PASS)
|
||||
demux.deliver(wire, _ADDR, local_port=62044, transport=transport)
|
||||
assert len(receiver.packets) == 1
|
||||
|
||||
|
||||
def test_obp_fanin_demux_control_on_listen_port() -> None:
|
||||
receiver = _RecordingObp()
|
||||
transport = _RecordingTransport()
|
||||
reply = ObpIngressReplyTransport(transport)
|
||||
registry = ObpBridgeRegistry()
|
||||
registry.register(
|
||||
ObpBridgeEntry(
|
||||
system_name="OBP-CL",
|
||||
network_id=_NETWORK,
|
||||
passphrase=_PASS,
|
||||
sink=InProcessObpSink(receiver),
|
||||
reply_transport=reply,
|
||||
)
|
||||
)
|
||||
demux = ObpFanInDemux(registry)
|
||||
wire = build_bcka(_PASS)
|
||||
demux.deliver(wire, _ADDR, local_port=62032, transport=transport)
|
||||
assert len(receiver.packets) == 1
|
||||
|
||||
|
||||
def test_build_obp_bridge_registry_starts_inject_protocol() -> None:
|
||||
class _InjectProto:
|
||||
def __init__(self) -> None:
|
||||
self.transport = None
|
||||
self.started = False
|
||||
|
||||
def startProtocol(self) -> None:
|
||||
self.started = True
|
||||
|
||||
proto = _InjectProto()
|
||||
config = {
|
||||
"SYSTEMS": {
|
||||
"OBP-CL": {
|
||||
"MODE": "OPENBRIDGE",
|
||||
"ENABLED": True,
|
||||
"NETWORK_ID": _NETWORK,
|
||||
"PASSPHRASE": _PASS,
|
||||
}
|
||||
}
|
||||
}
|
||||
build_obp_bridge_registry(
|
||||
config,
|
||||
{"OBP-CL": proto},
|
||||
bind_legacy_ports=False,
|
||||
listen_port=62032,
|
||||
primary_transport=_RecordingTransport(),
|
||||
)
|
||||
assert proto.started is True
|
||||
assert proto.transport is not None
|
||||
Loading…
Reference in new issue