test: add HBP/proxy integration coverage and prune redundant suite

Add real-stack REPEAT and proxy fan-in E2E tests so regressions outside
DeterministicScenario are caught. Remove duplicate or fragile tests,
consolidate monitor remap cases, and document harness vs integration policy.
pull/1/head
Rodrigo Pérez 4 months ago
parent 614c882bcb
commit 46024c4da4

@ -33,6 +33,8 @@ pythonpath = ["."]
testpaths = ["tests"]
markers = [
"behavior: integration-style behavior tests (P0/P1)",
"integration: real HBP/proxy stack (not DeterministicScenario inject)",
"mqtt: requires MQTT broker or heavy mqtt mocks",
"smoke: quick routing smoke tests",
]

@ -357,7 +357,6 @@ class BridgeLcTaMixin:
settings = talker_alias_settings(self._config, source_system)
if not settings["enabled"]:
return
target_mode = self._config.get("SYSTEMS", {}).get(target_system, {}).get("MODE")
if settings["mode"] == "both" and not self._ta_capable_source(source_system):
force_inject = True
st.pop("TX_TA_EMB", None)

@ -8,7 +8,8 @@ python3 -m pip install --no-user -e ".[dev]"
python3 -m pytest tests/<path>/test_<name>.py -q # single file
python3 -m pytest tests/<path>/test_<name>.py::test_foo -q # single test
python3 -m pytest tests/bridge/ -q # whole domain
python3 -m pytest tests/ -q # full suite (177 + 1 skip)
python3 -m pytest tests/ -q # full suite
python3 -m pytest tests/ -q -m "not mqtt" # skip MQTT-heavy tests
```
Use the project interpreter, e.g. `/opt/.pyenv/versions/3.11.8/bin/python3`.
@ -21,18 +22,30 @@ Use the project interpreter, e.g. `/opt/.pyenv/versions/3.11.8/bin/python3`.
| `hbp/` | HBP ingress, loop control, rate limit, timeout/collision, master maintenance |
| `obp/` | OpenBridge loop, rate limit, unit-data loop |
| `voice/` | Announcements, TTS schedule, broadcast queue, disconnected voice, in-band signalling |
| `talker_alias/` | Encode/decode, passthrough, MMDVM wire, bridge inject |
| `talker_alias/` | Encode/decode, passthrough, MMDVM wire, bridge inject (DeterministicScenario) |
| `parrot/` | Recording timers, playback loop, seq preservation, ingress path |
| `replay/` | JSONL session replay (V2-P0-007) |
| `schemas/` | Report v2 JSON Schema validation (`jsonschema` dev dep) |
| `application/test_report_payloads.py` | Report payload builders + CSV parse |
| `infrastructure/test_report_server_wire.py` | Report server wire opcodes |
| `smoke/` | Quick routing smoke + packet builder |
| `infrastructure/` | Logging reload, bridge router index |
| `application/` | RuntimeContext holder / config proxy |
| `scripts/` | Config conversion helpers |
| `application/` | Report payloads, monitor topology, proxy use cases, runtime context |
| `infrastructure/` | Logging reload, bridge router, **HBP REPEAT + proxy fan-in integration**, MQTT |
| `smoke/` | Quick routing smoke |
| `support/` | Shared stacks (`hbp_repeat_stack`, monitor sim) — not run as tests |
| `harness/` | Shared fakes (`DeterministicScenario`, assertions) — not run as tests |
## Integration vs harness
Most bridge/voice tests inject packets via **`DeterministicScenario`** (`bridge.dmrd_received` on fakes). That is fast but **skips** `udp_hbp` REPEAT rewrite and proxy UDP fan-in.
For regressions on those paths, use:
| File | Topic |
|------|-------|
| `infrastructure/test_hbp_repeat_talker_alias.py` | Real `HBPProtocol` REPEAT + embedded TA |
| `infrastructure/test_proxy_repeat_e2e.py` | Proxy fan-in → REPEAT downlink |
| `infrastructure/test_proxy_reload.py` | Hot reload keeps UDP listener |
Mark new stack tests with `@pytest.mark.integration`.
## Files by domain
### bridge/
@ -74,7 +87,6 @@ Use the project interpreter, e.g. `/opt/.pyenv/versions/3.11.8/bin/python3`.
| `test_announcement_anticollision.py` | 3 | Busy slot skip / abort |
| `test_broadcast_queue.py` | 2 | Same-TG broadcast queue |
| `test_disconnected_voice.py` | 3 | Not-linked / reflector prompts |
| `test_embed_ta_forward.py` | 2 | Embed TA state on bridge |
| `test_in_band_signalling.py` | 5 | Reflector / single-mode VTERM |
| `test_play_file_on_request.py` | 3 | On-demand file playback |
| `test_scheduled_announcement.py` | 4 | File announcements (AMBE) |
@ -85,7 +97,7 @@ Use the project interpreter, e.g. `/opt/.pyenv/versions/3.11.8/bin/python3`.
| File | Tests | Topic |
|------|-------|-------|
| `test_bridge_inject.py` | 5 | TA inject on bridge VHEAD |
| `test_bridge_inject.py` | 3 | TA inject on bridge VHEAD (harness) |
| `test_embed_ta.py` | 4 | Embedded LC modes |
| `test_encode_decode.py` | 8 | Domain encode/decode |
| `test_format.py` | 2 | Format from subscriber profile |
@ -102,16 +114,25 @@ Use the project interpreter, e.g. `/opt/.pyenv/versions/3.11.8/bin/python3`.
| `test_recording_timers.py` | 2 | Idle timeout, VTERM commit |
| `test_rekey_playback.py` | 4 | Seq preservation past 255 / 30s |
### smoke/ · infrastructure/ · scripts/
### infrastructure (integration highlights)
| File | Tests | Topic |
|------|-------|-------|
| `smoke/test_bridge_routing.py` | 1 | Static TG forward smoke |
| `smoke/test_packet_builder.py` | 1 | PacketSpec builder |
| `infrastructure/test_logging_reload.py` | 2 | Log level reload |
| `infrastructure/test_bridge_router_index.py` | 5 | BRIDGES O(1) index vs legacy scan |
| `application/test_runtime_context.py` | 5 | RuntimeContext holder, SIGHUP swap prep |
| `scripts/test_freedmr_cfg_to_yaml.py` | 5 | Legacy cfg → YAML |
| File | Topic |
|------|-------|
| `test_hbp_repeat_talker_alias.py` | REPEAT embed TA through real HBP |
| `test_proxy_repeat_e2e.py` | Proxy → REPEAT E2E |
| `test_proxy_reload.py` | Proxy hot reload |
| `test_udp_fanin.py` | UDP fan-in routing |
| `test_report_server_wire.py` | Report server wire opcodes |
| `test_logging_reload.py` | Log level reload |
| `test_bridge_router_index.py` | BRIDGES O(1) index vs legacy scan |
### smoke/ · application/
| File | Topic |
|------|-------|
| `smoke/test_bridge_routing.py` | Static TG forward smoke |
| `application/test_monitor_topology.py` | Inject-only proxy monitor remap |
| `application/test_runtime_context.py` | RuntimeContext holder, SIGHUP swap prep |
## Examples (copy-paste)
@ -122,6 +143,9 @@ python3 -m pytest tests/bridge/test_unit_data_routing.py -q
# After parrot seq fix
python3 -m pytest tests/parrot/test_rekey_playback.py -q
# Talker Alias REPEAT (real stack)
python3 -m pytest tests/infrastructure/test_hbp_repeat_talker_alias.py -q
# HBP loop + rate (common RF regressions)
python3 -m pytest tests/hbp/test_hbp_loop_control.py tests/hbp/test_hbp_rate_limit.py -q
@ -132,5 +156,6 @@ python3 -m pytest tests/bridge/test_startup_bridges.py::test_startup_bridge_rout
## Policy
- **New tests:** add a new file (or extend the smallest existing file for the same topic). Avoid large multi-topic modules.
- **Harness:** shared code lives in `harness/` only.
- **Harness:** shared code lives in `harness/` and `support/` only.
- **Stack regressions:** prefer `infrastructure/test_*_e2e.py` with real adapters over duplicating in `DeterministicScenario`.
- Full audit: `docs-priv/en/test-audit.md` (maintainer checkout).

@ -2,6 +2,8 @@
from __future__ import annotations
import pytest
from adn_server.application.report.monitor_topology import (
expand_inject_proxy_systems,
remap_inject_proxy_voice_event,
@ -22,25 +24,31 @@ def _peer(*, connected: bool = True) -> dict:
}
def test_expand_inject_proxy_fans_peers_into_system_n() -> None:
peer_a = bytes_4(730039101)
peer_b = bytes_4(7301896)
config = {
def _proxy_config(
peers: dict[bytes, dict],
*,
max_peers: int = 102,
base_port: int = 56400,
) -> dict:
return {
"PROXY": {"TARGET_SYSTEM": "SYSTEM"},
"SYSTEMS": {
"SYSTEM": {
"MODE": "MASTER",
"ENABLED": True,
"MAX_PEERS": 102,
"_REPORT_BASE_PORT": 56400,
"PEERS": {
peer_a: _peer(),
peer_b: _peer(),
},
},
"ECHO": {"MODE": "MASTER", "ENABLED": True, "PEERS": {}},
"MAX_PEERS": max_peers,
"_REPORT_BASE_PORT": base_port,
"PEERS": peers,
}
},
}
def test_expand_inject_proxy_fans_peers_into_system_n() -> None:
peer_a = bytes_4(730039101)
peer_b = bytes_4(7301896)
config = _proxy_config({peer_a: _peer(), peer_b: _peer()})
config["SYSTEMS"]["ECHO"] = {"MODE": "MASTER", "ENABLED": True, "PEERS": {}}
expanded = expand_inject_proxy_systems(
config,
config["SYSTEMS"],
@ -56,18 +64,7 @@ def test_expand_inject_proxy_fans_peers_into_system_n() -> None:
def test_build_topology_after_expand_matches_monitor_shape() -> None:
peer = bytes_4(730039101)
config = {
"PROXY": {"TARGET_SYSTEM": "SYSTEM"},
"SYSTEMS": {
"SYSTEM": {
"MODE": "MASTER",
"ENABLED": True,
"MAX_PEERS": 102,
"_REPORT_BASE_PORT": 56400,
"PEERS": {peer: _peer()},
}
},
}
config = _proxy_config({peer: _peer()})
expanded = expand_inject_proxy_systems(config, config["SYSTEMS"], {peer: 2})
doc = build_topology(expanded, seq=1)
names = {system["name"] for system in doc["systems"]}
@ -85,18 +82,7 @@ def test_build_topology_after_expand_matches_monitor_shape() -> None:
def test_expand_inject_proxy_emits_all_virtual_masters() -> None:
peer = bytes_4(730039101)
config = {
"PROXY": {"TARGET_SYSTEM": "SYSTEM"},
"SYSTEMS": {
"SYSTEM": {
"MODE": "MASTER",
"ENABLED": True,
"MAX_PEERS": 4,
"_REPORT_BASE_PORT": 56400,
"PEERS": {peer: _peer()},
}
},
}
config = _proxy_config({peer: _peer()}, max_peers=4)
expanded = expand_inject_proxy_systems(config, config["SYSTEMS"], {peer: 2})
for slot in range(4):
assert f"SYSTEM-{slot}" in expanded
@ -106,138 +92,74 @@ def test_expand_inject_proxy_emits_all_virtual_masters() -> None:
assert expanded["SYSTEM-1"]["PEERS"] == {}
def test_remap_voice_event_system_to_virtual_master() -> None:
peer = bytes_4(730039101)
config = {
"PROXY": {"TARGET_SYSTEM": "SYSTEM"},
"SYSTEMS": {
"SYSTEM": {
"MODE": "MASTER",
"ENABLED": True,
"MAX_PEERS": 102,
"PEERS": {peer: _peer()},
}
},
}
raw = "GROUP VOICE,START,RX,SYSTEM,3262598598,730039101,730039101,2,730444"
remapped = remap_inject_proxy_voice_event(
raw, config, config["SYSTEMS"], {peer: 4}
)
assert remapped.startswith("GROUP VOICE,START,RX,SYSTEM-4,")
def test_remap_voice_event_tx_keeps_echo_peer_id_for_hotspot_rx_display() -> None:
"""TX echo→hotspot: field 5 stays 9990 so hotspot chip is TX/green while receiving."""
peer = bytes_4(730039101)
config = {
"PROXY": {"TARGET_SYSTEM": "SYSTEM"},
"SYSTEMS": {
"SYSTEM": {
"MODE": "MASTER",
"ENABLED": True,
"MAX_PEERS": 102,
"PEERS": {peer: _peer()},
}
},
}
raw = "GROUP VOICE,START,TX,SYSTEM,4100887026,9990,730039101,2,730444"
remapped = remap_inject_proxy_voice_event(
raw, config, config["SYSTEMS"], {peer: 4}
)
parts = remapped.split(",")
assert parts[3] == "SYSTEM-4"
assert parts[5] == "9990"
def test_remap_voice_event_rx_normalizes_field5_to_hotspot_radio_id() -> None:
peer = bytes_4(730039101)
config = {
"PROXY": {"TARGET_SYSTEM": "SYSTEM"},
"SYSTEMS": {
"SYSTEM": {
"MODE": "MASTER",
"ENABLED": True,
"MAX_PEERS": 102,
"PEERS": {peer: _peer()},
}
},
}
raw = "GROUP VOICE,START,RX,SYSTEM,4100887026,73003,7300392,2,9990"
remapped = remap_inject_proxy_voice_event(
raw, config, config["SYSTEMS"], {peer: 4}
)
parts = remapped.split(",")
assert parts[3] == "SYSTEM-4"
assert parts[5] == "730039101"
def test_remap_voice_event_matches_user_prefix_when_single_hotspot_online() -> None:
"""User 7300391 with only HS1 online: legacy short subscriber id can resolve."""
peer = bytes_4(730039101)
config = {
"PROXY": {"TARGET_SYSTEM": "SYSTEM"},
"SYSTEMS": {
"SYSTEM": {
"MODE": "MASTER",
"ENABLED": True,
"MAX_PEERS": 102,
"PEERS": {peer: _peer()},
}
},
}
raw = "GROUP VOICE,START,TX,SYSTEM,2693411696,9990,7300392,2,9990"
remapped = remap_inject_proxy_voice_event(
raw, config, config["SYSTEMS"], {peer: 4}
)
parts = remapped.split(",")
assert parts[3] == "SYSTEM-4"
assert parts[5] == "9990"
def test_remap_voice_event_skips_ambiguous_user_with_multiple_hotspots() -> None:
"""User 7300391 + HS1/HS2 online: do not pick the wrong radio from rf_src alone."""
hs1 = bytes_4(730039101)
hs2 = bytes_4(730039102)
config = {
"PROXY": {"TARGET_SYSTEM": "SYSTEM"},
"SYSTEMS": {
"SYSTEM": {
"MODE": "MASTER",
"ENABLED": True,
"MAX_PEERS": 102,
"PEERS": {hs1: _peer(), hs2: _peer()},
}
},
}
raw = "GROUP VOICE,START,TX,SYSTEM,2693411696,9990,7300391,2,9990"
remapped = remap_inject_proxy_voice_event(
raw, config, config["SYSTEMS"], {hs1: 4, hs2: 5}
)
assert remapped == raw
def test_remap_voice_event_resolves_full_radio_id_with_sibling_hotspots_online() -> None:
"""Full radio id in rf_src still maps correctly when sibling HS are connected."""
hs1 = bytes_4(730039101)
hs2 = bytes_4(730039102)
config = {
"PROXY": {"TARGET_SYSTEM": "SYSTEM"},
"SYSTEMS": {
"SYSTEM": {
"MODE": "MASTER",
"ENABLED": True,
"MAX_PEERS": 102,
"PEERS": {hs1: _peer(), hs2: _peer()},
}
},
}
raw = "GROUP VOICE,START,TX,SYSTEM,4100887026,9990,730039102,2,730444"
@pytest.mark.parametrize(
"raw,peer_specs,slot_map,expect",
[
pytest.param(
"GROUP VOICE,START,RX,SYSTEM,3262598598,730039101,730039101,2,730444",
[(730039101,)],
{730039101: 4},
{"startswith": "GROUP VOICE,START,RX,SYSTEM-4,"},
id="rx_to_virtual_master",
),
pytest.param(
"GROUP VOICE,START,TX,SYSTEM,4100887026,9990,730039101,2,730444",
[(730039101,)],
{730039101: 4},
{"parts": {3: "SYSTEM-4", 5: "9990"}},
id="tx_echo_keeps_9990_for_hotspot_rx",
),
pytest.param(
"GROUP VOICE,START,RX,SYSTEM,4100887026,73003,7300392,2,9990",
[(730039101,)],
{730039101: 4},
{"parts": {3: "SYSTEM-4", 5: "730039101"}},
id="rx_normalizes_field5_to_hotspot_radio_id",
),
pytest.param(
"GROUP VOICE,START,TX,SYSTEM,2693411696,9990,7300392,2,9990",
[(730039101,)],
{730039101: 4},
{"parts": {3: "SYSTEM-4", 5: "9990"}},
id="single_hotspot_user_prefix",
),
pytest.param(
"GROUP VOICE,START,TX,SYSTEM,2693411696,9990,7300391,2,9990",
[(730039101,), (730039102,)],
{730039101: 4, 730039102: 5},
{"unchanged": True},
id="ambiguous_user_multiple_hotspots",
),
pytest.param(
"GROUP VOICE,START,TX,SYSTEM,4100887026,9990,730039102,2,730444",
[(730039101,), (730039102,)],
{730039101: 4, 730039102: 5},
{"parts": {3: "SYSTEM-5", 5: "9990"}},
id="full_radio_id_with_sibling_hotspots",
),
],
)
def test_remap_voice_event_inject_proxy(
raw: str,
peer_specs: list[tuple[int, ...]],
slot_map: dict[int, int],
expect: dict,
) -> None:
peers = {bytes_4(rid): _peer() for spec in peer_specs for rid in spec}
peer_slots = {bytes_4(rid): slot_map[rid] for spec in peer_specs for rid in spec}
config = _proxy_config(peers)
remapped = remap_inject_proxy_voice_event(
raw, config, config["SYSTEMS"], {hs1: 4, hs2: 5}
raw, config, config["SYSTEMS"], peer_slots
)
if expect.get("unchanged"):
assert remapped == raw
return
if prefix := expect.get("startswith"):
assert remapped.startswith(prefix)
return
parts = remapped.split(",")
assert parts[3] == "SYSTEM-5"
assert parts[5] == "9990"
for idx, value in expect.get("parts", {}).items():
assert parts[idx] == value
def test_remap_voice_event_passes_through_non_proxy_systems() -> None:

@ -1,19 +0,0 @@
"""Regression: configs that must keep loading after proxy integration."""
from __future__ import annotations
from adn_server.infrastructure import YamlConfigLoader
def test_adn_parrot_yaml_loads_without_proxy_block() -> None:
loader = YamlConfigLoader("/opt/new-adn-server")
config = loader.load("/opt/new-adn-server/adn-parrot.yaml")
assert "PROXY" not in config or config.get("PROXY") == {}
assert config["SYSTEMS"]["PARROT"]["MODE"] == "PEER"
def test_adn_server_yaml_loads_with_proxy() -> None:
loader = YamlConfigLoader("/opt/new-adn-server")
config = loader.load("/opt/new-adn-server/adn-server.yaml")
assert config["PROXY"]["TARGET_SYSTEM"] == "SYSTEM"
assert config["PROXY"]["LISTEN_PORT"] == 62031

@ -0,0 +1,124 @@
"""REPEAT path through real HBPProtocol — catches missing embedded Talker Alias."""
from __future__ import annotations
import pytest
from bitarray import bitarray
from adn_server.domain import bytes_4
from adn_server.infrastructure.hbp_constants import DMRA
from tests.harness.deterministic import DeterministicScenario, PacketSpec
from tests.support.hbp_repeat_stack import build_hbp_repeat_stack
pytestmark = pytest.mark.integration
_PEER_TX = bytes_4(730039210)
_PEER_RX = bytes_4(730039101)
_ADDR_TX = ("10.0.0.1", 62001)
_ADDR_RX = ("10.0.0.2", 62002)
_EMB_SLICE = slice(116, 148)
def _embed_bits(dmrpkt: bytes) -> bitarray:
bits = bitarray(endian="big")
bits.frombytes(dmrpkt)
return bits[_EMB_SLICE]
def _base_spec() -> PacketSpec:
return PacketSpec(
peer_id=730039210,
rf_src=7300392,
dst_id=7304,
slot=2,
stream_id=0xA1B2C3D4,
payload=b"\x00" * 33,
)
def test_repeat_downlink_embeds_group_lc_when_talker_alias_enabled() -> None:
"""Regression: raw REPEAT copy left embed LC all-zero; MMDVM never saw TA."""
stack = build_hbp_repeat_stack(talker_alias=True)
stack.register_peer(_PEER_TX, _ADDR_TX)
stack.register_peer(_PEER_RX, _ADDR_RX)
base = _base_spec()
uplink_burst = DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
stack.inject_spec(uplink_burst, _ADDR_TX)
downlink = stack.transport.for_addr(_ADDR_RX)
assert len(downlink) == 1
assert downlink[0][11:15] == _PEER_RX
uplink_emb = _embed_bits(uplink_burst.payload)
downlink_emb = _embed_bits(downlink[0][20:53])
assert downlink_emb != uplink_emb
slot_st = stack.hbp.STATUS[2]
assert slot_st.get("TX_TA_EMB") is not None
assert downlink_emb == slot_st["REP_EMB_LC"][1]
def test_repeat_sends_dmra_and_embed_on_vhead() -> None:
"""MMDVM may ignore DMRA UDP; embed in DMRD must still be prepared on VHEAD."""
stack = build_hbp_repeat_stack(talker_alias=True)
stack.register_peer(_PEER_TX, _ADDR_TX)
stack.register_peer(_PEER_RX, _ADDR_RX)
base = _base_spec()
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
assert stack.dmra_capture, "expected inject DMRA on VHEAD"
packets, exclude = stack.dmra_capture[0]
assert exclude == _PEER_TX
assert all(p[:4] == DMRA for p in packets)
dmra_to_rx = [p for p, _ in stack.transport.sent if p[:4] == DMRA and _ == _ADDR_RX]
assert dmra_to_rx, "DMRA should reach listening peer (legacy path)"
assert stack.hbp.STATUS[2].get("REP_STREAM_ID") == bytes_4(base.stream_id)
def test_repeat_leaves_burst_payload_unchanged_when_talker_alias_disabled() -> None:
stack = build_hbp_repeat_stack(talker_alias=False)
stack.register_peer(_PEER_TX, _ADDR_TX)
stack.register_peer(_PEER_RX, _ADDR_RX)
base = _base_spec()
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
uplink = DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1)
stack.inject_spec(uplink, _ADDR_TX)
downlink = stack.transport.for_addr(_ADDR_RX)
assert len(downlink) == 1
assert downlink[0][20:53] == uplink.payload
assert stack.hbp.STATUS[2].get("TX_TA_EMB") is None
def test_repeat_overlays_ta_on_second_superframe_cycle() -> None:
"""After bursts B–E, the next B should carry TA embed (alternate superframes)."""
stack = build_hbp_repeat_stack(talker_alias=True)
stack.register_peer(_PEER_TX, _ADDR_TX)
stack.register_peer(_PEER_RX, _ADDR_RX)
base = _base_spec()
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
for seq, dtype in enumerate((1, 2, 3, 4), start=1):
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=seq, dtype_vseq=dtype),
_ADDR_TX,
)
stack.transport.clear()
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=5, dtype_vseq=1),
_ADDR_TX,
)
downlink = stack.transport.for_addr(_ADDR_RX)
assert downlink
slot_st = stack.hbp.STATUS[2]
ta_emb = slot_st.get("TX_TA_EMB")
assert ta_emb is not None
bits = bitarray(endian="big")
bits.frombytes(downlink[0][20:53])
phase = slot_st.get("TX_TA_PHASE", 0)
assert bits[_EMB_SLICE] == ta_emb[phase][1]

@ -3,8 +3,10 @@
from __future__ import annotations
import logging
from unittest.mock import MagicMock
from adn_server.application.proxy import ProxyUseCases
from adn_server.infrastructure.config_reload import merge_top_level_config
from adn_server.domain.proxy import ClientEndpoint, ClientSlot
from adn_server.infrastructure.proxy.ip_blacklist import InMemoryProxyIpBlacklist
from adn_server.infrastructure.proxy.rpto_queue import InMemoryPendingRptoQueue
@ -75,6 +77,32 @@ def test_apply_proxy_config_reload_keeps_sessions_and_updates_timeout() -> None:
assert state.use_cases._max_peers == 20 # noqa: SLF001
def test_hot_reload_never_closes_udp_listener() -> None:
"""Regression: SIGHUP must not stopListening() on PROXY (EADDRINUSE / dropped sessions)."""
state = _minimal_proxy_state()
stop_mock = MagicMock(return_value=None)
state.udp_port = MagicMock()
state.udp_port.stopListening = stop_mock
apply_proxy_config_reload(
state,
{
"PROXY": {"LISTEN_PORT": 62031, "TARGET_SYSTEM": "SYSTEM", "TIMEOUT": 60},
"SYSTEMS": {"SYSTEM": {"MAX_PEERS": 20}},
},
logger=logging.getLogger("test.proxy.reload"),
)
stop_mock.assert_not_called()
assert len(state.use_cases.list_slots()) == 1
def test_merge_top_level_config_updates_proxy_section() -> None:
live = {"GLOBAL": {}, "PROXY": {"TIMEOUT": 30, "LISTEN_PORT": 62031}}
merge_top_level_config(live, {"PROXY": {"TIMEOUT": 45, "LISTEN_PORT": 62031}})
assert live["PROXY"]["TIMEOUT"] == 45
def test_proxy_use_cases_apply_runtime_settings() -> None:
uc = ProxyUseCases(
InMemoryProxySlotStore(),

@ -0,0 +1,88 @@
"""Proxy fan-in → MASTER inject → REPEAT to second hotspot (DroidStar path)."""
from __future__ import annotations
import pytest
from bitarray import bitarray
from adn_server.application.proxy import ProxyUseCases
from adn_server.domain import bytes_4
from adn_server.infrastructure.config_normalizer import ensure_system_runtime_config
from adn_server.infrastructure.hbp_constants import DMRD
from adn_server.infrastructure.proxy import (
InMemoryPendingRptoQueue,
InMemoryProxySlotStore,
InProcessHbpSink,
ProxyFanInProtocol,
ProxyReplyTransport,
)
from tests.harness.deterministic import DeterministicScenario, PacketSpec
from tests.support.hbp_repeat_stack import HbpRepeatStack, build_hbp_repeat_stack
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)
_EMB_SLICE = slice(116, 148)
def _proxy_fanin_stack() -> tuple[ProxyFanInProtocol, HbpRepeatStack]:
stack = build_hbp_repeat_stack(talker_alias=True, system_name="MASTER-A")
stack.register_peer(_PEER_TX, _ADDR_TX)
stack.register_peer(_PEER_RX, _ADDR_RX)
proxy = ProxyUseCases(
InMemoryProxySlotStore(),
InMemoryPendingRptoQueue(),
max_peers=8,
)
sink = InProcessHbpSink(stack.hbp)
fanin = ProxyFanInProtocol(proxy, sink)
fanin.transport = stack.transport # type: ignore[assignment]
stack.hbp.transport = ProxyReplyTransport(stack.transport)
return fanin, stack
def test_proxy_inject_repeat_reaches_listener_with_embed_lc() -> None:
"""TX via proxy (DroidStar) must REPEAT to other peer with TA embed, not raw copy."""
fanin, stack = _proxy_fanin_stack()
base = PacketSpec(
peer_id=730039210,
rf_src=7300392,
dst_id=7304,
slot=2,
stream_id=0x11223344,
payload=b"\x00" * 33,
)
vhead = DeterministicScenario.voice_head_spec(base)
burst = DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1)
fanin.datagramReceived(vhead.data(), _ADDR_TX)
stack.transport.clear()
fanin.datagramReceived(burst.data(), _ADDR_TX)
listener_pkts = [p for p in stack.transport.for_addr(_ADDR_RX) if p[:4] == DMRD]
assert len(listener_pkts) == 1
bits = bitarray(endian="big")
bits.frombytes(listener_pkts[0][20:53])
slot_st = stack.hbp.STATUS[2]
assert slot_st.get("TX_TA_EMB") is not None
assert bits[_EMB_SLICE] == slot_st["REP_EMB_LC"][1]
assert listener_pkts[0][20:53] != burst.payload
def test_proxy_attach_binds_sockaddr_used_for_master_ingress() -> None:
"""Peer must be registered at proxy client addr so DMRD is accepted and REPEAT fans out."""
fanin, stack = _proxy_fanin_stack()
ensure_system_runtime_config(stack.config)
packet = DeterministicScenario.voice_head_spec(
PacketSpec(peer_id=730039210, stream_id=0x55667788)
).data()
fanin.datagramReceived(packet, _ADDR_TX)
assert stack.hbp._peers[_PEER_TX]["SOCKADDR"] == _ADDR_TX
rx_traffic = stack.transport.for_addr(_ADDR_RX)
assert rx_traffic, "REPEAT must reach the other logged-in hotspot"
assert any(p[:4] == DMRD for p in rx_traffic)

@ -1,91 +0,0 @@
"""Live UDP smoke test for integrated proxy (isolated port; no production restart)."""
from __future__ import annotations
import socket
import threading
import pytest
from twisted.internet import reactor
from adn_server.application.proxy import ProxyUseCases
from adn_server.domain.value_objects import bytes_4
from adn_server.infrastructure.config_normalizer import ensure_system_runtime_config
from adn_server.infrastructure.hbp_constants import RPTACK, RPTL
from adn_server.infrastructure.proxy import (
InMemoryPendingRptoQueue,
InMemoryProxySlotStore,
InProcessHbpSink,
ProxyFanInProtocol,
ProxyReplyTransport,
)
from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol
_SMOKE_PORT = 62032
_PEER = bytes_4(1234567)
class _AclRouter:
def acl_check(self, peer_id: bytes, acl: object) -> bool:
return True
def _build_fanin() -> ProxyFanInProtocol:
config = {
"GLOBAL": {"PING_TIME": 10, "MAX_MISSED": 3, "USE_ACL": False},
"SYSTEMS": {
"HOTSPOT": {
"MODE": "MASTER",
"ENABLED": True,
"MAX_PEERS": 8,
"OPTIONS": "TS2=9990;",
}
},
}
ensure_system_runtime_config(config)
hbp = HBPProtocol("HOTSPOT", config)
hbp._router = _AclRouter() # type: ignore[assignment]
proxy = ProxyUseCases(
InMemoryProxySlotStore(),
InMemoryPendingRptoQueue(),
max_peers=8,
)
sink = InProcessHbpSink(hbp)
fanin = ProxyFanInProtocol(proxy, sink)
return fanin, hbp, sink
@pytest.mark.smoke
def test_live_udp_rptl_rptack_on_isolated_port() -> None:
"""RPTL in → inject → RPTACK out on 127.0.0.1:62032 (does not use production 62031)."""
fanin, hbp, _ = _build_fanin()
result: dict[str, bytes | None] = {"reply": None, "error": None}
def _run_client() -> None:
try:
client = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
client.settimeout(2.0)
client.bind(("127.0.0.1", 0))
client.sendto(RPTL + _PEER, ("127.0.0.1", _SMOKE_PORT))
data, _ = client.recvfrom(4096)
result["reply"] = data
client.close()
except OSError as exc:
result["error"] = str(exc).encode()
finally:
reactor.callFromThread(reactor.stop)
listener = reactor.listenUDP(_SMOKE_PORT, fanin, interface="127.0.0.1")
hbp.transport = ProxyReplyTransport(fanin.transport)
reactor.callWhenRunning(
lambda: threading.Thread(target=_run_client, daemon=True).start()
)
reactor.callLater(5.0, reactor.stop)
reactor.run()
listener.stopListening()
if result["error"]:
pytest.fail(result["error"].decode())
reply = result["reply"]
assert reply is not None, "no UDP reply (timeout?)"
assert reply.startswith(RPTACK), f"expected RPTACK, got {reply[:8]!r}"

@ -1,10 +0,0 @@
"""Report wire encoder factory."""
from __future__ import annotations
from adn_server.infrastructure.twisted_adapters.report import ReportWire, create_report_wire
def test_create_report_wire_returns_encoder() -> None:
wire = create_report_wire({"REPORTS": {}})
assert isinstance(wire, ReportWire)

@ -1,174 +0,0 @@
"""Tests for scripts/freedmr_cfg_to_yaml.py."""
from __future__ import annotations
import sys
import textwrap
from pathlib import Path
import yaml
_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(_ROOT / "scripts"))
from freedmr_cfg_to_yaml import dump_yaml, parse_freedmr_cfg # noqa: E402
def _write_cfg(tmp_path: Path, body: str) -> Path:
path = tmp_path / "test.cfg"
path.write_text(textwrap.dedent(body).strip() + "\n", encoding="utf-8")
return path
def test_converts_global_echo_and_obp(tmp_path: Path) -> None:
cfg = _write_cfg(
tmp_path,
"""
[GLOBAL]
SERVER_ID: 73010
USE_ACL: True
TGID_TS2_ACL: PERMIT:ALL
[REPORTS]
REPORT: True
REPORT_CLIENTS: 127.0.0.1
[LOGGER]
LOG_FILE: /var/log/FreeDMR/FreeDMR.log
LOG_NAME: FreeDMR
[ALIASES]
TRY_DOWNLOAD: True
STALE_DAYS: 1
[ECHO]
MODE: MASTER
ENABLED: True
PORT: 54917
TS2_STATIC: 9990
TGID_TS2_ACL: PERMIT:9990
GENERATOR: 0
[OBP-CR]
MODE: OPENBRIDGE
ENABLED: False
PORT: 62052
NETWORK_ID: 71210
PASSPHRASE: passw0rd
TARGET_IP: freedmrcr.net
TARGET_PORT: 62026
TGID_ACL: DENY :0-82,9990-9999
TGID_TS1_ACL: DENY :0-89
PROTO_VER: 5
""",
)
out = parse_freedmr_cfg(cfg)
assert out["GLOBAL"]["SERVER_ID"] == 73010
assert out["GLOBAL"]["TALKER_ALIAS"] is False
assert out["REPORTS"]["REPORT_CLIENTS"] == "127.0.0.1"
assert out["LOGGER"]["LOG_FILE"] == "/var/log/adn-server/adn-server.log"
assert out["LOGGER"]["LOG_NAME"] == "adn-server"
assert out["ALIASES"]["KEYS_FILE"] == "keys.json"
echo = out["SYSTEMS"]["ECHO"]
assert echo["MODE"] == "MASTER"
assert echo["TS2_STATIC"] == "9990"
assert echo["TGID_TS2_ACL"] == "PERMIT:9990"
obp = out["SYSTEMS"]["OBP-CR"]
assert obp["ENABLED"] is False
assert obp["TGID_ACL"] == "DENY:0-82,9990-9999"
assert obp["TGID_TS1_ACL"] == "DENY:0-89"
assert obp["PROTO_VER"] == 5
def test_includes_disabled_systems(tmp_path: Path) -> None:
cfg = _write_cfg(
tmp_path,
"""
[OBP-OFF]
MODE: OPENBRIDGE
ENABLED: False
PORT: 62099
NETWORK_ID: 1
PASSPHRASE: x
TARGET_IP: 1.2.3.4
TARGET_PORT: 62099
""",
)
out = parse_freedmr_cfg(cfg)
assert "OBP-OFF" in out["SYSTEMS"]
assert out["SYSTEMS"]["OBP-OFF"]["ENABLED"] is False
def test_preserves_section_and_key_order(tmp_path: Path) -> None:
cfg = _write_cfg(
tmp_path,
"""
[GLOBAL]
PATH: ./
PING_TIME: 10
SERVER_ID: 73010
[REPORTS]
REPORT: True
[ALLSTAR]
ENABLED: False
[ECHO]
MODE: MASTER
ENABLED: True
PORT: 54917
TS2_STATIC: 9990
[OBP-CR]
MODE: OPENBRIDGE
ENABLED: False
PORT: 62052
NETWORK_ID: 71210
PASSPHRASE: passw0rd
TARGET_IP: freedmrcr.net
TARGET_PORT: 62026
PROTO_VER: 5
""",
)
out = parse_freedmr_cfg(cfg)
assert list(out.keys()) == ["GLOBAL", "REPORTS", "ALLSTAR", "SYSTEMS"]
assert list(out["GLOBAL"].keys())[:3] == ["PATH", "PING_TIME", "SERVER_ID"]
assert list(out["SYSTEMS"].keys()) == ["ECHO", "OBP-CR"]
assert list(out["SYSTEMS"]["ECHO"].keys())[:4] == ["MODE", "ENABLED", "PORT", "TS2_STATIC"]
text = dump_yaml(out)
yaml_keys = [line.rstrip(":") for line in text.splitlines() if line.endswith(":") and not line.startswith("#")]
assert yaml_keys[:4] == ["GLOBAL", "REPORTS", "ALLSTAR", "SYSTEMS"]
def test_fixture_matches_cfg_section_order() -> None:
fixture = Path(__file__).resolve().parents[1] / "fixtures" / "sample_freedmr.cfg"
out = parse_freedmr_cfg(fixture)
expected_top = ["GLOBAL", "REPORTS", "LOGGER", "ALIASES", "ALLSTAR", "SYSTEMS"]
assert list(out.keys()) == expected_top
assert list(out["SYSTEMS"].keys()) == ["SYSTEM", "ECHO", "OBP-ES", "OBP-UY"]
assert list(out["GLOBAL"].keys())[:5] == ["PATH", "PING_TIME", "MAX_MISSED", "USE_ACL", "REG_ACL"]
def test_dump_yaml_is_parseable(tmp_path: Path) -> None:
cfg = _write_cfg(
tmp_path,
"""
[GLOBAL]
SERVER_ID: 1
[SYSTEM]
MODE: MASTER
ENABLED: True
PORT: 56400
PASSPHRASE: secret
GENERATOR: 1
""",
)
text = dump_yaml(parse_freedmr_cfg(cfg))
loaded = yaml.safe_load(text.split("\n", 2)[2])
assert loaded["SYSTEMS"]["SYSTEM"]["PORT"] == 56400

@ -1,20 +0,0 @@
"""Smoke tests for PacketSpec and DMR header parsing."""
from __future__ import annotations
from tests.harness.deterministic import PacketSpec, parse_dmr_fields
from adn_server.domain import int_id
def test_packet_spec_builds_valid_dmr_header() -> None:
spec = PacketSpec(dst_id=91, rf_src=3120001, peer_id=1001, seq=7)
packet = spec.data()
fields = parse_dmr_fields(packet)
assert fields["opcode"] == b"DMRD"
assert fields["seq"] == 7
assert int_id(fields["dst_id"]) == 91
assert int_id(fields["rf_src"]) == 3120001
assert fields["slot"] == 2
assert fields["call_type"] == "group"

@ -0,0 +1,128 @@
"""HBP MASTER + BridgeUseCases stack for REPEAT / talker-alias integration tests."""
from __future__ import annotations
import copy
from dataclasses import dataclass, field
from typing import Any
from adn_server.application.bridge_use_cases import BridgeUseCases
from adn_server.application.reporting_use_cases import ReportingUseCases
from adn_server.domain.dmr.bptc import encode_emblc
from adn_server.infrastructure.bridge_router_impl import InMemoryBridgeRouter
from adn_server.infrastructure.config_normalizer import (
apply_talker_alias_defaults,
ensure_system_runtime_config,
)
from adn_server.infrastructure.talker_alias_emblc import default_ta_emblc_encoder
from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol
from tests.harness.deterministic import FakeReportFactory, FakeReportSender, PacketSpec
from tests.harness.scenarios import talker_alias_config
class RecordingTransport:
"""Capture MASTER downlink UDP writes (REPEAT and DMRA)."""
def __init__(self) -> None:
self.sent: list[tuple[bytes, tuple[str, int]]] = []
def write(self, data: bytes, addr: tuple[str, int]) -> None:
self.sent.append((data, addr))
def for_addr(self, addr: tuple[str, int]) -> list[bytes]:
return [pkt for pkt, target in self.sent if target == addr]
def clear(self) -> None:
self.sent.clear()
class _AclRouter:
def acl_check(self, _value: bytes, _acl: object) -> bool:
return True
@dataclass
class HbpRepeatStack:
system_name: str
config: dict[str, Any]
hbp: HBPProtocol
bridge: BridgeUseCases
transport: RecordingTransport
dmra_capture: list[tuple[list[bytes], bytes | None]] = field(default_factory=list)
def register_peer(self, peer_id: bytes, sockaddr: tuple[str, int]) -> None:
self.hbp._peers[peer_id] = {
"CONNECTION": "YES",
"CONNECTED": 1_700_000_000.0,
"LAST_PING": 1_700_000_000.0,
"SOCKADDR": sockaddr,
"CALLSIGN": b"CE5RPY ",
"RADIO_ID": str(int.from_bytes(peer_id, "big")),
}
self.config["SYSTEMS"][self.system_name].setdefault("PEERS", {})[peer_id] = (
self.hbp._peers[peer_id]
)
def inject(self, packet: bytes, sockaddr: tuple[str, int]) -> None:
self.hbp.datagramReceived(packet, sockaddr)
def inject_spec(self, spec: PacketSpec, sockaddr: tuple[str, int]) -> None:
self.inject(spec.data(), sockaddr)
def build_hbp_repeat_stack(
*,
talker_alias: bool = True,
system_name: str = "MASTER-A",
) -> HbpRepeatStack:
"""Real HBPProtocol with BridgeUseCases TA callbacks (not FakeHbpProtocol)."""
config = copy.deepcopy(talker_alias_config())
if not talker_alias:
config["GLOBAL"]["TALKER_ALIAS"] = False
apply_talker_alias_defaults(config)
ensure_system_runtime_config(config)
sys_cfg = config["SYSTEMS"][system_name]
sys_cfg["REPEAT"] = True
sys_cfg["MAX_PEERS"] = 8
sys_cfg["USE_ACL"] = False
transport = RecordingTransport()
hbp = HBPProtocol(system_name, config, router=_AclRouter())
hbp.transport = transport # type: ignore[assignment]
protocols: dict[str, Any] = {system_name: hbp}
dmra_capture: list[tuple[list[bytes], bytes | None]] = []
def _send_dmra(target_system: str, packets: list[bytes], exclude_peer: bytes | None = None) -> int:
proto = protocols[target_system]
dmra_capture.append((list(packets), exclude_peer))
return proto.send_dmra_to_peers(packets, exclude_peer=exclude_peer)
def _get_dmra_blocks(_system: str, stream_id: bytes) -> dict[int, bytes] | None:
return hbp.get_dmra_blocks(stream_id)
report_factory = FakeReportFactory()
bridge = BridgeUseCases(
InMemoryBridgeRouter(),
config,
send_to_system=lambda *_a, **_k: None,
get_protocols=lambda: protocols,
reporting=ReportingUseCases(FakeReportSender(report_factory), config),
send_dmra_to_system=_send_dmra,
get_dmra_blocks=_get_dmra_blocks,
encode_emblc=encode_emblc,
ta_emblc_encoder=default_ta_emblc_encoder,
)
hbp._on_talker_alias_repeat_prepare = bridge.prepare_talker_alias_local_repeat
hbp._on_talker_alias_repeat_burst = bridge.rewrite_repeat_voice_burst
hbp._on_talker_alias_stream_end = bridge.clear_talker_alias_stream
return HbpRepeatStack(
system_name=system_name,
config=config,
hbp=hbp,
bridge=bridge,
transport=transport,
dmra_capture=dmra_capture,
)

@ -10,7 +10,6 @@ from tests.harness.deterministic import (
)
from tests.harness.scenarios import talker_alias_config
from adn_server.domain import bytes_3, bytes_4
from adn_server.domain.talker_alias import DMRA_BLOCK_COUNT, DMRA_OPCODE
@ -94,48 +93,3 @@ def test_both_mode_without_ta_still_rewrites_group_embedded_lc() -> None:
assert emb.any() # ...not the preserved all-zero source embedded LC
def test_talker_alias_local_repeat_excludes_source_peer() -> None:
scenario = DeterministicScenario(
config=talker_alias_config(),
bridges=active_bridge(91, (("MASTER-A", 2),)),
)
stream_id = bytes_4(0x92929292)
peer = bytes_4(1001)
rf_src = bytes_3(3120001)
scenario.bridge.send_talker_alias_local_repeat("MASTER-A", peer, rf_src, stream_id)
assert len(scenario.dmra_capture) == 1
assert scenario.dmra_capture[0].exclude_peer == peer
def test_repeat_embed_prepares_slot_state_for_mmdvm() -> None:
"""REPEAT must set TX_TA_EMB on the source MASTER slot for embedded downlink DMRD."""
from bitarray import bitarray
scenario = DeterministicScenario(
config=talker_alias_config(),
bridges=active_bridge(7304, (("MASTER-A", 2),)),
)
stream_id = bytes_4(0xA1A2A3A4)
peer = bytes_4(1001)
rf_src = bytes_3(7300392)
dst_id = bytes_3(7304)
scenario.bridge.prepare_talker_alias_local_repeat(
"MASTER-A", peer, rf_src, dst_id, 2, stream_id,
)
slot_st = scenario.protocols["MASTER-A"].STATUS[2]
assert slot_st.get("REP_STREAM_ID") == stream_id
assert "REP_EMB_LC" in slot_st
assert slot_st.get("TX_TA_EMB") is not None
raw = b"\x00" * 33
out = scenario.bridge.rewrite_repeat_voice_burst("MASTER-A", 2, stream_id, 1, raw)
bits = bitarray(endian="big")
bits.frombytes(out)
assert bits[116:148] == slot_st["REP_EMB_LC"][1]
scenario.bridge.clear_talker_alias_stream("MASTER-A", stream_id)
assert "REP_STREAM_ID" not in slot_st
assert "TX_TA_EMB" not in slot_st

@ -1,39 +0,0 @@
"""Talker Alias embedded LC alternation on voice bursts."""
from __future__ import annotations
from tests.harness.deterministic import DeterministicScenario, PacketSpec, active_bridge
from tests.harness.scenarios import talker_alias_config
def test_talker_alias_embed_state_prepared_on_bridge_vhead() -> None:
bridges = active_bridge(91, (("MASTER-A", 2), ("MASTER-B", 2)))
scenario = DeterministicScenario(config=talker_alias_config(), bridges=bridges)
base = PacketSpec(dst_id=91, rf_src=3120001, stream_id=0xABABABAB)
scenario.inject_hbp("MASTER-A", DeterministicScenario.voice_head_spec(base))
ts_st = scenario.protocols["MASTER-B"].STATUS.get(2, {})
assert "TX_TA_EMB" in ts_st
assert ts_st.get("TX_TA_ON") is False
assert ts_st.get("TX_TA_PHASE") == 0
def test_talker_alias_embed_cleared_on_vterm() -> None:
bridges = active_bridge(91, (("MASTER-A", 2), ("MASTER-B", 2)))
scenario = DeterministicScenario(config=talker_alias_config(), bridges=bridges)
base = PacketSpec(dst_id=91, rf_src=3120001, stream_id=0xCDCDCDCD)
scenario.inject_hbp("MASTER-A", DeterministicScenario.voice_head_spec(base))
scenario.inject_hbp(
"MASTER-A",
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1),
)
scenario.inject_hbp(
"MASTER-A",
DeterministicScenario.voice_term_spec(base, seq=99),
)
ts_st = scenario.protocols["MASTER-B"].STATUS.get(2, {})
assert "TX_TA_EMB" not in ts_st
assert "TX_TA_ON" not in ts_st
Loading…
Cancel
Save

Powered by TurnKey Linux.