fix: audit refactors, report contract, and concurrent OBP downlink (#36)

* fix: audit wave 1 server hygiene refactors

Inject call_later into PlaybackUseCases, move voice config mtime watch to
bootstrap, complete SubscriptionStore port methods, and relocate echo
routing seed to the application layer.

* fix: document dashboard_state in report-v2 schema

Add dashboard_state to report-v2.json with a two-master example fixture,
update public protocol docs for HELLO → STATE_SND connect flow, and drop
the unused TOPOLOGY_JSON HELLO feature token.

* fix: add ReportWire contract tests against report-v2 schema

Assert state_frames and bridge_event_frames output validates against
committed example fixtures and the report-v2 JSON schema.

Fix test imports for echo_seed module relocation.

* fix: align OPTIONS static validity checks across routing and report

Delegate subscription_table validation to peer_options_static_valid so empty
OPTIONS is valid and PASS-mixed strings are rejected consistently.

* fix: OBP DMRE source-server validation without ALLOW_UNREG_ID bypass

Port OPENBRIDGE.validate_id lookup for 6-7 digit source servers so OBP
ingress matches legacy production config (VALIDATE_SERVER_IDS=True).

* fix: audit items 11-13 coverage, infra tests, and warning logs

* fix: add tests/fakes shim for application test decoupling

* fix: per-stream OBP bridge TX legs for concurrent MASTER downlink

When two OBP voice streams share the same MASTER timeslot, stop
flip-flopping the flat TX row so per-peer downlink gates stay stable.

* fix: ruff lint in OBP concurrent streams downlink test

Remove dead code after return and unused start_tx_events variable.
pull/37/head
ce5rpy 3 months ago committed by GitHub
parent 0b9e536279
commit c15430a3e0
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

2
.gitignore vendored

@ -49,6 +49,8 @@ site/
# -----------------------------------------------------------------------------
# Python
# -----------------------------------------------------------------------------
htmlcov/
.coverage
__pycache__/
*.py[cod]
*$py.class

@ -1,12 +1,12 @@
# Report protocol v2 (JSON)
**Status:** schema draft. Wire encoding and server emission are in progress; monitor v2 consumer is a separate deliverable.
**Status:** schema + TCP slim wire (`STATE_SND`) shipped on adn-server 2.x; monitor 2.x consumer required for production pairing.
## Goals
Replace monitor snapshots that today use **pickle** (`CONFIG_SND`, `BRIDGE_SND`) and **CSV strings** (`BRDG_EVENT`) with **typed JSON** that any client can decode.
**Release policy:** **adn-server 1.0.x** + **adn-monitor 1.0.x** = report v1 (frozen tag pair). **2.x** server sends **HELLO** with `report_protocol: 2` and **v2 JSON only** (`TOPOLOGY_SND`, `ROUTING_TABLE_SND`, `VOICE_EVENT_SND`, `DELTA_SND`) — no duplicate pickle/CSV frames. **adn-monitor 2.x** is the polyglot consumer: it decodes v1 wire from legacy peers and v2 JSON from **adn-server 2.x**, mapping both into the same dashboard state for the UI. **adn-monitor 1.0.0** cannot decode v2 opcodes — upgrade the monitor line for **adn-server 2.x**.
**Release policy:** **adn-server 1.0.x** + **adn-monitor 1.0.x** = report v1 (frozen tag pair). **2.x** server sends **HELLO** with `report_protocol: 2` and v2 JSON on TCP. The monitor connect path is **`HELLO` → `STATE_SND` (`dashboard_state`)** plus async **`ROUTING_TABLE_SND` / `DELTA_SND`** for UA/SINGLE_TS chips and **`VOICE_EVENT_SND`** for live calls. **`topology` is internal-only** (used to build `dashboard_state`, never sent on TCP). **adn-monitor 2.x** decodes v1 wire from legacy peers and v2 JSON from **adn-server 2.x**. **adn-monitor 1.0.0** cannot decode v2 opcodes — upgrade the monitor line for **adn-server 2.x**.
## Transport (unchanged)
@ -18,16 +18,15 @@ Replace monitor snapshots that today use **pickle** (`CONFIG_SND`, `BRIDGE_SND`)
| Opcode | Hex | v1 payload | v2 payload |
|--------|-----|------------|------------|
| `HELLO` | `0xFF` | JSON hello (`protocol`: 1) | JSON hello (`report_protocol`: 2) |
| `CONFIG_SND` | `0x01` | pickle SYSTEMS | — (use `TOPOLOGY_SND`) |
| `CONFIG_SND` | `0x01` | pickle SYSTEMS | — (use `STATE_SND`) |
| `BRIDGE_SND` | `0x03` | pickle BRIDGES | — (use `ROUTING_TABLE_SND`) |
| `BRDG_EVENT` | `0x07` | CSV text | — (use `VOICE_EVENT_SND`) |
| `TOPOLOGY_SND` | `0x10` | — | JSON `topology` |
| `STATE_SND` | `0x14` | — | JSON `dashboard_state` |
| `TOPOLOGY_SND` | `0x10` | — | *(internal schema only — not emitted on TCP)* |
| `ROUTING_TABLE_SND` | `0x11` | — | JSON `routing_table` |
| `VOICE_EVENT_SND` | `0x12` | — | JSON `voice_event` |
| `DELTA_SND` | `0x13` | — | JSON `delta` |
Proposed opcodes `0x10`–`0x13` are reserved in the schema phase; exact values may change before P1-002 ships.
```mermaid
flowchart TB
subgraph v1 [Report v1 — adn-server 1.0.x + monitor 1.0.x]
@ -37,8 +36,8 @@ flowchart TB
end
subgraph v2 [Report v2 — adn-server 2.x + monitor 2.x]
H2[HELLO report_protocol 2] --> T2[TOPOLOGY_SND JSON]
T2 --> R2[ROUTING_TABLE_SND JSON]
H2[HELLO report_protocol 2] --> S2[STATE_SND dashboard_state]
S2 --> R2[ROUTING_TABLE_SND JSON async]
R2 --> V2[VOICE_EVENT_SND JSON]
R2 -.-> D2[DELTA_SND optional]
end
@ -54,7 +53,7 @@ On connect the server sends **`HELLO` (`0xFF`)** first (same as today). v2 clien
"server": "adn-server",
"version": "2.0.0-alpha.1",
"report_protocol": 2,
"features": ["INGRESS", "END_TX_FORWARD", "PUSH_ON_CONNECT", "REPORT_V2", "TOPOLOGY_JSON", "ROUTING_TABLE_JSON", "VOICE_EVENT_JSON", "DELTA_UPDATES"],
"features": ["INGRESS", "END_TX_FORWARD", "PUSH_ON_CONNECT", "REPORT_V2", "ROUTING_TABLE_JSON", "VOICE_EVENT_JSON", "DELTA_UPDATES"],
"systems": ["MASTER-A", "OBP-CL"]
}
```
@ -70,15 +69,64 @@ Monitor **1.0.x** does not speak this wire; use the **1.0.x** server tag for tha
| `type` | Replaces | Purpose |
|--------|----------|---------|
| `topology` | `CONFIG_SND` | Systems, peers, OpenBridge legs (no secrets). |
| `routing_table` | `BRIDGE_SND` | Active bridge legs per talkgroup / reflector key. |
| `dashboard_state` | `CONFIG_SND` (slim) | Linked masters/peers/openbridges snapshot (`STATE_SND`). **Full snapshot** — monitor prunes absent masters. |
| `topology` | `CONFIG_SND` (full) | Internal schema / ops only — **not** sent on TCP monitor wire. |
| `routing_table` | `BRIDGE_SND` | Active bridge legs per talkgroup / reflector key (UA/SINGLE_TS chips). |
| `voice_event` | `BRDG_EVENT` | Structured call start/end/ingress. |
| `delta` | — | Incremental `topology` or `routing_table` patch since `since_seq`. |
| `delta` | — | Incremental `routing_table` patch since `since_seq`. |
## Reference payloads
Each frame payload is one JSON object with a required `type`. Examples below use anonymized IDs.
### `dashboard_state`
```json
{
"type": "dashboard_state",
"ts": 1717555200.0,
"server_id": "7302",
"ctable": {
"MASTERS": {
"MASTER-A": {
"mode": "MASTER",
"ip": "10.0.0.1",
"port": 62030,
"peers": {
"3120001": {
"id": 3120001,
"connected": true,
"ip": "10.0.0.50",
"port": 62031,
"callsign": "CE5RPY"
}
}
}
},
"PEERS": {},
"OPENBRIDGES": {
"OBP-CL": {
"mode": "OPENBRIDGE",
"network_id": 73010,
"ip": "44.31.61.68",
"port": 62999,
"streams": {}
}
}
}
}
```
| Field | Type | Required | Notes |
|-------|------|----------|-------|
| `ts` | float | yes | Epoch time of snapshot. |
| `server_id` | string | no | `GLOBAL.SERVER_ID` (MQTT retained state). |
| `ctable.MASTERS` | object | yes | Masters with **connected** peers only. |
| `ctable.PEERS` | object | yes | Connected upstream PEER/XLXPEER systems. |
| `ctable.OPENBRIDGES` | object | yes | Enabled OpenBridge legs (`streams` empty; live chips from `voice_event`). |
Committed example: `schemas/examples/dashboard_state.json`.
### `topology`
```json
@ -133,7 +181,7 @@ No passwords or encryption material are included (unlike legacy pickle).
"ts": 1717555260.5,
"routes": [
{
"bridge_key": "52090",
"relay_table_key": "52090",
"legs": [
{
"system": "MASTER-A",
@ -153,7 +201,7 @@ No passwords or encryption material are included (unlike legacy pickle).
]
},
{
"bridge_key": "#310",
"relay_table_key": "#310",
"legs": [
{
"system": "MASTER-A",
@ -170,7 +218,7 @@ No passwords or encryption material are included (unlike legacy pickle).
| Field | Type | Required | Notes |
|-------|------|----------|-------|
| `routes[].bridge_key` | string | yes | Talkgroup id or reflector key (`#nnn`). |
| `routes[].relay_table_key` | string | yes | Talkgroup id or reflector key (`#nnn`). |
| `legs[].system` | string | yes | Target system name. |
| `legs[].ts` | 1 \| 2 | yes | Timeslot (legacy `TS`). |
| `legs[].tgid` | int | yes | Talkgroup (1–16777215). |
@ -231,7 +279,7 @@ OpenBridge **INGRESS** vs **START** semantics match [Monitoring and reports](../
"ts": 1717555261.0,
"routes": [
{
"bridge_key": "52090",
"relay_table_key": "52090",
"legs": [
{
"system": "MASTER-A",

@ -1,12 +1,12 @@
# Protocolo de informes v2 (JSON)
**Estado:** borrador de esquema. La codificación en wire y la emisión en el servidor están en progreso; el consumidor monitor v2 es un entregable aparte.
**Estado:** esquema + wire TCP slim (`STATE_SND`) en adn-server 2.x; pareja de producción requiere adn-monitor 2.x.
## Objetivos
Sustituir instantáneas al monitor que hoy usan **pickle** (`CONFIG_SND`, `BRIDGE_SND`) y **CSV** (`BRDG_EVENT`) por **JSON tipado**.
**Política de releases:** **adn-server 1.0.x** + **adn-monitor 1.0.x** = report v1 (tags congelados). El servidor **2.x** emite **solo report v2** (sin shim pickle ni `dual`); requiere **adn-monitor 2.x** en la misma línea.
**Política de releases:** **adn-server 1.0.x** + **adn-monitor 1.0.x** = report v1 (tags congelados). El servidor **2.x** emite **HELLO** con `report_protocol: 2` y JSON v2 en TCP. El flujo de conexión del monitor es **`HELLO` → `STATE_SND` (`dashboard_state`)** más **`ROUTING_TABLE_SND` / `DELTA_SND`** asíncronos para chips UA/SINGLE_TS y **`VOICE_EVENT_SND`** para llamadas en vivo. **`topology` es solo interno** (construye `dashboard_state`, no se emite en TCP). **adn-monitor 2.x** decodifica v1 de peers legacy y v2 de **adn-server 2.x**.
## Transporte (sin cambios)
@ -18,10 +18,11 @@ Sustituir instantáneas al monitor que hoy usan **pickle** (`CONFIG_SND`, `BRIDG
| Opcode | Hex | Payload v1 | Payload v2 |
|--------|-----|------------|------------|
| `HELLO` | `0xFF` | JSON hello (`protocol`: 1) | JSON hello (`report_protocol`: 2) |
| `CONFIG_SND` | `0x01` | pickle SYSTEMS | — (`TOPOLOGY_SND`) |
| `CONFIG_SND` | `0x01` | pickle SYSTEMS | — (`STATE_SND`) |
| `BRIDGE_SND` | `0x03` | pickle BRIDGES | — (`ROUTING_TABLE_SND`) |
| `BRDG_EVENT` | `0x07` | texto CSV | — (`VOICE_EVENT_SND`) |
| `TOPOLOGY_SND` | `0x10` | — | JSON `topology` |
| `STATE_SND` | `0x14` | — | JSON `dashboard_state` |
| `TOPOLOGY_SND` | `0x10` | — | *(solo esquema interno — no se emite en TCP)* |
| `ROUTING_TABLE_SND` | `0x11` | — | JSON `routing_table` |
| `VOICE_EVENT_SND` | `0x12` | — | JSON `voice_event` |
| `DELTA_SND` | `0x13` | — | JSON `delta` |
@ -37,8 +38,8 @@ flowchart TB
end
subgraph v2 [Informe v2 — adn-server 2.x + monitor 2.x]
H2[HELLO report_protocol 2] --> T2[TOPOLOGY_SND JSON]
T2 --> R2[ROUTING_TABLE_SND JSON]
H2[HELLO report_protocol 2] --> S2[STATE_SND dashboard_state]
S2 --> R2[ROUTING_TABLE_SND JSON async]
R2 --> V2[VOICE_EVENT_SND JSON]
R2 -.-> D2[DELTA_SND opcional]
end
@ -54,7 +55,7 @@ Al conectar, el servidor envía **`HELLO` (`0xFF`)** primero. Clientes v2 miran
"server": "adn-server",
"version": "2.0.0-alpha.1",
"report_protocol": 2,
"features": ["INGRESS", "END_TX_FORWARD", "PUSH_ON_CONNECT", "REPORT_V2", "TOPOLOGY_JSON", "ROUTING_TABLE_JSON", "VOICE_EVENT_JSON", "DELTA_UPDATES"],
"features": ["INGRESS", "END_TX_FORWARD", "PUSH_ON_CONNECT", "REPORT_V2", "ROUTING_TABLE_JSON", "VOICE_EVENT_JSON", "DELTA_UPDATES"],
"systems": ["MASTER-A", "OBP-CL"]
}
```
@ -68,15 +69,20 @@ Al conectar, el servidor envía **`HELLO` (`0xFF`)** primero. Clientes v2 miran
| `type` | Sustituye | Uso |
|--------|-----------|-----|
| `topology` | `CONFIG_SND` | Sistemas, peers, piernas OBP (sin secretos). |
| `routing_table` | `BRIDGE_SND` | Piernas de bridge por TG / reflector. |
| `dashboard_state` | `CONFIG_SND` (slim) | Instantánea de sistemas enlazados (`STATE_SND`). **Snapshot completo** — el monitor elimina masters ausentes. |
| `topology` | `CONFIG_SND` (completo) | Solo esquema interno / ops — **no** se envía en TCP al monitor. |
| `routing_table` | `BRIDGE_SND` | Piernas de bridge por TG / reflector (chips UA/SINGLE_TS). |
| `voice_event` | `BRDG_EVENT` | Inicio/fin/ingress de llamadas. |
| `delta` | — | Parche incremental desde `since_seq`. |
| `delta` | — | Parche incremental de `routing_table` desde `since_seq`. |
## Payloads de referencia
Cada trama lleva un único objeto JSON con `type` obligatorio. Los ejemplos usan IDs anonimizados.
### `dashboard_state`
Ver ejemplo completo en `schemas/examples/dashboard_state.json` (mismo objeto que MQTT `{prefix}/state`).
### `topology`
```json

@ -1,23 +1,23 @@
[pytest]
testpaths = tests
pythonpath = src
pythonpath = src tests
addopts = -ra
markers =
smoke: minimal wiring checks; not sufficient alone for regression
behavior: observable regression test (preferred for bridge/voice integration)
unit: pure domain or isolated helper; private methods acceptable
udp_blackbox: opt-in subprocess UDP integration tests (not yet implemented)
integration: real HBP/proxy stack (not DeterministicScenario inject)
mqtt: requires MQTT broker or heavy mqtt mocks
# Coverage headline excludes the rows below (I/O-heavy adapters still need lab/integration).
# Regenerate: python -m pytest --cov=adn_server --cov-report=html
# htmlcov/ is gitignored; delete stale trees before trusting the report.
[coverage:run]
source = adn_server
omit =
*/main.py
*/parrot_main.py
*/infrastructure/twisted_adapters/*
*/infrastructure/persistence/*
*/infrastructure/security/*
*/infrastructure/voice/*
*/infrastructure/logging_config.py
*/infrastructure/config_loader.py

@ -0,0 +1,64 @@
{
"type": "dashboard_state",
"ts": 1717555200.0,
"server_id": "7302",
"ctable": {
"MASTERS": {
"MASTER-A": {
"mode": "MASTER",
"ip": "10.0.0.1",
"port": 62030,
"single_mode": false,
"default_ua_timer": 10.0,
"peers": {
"3120001": {
"id": 3120001,
"connected": true,
"ip": "10.0.0.50",
"port": 62031,
"callsign": "CE5RPY",
"connected_at": 1717555200,
"ts1_static": ["91", "92"],
"ts2_static": ["7302"],
"single_mode": false,
"ua_timer_min": 10.0,
"ua_sessions": {},
"ua_multi_tgs": {},
"rf_mode": "duplex"
}
}
},
"MASTER-B": {
"mode": "MASTER",
"ip": "10.0.0.2",
"port": 62032,
"peers": {
"3120002": {
"id": 3120002,
"connected": true,
"ip": "10.0.0.51",
"port": 62033,
"callsign": "EA5GVK",
"connected_at": 1717555200,
"single_mode": false,
"ua_timer_min": 10.0,
"ua_sessions": {},
"ua_multi_tgs": {},
"rf_mode": "duplex"
}
}
}
},
"PEERS": {},
"OPENBRIDGES": {
"OBP-CL": {
"mode": "OPENBRIDGE",
"network_id": 73010,
"ip": "44.31.61.68",
"port": 62999,
"enhanced_obp": true,
"streams": {}
}
}
}
}

@ -8,7 +8,6 @@
"END_TX_FORWARD",
"PUSH_ON_CONNECT",
"REPORT_V2",
"TOPOLOGY_JSON",
"ROUTING_TABLE_JSON",
"VOICE_EVENT_JSON",
"DELTA_UPDATES"

@ -7,6 +7,7 @@
"required": ["type"],
"oneOf": [
{ "$ref": "#/$defs/hello" },
{ "$ref": "#/$defs/dashboard_state" },
{ "$ref": "#/$defs/topology" },
{ "$ref": "#/$defs/routing_table" },
{ "$ref": "#/$defs/voice_event" },
@ -55,7 +56,7 @@
"type": "array",
"items": { "type": "string" },
"uniqueItems": true,
"description": "Capability tokens. v1: INGRESS, END_TX_FORWARD, PUSH_ON_CONNECT. v2 adds REPORT_V2, TOPOLOGY_JSON, ROUTING_TABLE_JSON, VOICE_EVENT_JSON, DELTA_UPDATES."
"description": "Capability tokens. v1: INGRESS, END_TX_FORWARD, PUSH_ON_CONNECT. v2 adds REPORT_V2, ROUTING_TABLE_JSON, VOICE_EVENT_JSON, DELTA_UPDATES. TCP monitor uses dashboard_state (STATE_SND), not topology."
},
"systems": {
"type": "array",
@ -64,6 +65,132 @@
}
}
},
"dashboard_peer": {
"type": "object",
"additionalProperties": true,
"required": ["id", "connected"],
"properties": {
"id": { "$ref": "#/$defs/dmr_id" },
"connected": { "type": "boolean" },
"ip": { "type": "string" },
"port": { "type": "integer", "minimum": 0, "maximum": 65535 },
"callsign": { "type": "string" },
"connected_at": { "type": "integer" },
"options": { "type": "string" },
"ts1_static": {
"type": "array",
"items": { "type": "string" }
},
"ts2_static": {
"type": "array",
"items": { "type": "string" }
},
"single_mode": { "type": "boolean" },
"ua_timer_min": { "type": "number", "minimum": 0 },
"ua_sessions": {
"type": "object",
"additionalProperties": {
"type": "object",
"additionalProperties": false,
"required": ["tgid", "expires_at"],
"properties": {
"tgid": { "$ref": "#/$defs/dmr_id" },
"expires_at": { "type": "number" }
}
}
},
"ua_multi_tgs": {
"type": "object",
"additionalProperties": {
"type": "array",
"items": { "type": "string" }
}
},
"rf_mode": {
"type": "string",
"enum": ["simplex", "duplex"]
}
}
},
"dashboard_master": {
"type": "object",
"additionalProperties": false,
"required": ["mode", "peers"],
"properties": {
"mode": { "const": "MASTER" },
"peers": {
"type": "object",
"additionalProperties": { "$ref": "#/$defs/dashboard_peer" }
},
"single_mode": { "type": "boolean" },
"default_ua_timer": { "type": "number", "minimum": 0 },
"ip": { "type": "string" },
"port": { "type": "integer", "minimum": 0, "maximum": 65535 }
}
},
"dashboard_upstream_peer": {
"type": "object",
"additionalProperties": false,
"required": ["mode", "connected"],
"properties": {
"mode": { "type": "string", "enum": ["PEER", "XLXPEER"] },
"connected": { "type": "boolean" },
"callsign": { "type": "string" },
"location": { "type": "string" },
"description": { "type": "string" },
"url": { "type": "string" },
"master_ip": { "type": "string" },
"master_port": { "type": "integer", "minimum": 0, "maximum": 65535 },
"radio_id": { "$ref": "#/$defs/dmr_id" },
"connected_at": { "type": "integer" }
}
},
"dashboard_openbridge": {
"type": "object",
"additionalProperties": false,
"required": ["mode", "streams"],
"properties": {
"mode": { "const": "OPENBRIDGE" },
"streams": { "type": "object" },
"network_id": { "$ref": "#/$defs/dmr_id" },
"ip": { "type": "string" },
"port": { "type": "integer", "minimum": 0, "maximum": 65535 },
"enhanced_obp": { "type": "boolean" }
}
},
"dashboard_ctable": {
"type": "object",
"additionalProperties": false,
"required": ["MASTERS", "PEERS", "OPENBRIDGES"],
"properties": {
"MASTERS": {
"type": "object",
"additionalProperties": { "$ref": "#/$defs/dashboard_master" }
},
"PEERS": {
"type": "object",
"additionalProperties": { "$ref": "#/$defs/dashboard_upstream_peer" }
},
"OPENBRIDGES": {
"type": "object",
"additionalProperties": { "$ref": "#/$defs/dashboard_openbridge" }
}
}
},
"dashboard_state": {
"type": "object",
"additionalProperties": false,
"required": ["type", "ts", "ctable"],
"properties": {
"type": { "const": "dashboard_state" },
"ts": { "type": "number", "description": "Unix epoch seconds (fractional allowed)." },
"server_id": {
"type": "string",
"description": "GLOBAL.SERVER_ID when present (MQTT retained state)."
},
"ctable": { "$ref": "#/$defs/dashboard_ctable" }
}
},
"topology_system": {
"type": "object",
"additionalProperties": false,

@ -26,13 +26,11 @@
from __future__ import annotations
import logging
from collections.abc import Callable
from random import randint
from time import time
from typing import Any
from twisted.internet import reactor
from twisted.internet.base import DelayedCall
from ..domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, HBPF_VOICE, bytes_4, int_id
logger = logging.getLogger(__name__)
@ -49,8 +47,15 @@ _SOURCE_MAX_S = 180.0
class PlaybackUseCases:
"""Legacy playback class behaviour; playback is scheduled on the reactor (non-blocking)."""
def __init__(self, system_name: str, get_protocol: Any = None) -> None:
def __init__(
self,
system_name: str,
*,
call_later: Callable[..., Any],
get_protocol: Any = None,
) -> None:
self._system = system_name
self._call_later = call_later
self._get_protocol = get_protocol
self.STATUS: dict[str, Any] = {}
self.CALL_DATA: list[bytes] = []
@ -66,10 +71,10 @@ class PlaybackUseCases:
self._playback_index = 0
self._playback_stream_id = b""
self._playback_ta_from_stream = b""
self._delay_call: DelayedCall | None = None
self._packet_call: DelayedCall | None = None
self._idle_call: DelayedCall | None = None
self._max_call: DelayedCall | None = None
self._delay_call: Any = None
self._packet_call: Any = None
self._idle_call: Any = None
self._max_call: Any = None
def dmrd_received(
self,
@ -219,7 +224,7 @@ class PlaybackUseCases:
self._cancel_record_idle()
if not self._recording_active:
return
self._idle_call = reactor.callLater(_RECORD_IDLE_S, self._on_record_idle, proto)
self._idle_call = self._call_later(_RECORD_IDLE_S, self._on_record_idle, proto)
def _cancel_record_idle(self) -> None:
if self._idle_call is not None and self._idle_call.active():
@ -277,7 +282,7 @@ class PlaybackUseCases:
if not recorded:
return
self._playback_busy = True
self._delay_call = reactor.callLater(
self._delay_call = self._call_later(
_PLAYBACK_DELAY_S,
self._start_playback,
proto,
@ -327,7 +332,7 @@ class PlaybackUseCases:
self._cancel_record_max()
if not self._recording_active:
return
self._max_call = reactor.callLater(_SOURCE_MAX_S, self._on_record_max_duration, proto)
self._max_call = self._call_later(_SOURCE_MAX_S, self._on_record_max_duration, proto)
def _cancel_record_max(self) -> None:
if self._max_call is not None and self._max_call.active():
@ -454,7 +459,7 @@ class PlaybackUseCases:
else:
logger.warning("(%s) Playback packet %d dropped (no protocol/send_system)", self._system, self._playback_index)
self._playback_index += 1
self._packet_call = reactor.callLater(_PACKET_INTERVAL_S, self._send_next_packet, proto)
self._packet_call = self._call_later(_PACKET_INTERVAL_S, self._send_next_packet, proto)
def _finish_playback(self) -> None:
if self._delay_call is not None and self._delay_call.active():

@ -285,6 +285,23 @@ class SubscriptionStore(ABC):
def list_by_phase(self, phase: "SubscriptionPhase") -> tuple["Subscription", ...]:
...
@abstractmethod
def relay_tables_with_active_source(
self,
system: str,
slot: int,
dst_tgid: int,
) -> tuple[str, ...]:
...
@abstractmethod
def legs_in_table(self, table_key: str) -> tuple["Subscription", ...]:
...
@abstractmethod
def has_active_target_leg(self, system: str, slot: int, tgid: int) -> bool:
...
class ProxySlotStore(ABC):
"""Hotspot session registry keyed by peer_id (Phase 3)."""

@ -40,7 +40,6 @@ REPORT_FEATURES = (
"END_TX_FORWARD",
"PUSH_ON_CONNECT",
"REPORT_V2",
"TOPOLOGY_JSON",
"ROUTING_TABLE_JSON",
"VOICE_EVENT_JSON",
"DELTA_UPDATES",

@ -461,6 +461,112 @@ def inject_only_defer_obp_hbp_slot_contention(
return source_is_hbp and connected_count > 1
_OBP_FLAT_TX_KEYS = (
"TX_START",
"TX_TGID",
"TX_STREAM_ID",
"TX_RFS",
"TX_PEER",
"TX_H_LC",
"TX_T_LC",
"TX_EMB_LC",
"TX_TIME",
"TX_TYPE",
)
def obp_deferred_bridge_tx_leg(slot_st: dict[str, Any], stream_id: bytes) -> dict[str, Any]:
"""Per-stream bridge TX stamp for inject-only OBP→MASTER legs on a shared slot row."""
legs = slot_st.setdefault("TX_STREAMS", {})
if not isinstance(legs, dict):
legs = {}
slot_st["TX_STREAMS"] = legs
return legs.setdefault(stream_id, {})
def obp_flat_bridge_tx_idle(slot_st: dict[str, Any], pkt_time: float) -> bool:
"""True when flat ``TX_*`` on the slot is not carrying an active bridge leg."""
tx_type = slot_st.get("TX_TYPE")
if tx_type is None or tx_type == HBPF_SLT_VTERM:
return True
tx_time = float(slot_st.get("TX_TIME", 0) or 0)
return pkt_time >= tx_time + STREAM_TO
def obp_bridge_tx_leg_active(leg: dict[str, Any], pkt_time: float) -> bool:
tx_type = leg.get("TX_TYPE")
if tx_type is None or tx_type == HBPF_SLT_VTERM:
return False
tx_time = float(leg.get("TX_TIME", 0) or 0)
return pkt_time < tx_time + STREAM_TO
def obp_publish_flat_bridge_tx(slot_st: dict[str, Any], leg: dict[str, Any]) -> None:
"""Mirror one per-stream bridge leg onto the legacy flat ``TX_*`` slot row."""
for key in _OBP_FLAT_TX_KEYS:
if key in leg:
slot_st[key] = leg[key]
def obp_sync_flat_bridge_tx_times(
slot_st: dict[str, Any],
leg: dict[str, Any],
stream_id: bytes,
pkt_time: float,
dtype_vseq: int,
) -> None:
"""Refresh per-stream leg activity; mirror times to flat row only for the owner stream."""
leg["TX_TIME"] = pkt_time
leg["TX_TYPE"] = dtype_vseq
if slot_st.get("TX_STREAM_ID") == stream_id:
slot_st["TX_TIME"] = pkt_time
slot_st["TX_TYPE"] = dtype_vseq
def obp_pick_active_bridge_tx_leg(
slot_st: dict[str, Any],
pkt_time: float,
) -> tuple[bytes | None, dict[str, Any] | None]:
legs = slot_st.get("TX_STREAMS")
if not isinstance(legs, dict):
return None, None
best_sid: bytes | None = None
best_leg: dict[str, Any] | None = None
best_time = -1.0
for sid, leg in legs.items():
if not isinstance(leg, dict) or not obp_bridge_tx_leg_active(leg, pkt_time):
continue
t = float(leg.get("TX_TIME", 0) or 0)
if t > best_time:
best_time = t
best_sid = sid
best_leg = leg
return best_sid, best_leg
def obp_clear_flat_bridge_tx(slot_st: dict[str, Any]) -> None:
slot_st["TX_TYPE"] = HBPF_SLT_VTERM
slot_st["TX_STREAM_ID"] = b"\x00"
slot_st["TX_TIME"] = 0.0
def obp_clear_deferred_bridge_tx_leg(
slot_st: dict[str, Any],
stream_id: bytes,
pkt_time: float,
) -> None:
"""End one per-stream OBP bridge leg; republish flat ``TX_*`` for another active leg if any."""
legs = slot_st.get("TX_STREAMS")
if isinstance(legs, dict):
legs.pop(stream_id, None)
if slot_st.get("TX_STREAM_ID") == stream_id:
repl_sid, repl_leg = obp_pick_active_bridge_tx_leg(slot_st, pkt_time)
if repl_leg is not None and repl_sid is not None:
obp_publish_flat_bridge_tx(slot_st, repl_leg)
else:
obp_clear_flat_bridge_tx(slot_st)
def hbp_slot_blocks_group_voice_for_peer(
slot_st: dict[str, Any],
peer_id: bytes,

@ -434,14 +434,9 @@ class SubscriptionTableMixin:
def _options_static_lists_valid(self, opt_str: bytes | str) -> bool:
"""Legacy: malformed TS1/TS2 in OPTIONS aborts static bridge refresh."""
parsed = self._parse_options_string(opt_str)
if not parsed:
return False
for key in ("TS1_STATIC", "TS2_STATIC"):
val = str(parsed.get(key) or "").strip()
if val and re.search(r"[^\d,]", val):
return False
return True
from adn_server.application.report.payloads import peer_options_static_valid
return peer_options_static_valid(opt_str)
def _should_apply_system_single_from_options(self, system_name: str) -> bool:
"""Whether peer OPTIONS may overwrite system ``SINGLE_MODE`` (legacy single-hotspot only).

@ -48,6 +48,7 @@ import time
from typing import Any
from ...domain import HBPF_SLT_VTERM, int_id
from .helpers import obp_clear_deferred_bridge_tx_leg
logger = logging.getLogger(__name__)
@ -313,6 +314,26 @@ class RoutingTimerMixin:
)
if _slot.get("RX_TIME", 0) < now - 60:
_slot["RX_STREAM_ID"] = b"\x00"
tx_streams = _slot.get("TX_STREAMS")
if isinstance(tx_streams, dict):
for sid, leg in list(tx_streams.items()):
if not isinstance(leg, dict):
continue
if leg.get("TX_TYPE") != HBPF_SLT_VTERM and leg.get("TX_TIME", 0) < now - 5:
logger.debug(
"(%s) *TIME OUT* TX STREAM ID: %s SUB: %s TGID %s, TS %s, Duration: %.2f",
system_name, int_id(sid), int_id(leg.get("TX_RFS", b"")),
int_id(leg.get("TX_TGID", b"")), slot,
leg.get("TX_TIME", 0) - leg.get("TX_START", 0),
)
self._send_routing_event(
"GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}".format(
system_name, int_id(sid), int_id(leg.get("TX_PEER", b"")),
int_id(leg.get("TX_RFS", b"")), slot, int_id(leg.get("TX_TGID", b"")),
leg.get("TX_TIME", 0) - leg.get("TX_START", 0),
)
)
obp_clear_deferred_bridge_tx_leg(_slot, sid, now)
if _slot.get("TX_TYPE") != HBPF_SLT_VTERM and _slot.get("TX_TIME", 0) < now - 5:
_slot["TX_TYPE"] = HBPF_SLT_VTERM
logger.debug(

@ -46,6 +46,11 @@ from .routing.helpers import (
inject_only_defer_obp_hbp_slot_contention,
is_private_subscriber_dst,
is_unit_data_ingress,
obp_clear_deferred_bridge_tx_leg,
obp_deferred_bridge_tx_leg,
obp_flat_bridge_tx_idle,
obp_publish_flat_bridge_tx,
obp_sync_flat_bridge_tx_times,
obp_target_bcsq_quenches_stream,
resolve_voice_peer_id,
slot_has_active_voice,
@ -134,8 +139,8 @@ class RoutingUseCases(
return
try:
self._reporting.send_routing_table(self._routing_table_for_report(), incremental=incremental)
except Exception:
pass
except Exception as e:
logger.warning("(ROUTER) send_routing_table failed: %s", e)
def routing_table_for_report(self) -> dict[str, list[dict[str, Any]]]:
"""Return BRIDGES export for monitor/report only (routing uses ``SubscriptionStore``)."""
@ -611,15 +616,48 @@ class RoutingUseCases(
_src_stream_st = getattr(src_proto, "STATUS", {}).get(stream_id, {}) if src_proto else {}
else:
_src_stream_st = getattr(src_proto, "STATUS", {}).get(slot, {}) if src_proto else {}
_closing_bridge_leg = (
frame_type == HBPF_DATA_SYNC
and dtype_vseq == HBPF_SLT_VTERM
and _ts_st.get("TX_STREAM_ID") == stream_id
and _ts_st.get("TX_TYPE") != HBPF_SLT_VTERM
# Slot contention: active QSO blocks any other stream; post-VTERM uses GROUP_HANGTIME.
_group_hangtime = float(_target_system.get("GROUP_HANGTIME", 0) or 0)
_tgt_peers = getattr(tgt_proto, "_peers", None) if tgt_proto else None
if _tgt_peers is None:
_tgt_peers = _target_system.get("PEERS", {})
_target_connected = (
count_connected_peers(_tgt_peers) if isinstance(_tgt_peers, dict) else 0
)
_defer_slot_contention = inject_only_defer_obp_hbp_slot_contention(
self._config,
entry["SYSTEM"],
_target_system,
source_is_obp=source_is_obp,
source_is_hbp=not source_is_obp,
connected_count=_target_connected,
)
_obp_deferred_bridge = _defer_slot_contention and source_is_obp
_bridge_tx_leg: dict[str, Any] | None = None
if _obp_deferred_bridge:
_bridge_tx_leg = obp_deferred_bridge_tx_leg(_ts_st, stream_id)
_lc_row = _bridge_tx_leg
_closing_bridge_leg = (
frame_type == HBPF_DATA_SYNC
and dtype_vseq == HBPF_SLT_VTERM
and _bridge_tx_leg.get("TX_TYPE") is not None
and _bridge_tx_leg.get("TX_TYPE") != HBPF_SLT_VTERM
)
else:
_lc_row = _ts_st
_closing_bridge_leg = (
frame_type == HBPF_DATA_SYNC
and dtype_vseq == HBPF_SLT_VTERM
and _ts_st.get("TX_STREAM_ID") == stream_id
and _ts_st.get("TX_TYPE") != HBPF_SLT_VTERM
)
if _closing_bridge_leg:
call_duration = pkt_time - _ts_st.get("TX_START", pkt_time)
_end_peer = _ts_st.get("TX_PEER", peer_id)
if _obp_deferred_bridge and _bridge_tx_leg is not None:
call_duration = pkt_time - _bridge_tx_leg.get("TX_START", pkt_time)
_end_peer = _bridge_tx_leg.get("TX_PEER", peer_id)
else:
call_duration = pkt_time - _ts_st.get("TX_START", pkt_time)
_end_peer = _ts_st.get("TX_PEER", peer_id)
_end_report_peer = int_id(_end_peer)
if not source_is_obp:
_end_report_peer = int_id(
@ -638,23 +676,8 @@ class RoutingUseCases(
call_duration,
)
)
_ts_st["TX_TYPE"] = HBPF_SLT_VTERM
# Slot contention: active QSO blocks any other stream; post-VTERM uses GROUP_HANGTIME.
_group_hangtime = float(_target_system.get("GROUP_HANGTIME", 0) or 0)
_tgt_peers = getattr(tgt_proto, "_peers", None) if tgt_proto else None
if _tgt_peers is None:
_tgt_peers = _target_system.get("PEERS", {})
_target_connected = (
count_connected_peers(_tgt_peers) if isinstance(_tgt_peers, dict) else 0
)
_defer_slot_contention = inject_only_defer_obp_hbp_slot_contention(
self._config,
entry["SYSTEM"],
_target_system,
source_is_obp=source_is_obp,
source_is_hbp=not source_is_obp,
connected_count=_target_connected,
)
if not _obp_deferred_bridge:
_ts_st["TX_TYPE"] = HBPF_SLT_VTERM
if (
not _defer_slot_contention
and hbp_slot_blocks_group_voice(
@ -690,32 +713,56 @@ class RoutingUseCases(
int_id(_ts_st.get("RX_TGID", b"") or _ts_st.get("TX_TGID", b"")),
)
continue
# New stream detection — legacy OBP uses _target_status[TS]['TX_STREAM_ID'], legacy HBP uses self.STATUS[_slot]['RX_STREAM_ID']
if source_is_obp:
# New stream detection — legacy OBP uses flat TX_STREAM_ID; deferred OBP uses per-stream legs.
if _obp_deferred_bridge and _bridge_tx_leg is not None:
_is_new_stream = "TX_STREAM_ID" not in _bridge_tx_leg
elif source_is_obp:
_is_new_stream = (_ts_st.get("TX_STREAM_ID") != stream_id)
else:
src_status = getattr(src_proto, "STATUS", None) if src_proto else None
_is_new_stream = (src_status.get(slot, {}).get("RX_STREAM_ID") if src_status else b"") != stream_id
if _is_new_stream:
_ts_st["TX_START"] = pkt_time
_ts_st["TX_TGID"] = entry_tgid_b
_ts_st["TX_STREAM_ID"] = stream_id
_ts_st["TX_RFS"] = rf_src
_ts_st["TX_PEER"] = peer_id
dst_lc = source_lc[0:3] + entry_tgid_b + rf_src
_ts_st["TX_H_LC"] = bptc.encode_header_lc(dst_lc)
_ts_st["TX_T_LC"] = bptc.encode_terminator_lc(dst_lc)
_ts_st["TX_EMB_LC"] = self._encode_emblc(dst_lc)
self._dispatch_talker_alias_on_bridge_open(
_ts_st,
system_name,
entry["SYSTEM"],
rf_src,
stream_id,
peer_id,
int(entry_ts),
int_id(entry_tgid_b),
)
if _obp_deferred_bridge and _bridge_tx_leg is not None:
_bridge_tx_leg["TX_START"] = pkt_time
_bridge_tx_leg["TX_TGID"] = entry_tgid_b
_bridge_tx_leg["TX_STREAM_ID"] = stream_id
_bridge_tx_leg["TX_RFS"] = rf_src
_bridge_tx_leg["TX_PEER"] = peer_id
_bridge_tx_leg["TX_H_LC"] = bptc.encode_header_lc(dst_lc)
_bridge_tx_leg["TX_T_LC"] = bptc.encode_terminator_lc(dst_lc)
_bridge_tx_leg["TX_EMB_LC"] = self._encode_emblc(dst_lc)
if obp_flat_bridge_tx_idle(_ts_st, pkt_time) or _ts_st.get("TX_STREAM_ID") == stream_id:
obp_publish_flat_bridge_tx(_ts_st, _bridge_tx_leg)
self._dispatch_talker_alias_on_bridge_open(
_bridge_tx_leg,
system_name,
entry["SYSTEM"],
rf_src,
stream_id,
peer_id,
int(entry_ts),
int_id(entry_tgid_b),
)
else:
_ts_st["TX_START"] = pkt_time
_ts_st["TX_TGID"] = entry_tgid_b
_ts_st["TX_STREAM_ID"] = stream_id
_ts_st["TX_RFS"] = rf_src
_ts_st["TX_PEER"] = peer_id
_ts_st["TX_H_LC"] = bptc.encode_header_lc(dst_lc)
_ts_st["TX_T_LC"] = bptc.encode_terminator_lc(dst_lc)
_ts_st["TX_EMB_LC"] = self._encode_emblc(dst_lc)
self._dispatch_talker_alias_on_bridge_open(
_ts_st,
system_name,
entry["SYSTEM"],
rf_src,
stream_id,
peer_id,
int(entry_ts),
int_id(entry_tgid_b),
)
logger.info(
"(%s) Conference Bridge: %s, Call Bridged to HBP System: %s TS: %s, TGID: %s",
system_name, _relay_table_key, entry["SYSTEM"], entry_ts, int_id(entry_tgid_b),
@ -730,8 +777,13 @@ class RoutingUseCases(
int_id(entry_tgid_b),
)
)
_ts_st["TX_TIME"] = pkt_time
_ts_st["TX_TYPE"] = dtype_vseq
if _obp_deferred_bridge and _bridge_tx_leg is not None:
obp_sync_flat_bridge_tx_times(
_ts_st, _bridge_tx_leg, stream_id, pkt_time, dtype_vseq,
)
else:
_ts_st["TX_TIME"] = pkt_time
_ts_st["TX_TYPE"] = dtype_vseq
# Slot bit rewrite (legacy bridge.py 457-460 / 770-773)
_src_entry_ts = slot
if _src_entry_ts != entry_ts:
@ -743,12 +795,16 @@ class RoutingUseCases(
dmrbits = bitarray(endian="big")
dmrbits.frombytes(dmrpkt)
if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD:
dmrbits = _ts_st["TX_H_LC"][0:98] + dmrbits[98:166] + _ts_st["TX_H_LC"][98:197]
dmrbits = _lc_row["TX_H_LC"][0:98] + dmrbits[98:166] + _lc_row["TX_H_LC"][98:197]
elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM:
dmrbits = _ts_st["TX_T_LC"][0:98] + dmrbits[98:166] + _ts_st["TX_T_LC"][98:197]
dmrbits = _lc_row["TX_T_LC"][0:98] + dmrbits[98:166] + _lc_row["TX_T_LC"][98:197]
if not _closing_bridge_leg:
call_duration = pkt_time - _ts_st.get("TX_START", pkt_time)
_end_peer = _ts_st.get("TX_PEER", peer_id)
if _obp_deferred_bridge and _bridge_tx_leg is not None:
call_duration = pkt_time - _bridge_tx_leg.get("TX_START", pkt_time)
_end_peer = _bridge_tx_leg.get("TX_PEER", peer_id)
else:
call_duration = pkt_time - _ts_st.get("TX_START", pkt_time)
_end_peer = _ts_st.get("TX_PEER", peer_id)
_end_report_peer = int_id(_end_peer)
if not source_is_obp:
_end_report_peer = int_id(
@ -767,10 +823,11 @@ class RoutingUseCases(
call_duration,
)
)
_ts_st["TX_TYPE"] = HBPF_SLT_VTERM
if not _obp_deferred_bridge:
_ts_st["TX_TYPE"] = HBPF_SLT_VTERM
elif dtype_vseq in (1, 2, 3, 4):
self._rewrite_embed_lc(
dmrbits, _ts_st, dtype_vseq, "TX_EMB_LC",
dmrbits, _lc_row, dtype_vseq, "TX_EMB_LC",
)
dmrpkt_out = dmrbits.tobytes()
# bridge_master.routerOBP.to_target HBP branch: _tmp_data + dmrpkt only (~2041-2042);
@ -789,7 +846,11 @@ class RoutingUseCases(
except Exception as e:
logger.warning("(ROUTER) send_to_system %s failed: %s", entry.get("SYSTEM"), e)
if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM:
self._clear_talker_alias_embed(_ts_st)
if _obp_deferred_bridge and _bridge_tx_leg is not None:
self._clear_talker_alias_embed(_bridge_tx_leg)
obp_clear_deferred_bridge_tx_leg(_ts_st, stream_id, pkt_time)
else:
self._clear_talker_alias_embed(_ts_st)
self._talker_alias.clear_stream(entry["SYSTEM"], stream_id)
# Legacy bridge_master routerOBP ~2420-2434: after to_target, VTERM — CALL END log, END RX report, _fin, lastSeq
if (

@ -0,0 +1,87 @@
# ADN DMR Peer Server - application subscription echo seed
#
# 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
###############################################################################
"""Initial BRIDGES snapshot for the ECHO parrot system (bootstrap only)."""
from __future__ import annotations
import time
from typing import Any
from adn_server.domain import bytes_3
def seed_echo_routing_table(config: dict[str, Any]) -> dict[str, list[dict[str, Any]]]:
"""Initial BRIDGES for ECHO system (legacy make_bridges 9990 + MASTER expansion)."""
now = time.time()
timeout_sec = 2 * 60
tgid_b = bytes_3(9990)
bridges: dict[str, list[dict[str, Any]]] = {
"9990": [
{
"SYSTEM": "ECHO",
"TS": 2,
"TGID": tgid_b,
"ACTIVE": True,
"TIMEOUT": timeout_sec,
"TO_TYPE": "NONE",
"ON": [],
"OFF": [],
"RESET": [],
"TIMER": now + timeout_sec,
}
]
}
systems_cfg = config.get("SYSTEMS", {})
for _system, sys_cfg in systems_cfg.items():
if _system == "ECHO":
continue
if sys_cfg.get("MODE") != "MASTER":
continue
_tmout = float(sys_cfg.get("DEFAULT_UA_TIMER", 10))
bridges["9990"].append(
{
"SYSTEM": _system,
"TS": 1,
"TGID": tgid_b,
"ACTIVE": False,
"TIMEOUT": _tmout * 60,
"TO_TYPE": "ON",
"OFF": [],
"ON": [tgid_b],
"RESET": [],
"TIMER": now,
}
)
bridges["9990"].append(
{
"SYSTEM": _system,
"TS": 2,
"TGID": tgid_b,
"ACTIVE": False,
"TIMEOUT": _tmout * 60,
"TO_TYPE": "ON",
"OFF": [],
"ON": [tgid_b],
"RESET": [],
"TIMER": now,
}
)
return bridges

@ -31,9 +31,6 @@ class SubscriptionRouter:
def __init__(self, store: SubscriptionStore) -> None:
self._store = store
self._indexed = hasattr(store, "relay_tables_with_active_source") and hasattr(
store, "legs_in_table"
)
def resolve(self, ingress: VoiceIngress) -> tuple[ForwardLeg, ...]:
"""Return active forward legs when the source has an ACTIVE row on the dst TG (legacy to_target)."""
@ -50,16 +47,7 @@ class SubscriptionRouter:
legs: list[ForwardLeg] = []
seen_obp: set[tuple[str, int]] = set()
for table_key in tables:
subs = (
self._store.legs_in_table(table_key)
if self._indexed
else tuple(
sub
for sub in self._store.snapshot()
if sub.table_key() == table_key
)
)
for sub in subs:
for sub in self._store.legs_in_table(table_key):
if sub.system.value == ingress.source_system:
continue
if not sub.is_active():
@ -80,21 +68,4 @@ class SubscriptionRouter:
def relay_tables_with_active_source(self, system: str, slot: int, dst_tgid: int) -> tuple[str, ...]:
"""Mirror legacy ``relay_tables_with_active_source`` on subscription rows."""
if self._indexed:
return self._store.relay_tables_with_active_source(system, slot, dst_tgid)
tables: list[str] = []
seen: set[str] = set()
for sub in self._store.snapshot():
if sub.system.value != system:
continue
if int(sub.channel.slot) != int(slot):
continue
if not sub.is_active():
continue
if int(sub.target_tgid) != int(dst_tgid):
continue
key = sub.table_key()
if key not in seen:
seen.add(key)
tables.append(key)
return tuple(sorted(tables))
return self._store.relay_tables_with_active_source(system, slot, dst_tgid)

@ -22,7 +22,7 @@
Runtime hot paths mutate the store directly and publish via ``export_routing_table``; they must not
call ``replace_store_from_routing_table`` (would overwrite store authority with the shim).
Bootstrap and tests may import once — e.g. ``_seed_echo_routing_table`` in ``peer_server.py``.
Bootstrap and tests may import once — e.g. ``seed_echo_routing_table`` in ``application/subscription/echo_seed.py``.
"""
from __future__ import annotations

@ -37,19 +37,7 @@ def system_has_active_leg_in_store(
tgid: int,
) -> bool:
"""True when an ACTIVE subscription leg exists for ``(system, slot, tgid)``."""
indexed = getattr(store, "has_active_target_leg", None)
if callable(indexed):
return bool(indexed(system, slot, tgid))
for sub in store.snapshot():
if sub.system.value != system:
continue
if int(sub.channel.slot) != int(slot):
continue
if int(sub.target_tgid) != int(tgid):
continue
if sub.is_active():
return True
return False
return store.has_active_target_leg(system, slot, tgid)
def active_system_slots_for_tg_in_store(

@ -27,7 +27,6 @@
from __future__ import annotations
import logging
import os
import time
from datetime import datetime
from typing import Any, Callable
@ -72,7 +71,6 @@ class VoiceUseCases:
self._tts_running: dict[int, bool] = {}
self._announcement_last_hour: dict[int, int] = {}
self._tts_last_hour: dict[int, int] = {}
self._config_file_mtime: float = 0.0
self._broadcast_queue: list[dict[str, Any]] = []
self._broadcast_active_tgs: set[str] = set()
@ -249,8 +247,8 @@ class VoiceUseCases:
for sid in list(obj.STATUS.keys()):
if sid not in (1, 2):
del obj.STATUS[sid]
except Exception:
pass
except Exception as e:
logger.warning("(%s) slot STATUS cleanup failed: %s", label, e)
self._announcement_running[ann_idx] = False
if not targets:
logger.info(
@ -523,7 +521,8 @@ class VoiceUseCases:
self._tts_running[tts_idx] = False
try:
msg = failure.getErrorMessage()
except Exception:
except Exception as e:
logger.warning("(%s) failure.getErrorMessage unavailable: %s", label, e)
msg = str(failure)
logger.error("(%s) TTS conversion error: %s", label, msg)
@ -550,8 +549,8 @@ class VoiceUseCases:
for sid in list(obj.STATUS.keys()):
if sid not in (1, 2):
del obj.STATUS[sid]
except Exception:
pass
except Exception as e:
logger.warning("(%s) slot STATUS cleanup failed: %s", label, e)
self._tts_running[tts_idx] = False
if not targets:
logger.info(
@ -718,17 +717,8 @@ class VoiceUseCases:
self._call_from_reactor(protocol.send_voice_packet, pkt, bytes_3(5000), bytes_3(9), _slot)
logger.debug("(%s) disconnected voice thread end", system)
def check_voice_config_reload(self, config_file_path: str | None = None) -> None:
"""Check adn-voice.yaml mtime and reload if changed (15s loop). Start/stop announcement LoopingCalls."""
if config_file_path and os.path.isfile(config_file_path):
try:
mtime = os.path.getmtime(config_file_path)
except OSError:
return
if mtime == self._config_file_mtime:
return
self._config_file_mtime = mtime
logger.info("(VOICE-RELOAD) config file change detected, reloading configuration...")
def apply_voice_config(self) -> None:
"""Start/stop announcement and TTS LoopingCalls from ``config["VOICE"]``."""
g = self._config.get("VOICE", {})
if not self._start_looping_call:
return
@ -740,8 +730,8 @@ class VoiceUseCases:
try:
if getattr(self._ann_tasks[ann_idx], "running", False):
self._ann_tasks[ann_idx].stop()
except Exception:
pass
except Exception as e:
logger.warning("(VOICE-RELOAD) stop announcement task %s failed: %s", ann_idx + 1, e)
del self._ann_tasks[ann_idx]
logger.info("(VOICE-RELOAD) ANNOUNCEMENT-%s stopped", ann_idx + 1)
for ann_idx, item in enumerate(announcements):
@ -752,8 +742,8 @@ class VoiceUseCases:
try:
if getattr(self._ann_tasks[ann_idx], "running", False):
self._ann_tasks[ann_idx].stop()
except Exception:
pass
except Exception as e:
logger.warning("(VOICE-RELOAD) stop %s failed: %s", label, e)
del self._ann_tasks[ann_idx]
logger.info("(VOICE-RELOAD) %s stopped", label)
mode = item.get("MODE", "interval")
@ -772,8 +762,8 @@ class VoiceUseCases:
try:
if getattr(self._tts_tasks[tts_idx], "running", False):
self._tts_tasks[tts_idx].stop()
except Exception:
pass
except Exception as e:
logger.warning("(VOICE-RELOAD) stop TTS task %s failed: %s", tts_idx + 1, e)
del self._tts_tasks[tts_idx]
logger.info("(VOICE-RELOAD) TTS-%s stopped", tts_idx + 1)
for tts_idx, item in enumerate(tts_list):
@ -784,8 +774,8 @@ class VoiceUseCases:
try:
if getattr(self._tts_tasks[tts_idx], "running", False):
self._tts_tasks[tts_idx].stop()
except Exception:
pass
except Exception as e:
logger.warning("(VOICE-RELOAD) stop %s failed: %s", label, e)
del self._tts_tasks[tts_idx]
logger.info("(VOICE-RELOAD) %s stopped", label)
mode = item.get("MODE", "interval")
@ -796,5 +786,3 @@ class VoiceUseCases:
"(VOICE-RELOAD) %s enabled - mode: %s, file: %s, TG: %s",
label, mode, item.get("FILE"), item.get("TG"),
)
if config_file_path and self._config_file_mtime:
logger.info("(VOICE-RELOAD) config reload completed")

@ -44,7 +44,6 @@ from adn_server.application.proxy.deployment import (
proxy_target_system,
)
from adn_server.application.report.queue import BoundedReportQueue, QueuedReportSender
from adn_server.infrastructure.udp_rcvbuf import apply_udp_rcvbuf, udp_rcvbuf_bytes
from adn_server.application.runtime_context import (
ConfigProxy,
RuntimeContext,
@ -52,8 +51,8 @@ from adn_server.application.runtime_context import (
prepare_reload_config,
swap_runtime_config,
)
from adn_server.application.subscription.echo_seed import seed_echo_routing_table
from adn_server.application.subscription.store_sync import replace_store_from_routing_table
from adn_server.domain import bytes_3
from adn_server.domain.dmr.bptc import encode_emblc
from adn_server.infrastructure.acl_router import InMemoryAclRouter
from adn_server.infrastructure.config_normalizer import (
@ -92,6 +91,7 @@ from adn_server.infrastructure.twisted_adapters.report.mqtt_publisher import (
from adn_server.infrastructure.twisted_adapters.report.worker import start_report_queue_worker
from adn_server.infrastructure.twisted_adapters.report_server import ReportServerFactory
from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocolFactory
from adn_server.infrastructure.udp_rcvbuf import apply_udp_rcvbuf, udp_rcvbuf_bytes
from adn_server.infrastructure.voice import DefaultVoiceProvider, StubVoiceProvider
from adn_server.infrastructure.voice.recording import RecordingHandler
@ -163,39 +163,6 @@ def _wire_monitor_downlink_ctx(
report_factory.set_downlink_ctx_for_system(_ctx_for)
def _seed_echo_routing_table(config: dict) -> dict:
"""Initial BRIDGES for ECHO system (legacy make_bridges 9990 + MASTER expansion)."""
now = time.time()
timeout_sec = 2 * 60
tgid_b = bytes_3(9990)
bridges: dict = {
"9990": [
{
"SYSTEM": "ECHO",
"TS": 2,
"TGID": tgid_b,
"ACTIVE": True,
"TIMEOUT": timeout_sec,
"TO_TYPE": "NONE",
"ON": [],
"OFF": [],
"RESET": [],
"TIMER": now + timeout_sec,
}
]
}
systems_cfg = config.get("SYSTEMS", {})
for _system, sys_cfg in systems_cfg.items():
if _system == "ECHO":
continue
if sys_cfg.get("MODE") != "MASTER":
continue
_tmout = float(sys_cfg.get("DEFAULT_UA_TIMER", 10))
bridges["9990"].append({"SYSTEM": _system, "TS": 1, "TGID": tgid_b, "ACTIVE": False, "TIMEOUT": _tmout * 60, "TO_TYPE": "ON", "OFF": [], "ON": [tgid_b], "RESET": [], "TIMER": now})
bridges["9990"].append({"SYSTEM": _system, "TS": 2, "TGID": tgid_b, "ACTIVE": False, "TIMEOUT": _tmout * 60, "TO_TYPE": "ON", "OFF": [], "ON": [tgid_b], "RESET": [], "TIMER": now})
return bridges
def _looping_errback(logger: logging.Logger, failure):
"""Errback for LoopingCalls (legacy loopingErrHandle). Stops reactor to avoid memory leaks."""
try:
@ -279,7 +246,7 @@ def run_peer_server(
systems_cfg = config.get("SYSTEMS", {})
subscription_store = InMemorySubscriptionStore()
if "ECHO" in systems_cfg and systems_cfg["ECHO"].get("MODE") in ("PEER", "MASTER"):
replace_store_from_routing_table(subscription_store, _seed_echo_routing_table(config))
replace_store_from_routing_table(subscription_store, seed_echo_routing_table(config))
_ctx = runtime_holder.get()
runtime_holder.swap(
RuntimeContext(
@ -376,7 +343,7 @@ def run_peer_server(
start_looping_call=_start_voice_loop,
defer_to_thread=threads.deferToThread,
)
voice_use_cases.check_voice_config_reload(voice_config_path)
voice_use_cases.apply_voice_config()
ident_use_cases = IdentUseCases(
config,
voice_use_cases,
@ -670,10 +637,29 @@ def run_peer_server(
reactor.addSystemEventTrigger("before", "shutdown", shutdown_handler)
# Voice config reload (15s): re-read adn-voice.yaml, start/stop announcement LoopingCalls on change
voice_config_mtime = 0.0
if voice_config_path and os.path.isfile(voice_config_path):
try:
voice_config_mtime = os.path.getmtime(voice_config_path)
except OSError:
voice_config_mtime = 0.0
def voice_reload_loop():
nonlocal voice_config_mtime
try:
if voice_config_path and os.path.isfile(voice_config_path):
try:
mtime = os.path.getmtime(voice_config_path)
except OSError:
return
if mtime == voice_config_mtime:
return
voice_config_mtime = mtime
logger.info("(VOICE-RELOAD) config file change detected, reloading configuration...")
loader.reload_voice_config(config, voice_config_path)
voice_use_cases.check_voice_config_reload(voice_config_path)
voice_use_cases.apply_voice_config()
if voice_config_path and voice_config_mtime:
logger.info("(VOICE-RELOAD) config reload completed")
except Exception as e:
logger.debug("(VOICE-RELOAD) %s", e)

@ -100,7 +100,11 @@ def run_echo(config: dict[str, Any], *, logger: logging.Logger) -> None:
logger.debug("(ECHO) skip %s (MODE=%s; echo only runs PEER systems)", system_name, sys_cfg.get("MODE"))
continue
pb = PlaybackUseCases(system_name, get_protocol=lambda sn=system_name: protocols.get(sn))
pb = PlaybackUseCases(
system_name,
call_later=reactor.callLater,
get_protocol=lambda sn=system_name: protocols.get(sn),
)
protocol = HBPProtocolFactory(
system_name,

@ -975,6 +975,25 @@ class HBPProtocol(DatagramProtocol):
return _peer_ids[_int_peer_id]
return False
def validate_obp_source_server_id(self, peer_id: bytes):
"""OPENBRIDGE DMRE source-server lookup (legacy hblink.py OPENBRIDGE.validate_id).
Unlike ``validate_id`` on HBP systems, OBP ingress always resolves 6–7 digit
source servers against alias tables with no ``ALLOW_UNREG_ID`` bypass.
Returns a callsign string on match, or ``False``.
"""
_int_peer_id = int(int_id(peer_id))
_local_subscriber_ids = self._CONFIG.get("_LOCAL_SUBSCRIBER_IDS", {})
_subscriber_ids = self._CONFIG.get("_SUB_IDS", {})
_peer_ids = self._CONFIG.get("_PEER_IDS", {})
if _int_peer_id in _local_subscriber_ids:
return _local_subscriber_ids[_int_peer_id]
if _int_peer_id in _subscriber_ids:
return _subscriber_ids[_int_peer_id]
if _int_peer_id in _peer_ids:
return _peer_ids[_int_peer_id]
return False
def proxy_IPBlackList(self, peer_id: bytes, sockaddr: tuple[str, int]) -> None:
"""Legacy hblink.py proxy_IPBlackList: send PRBL to proxy to blacklist a peer's IP for 5 min."""
_bltime = str(time.time() + 300)
@ -1534,8 +1553,13 @@ class HBPProtocol(DatagramProtocol):
if self._on_options_received:
try:
self._on_options_received(self._system, _this_peer["OPTIONS"])
except Exception:
pass
except Exception as e:
logger.warning(
"(%s) on_options_received failed for peer %s: %s",
self._system,
int_id(_peer_id),
e,
)
if is_proxy_inject_only(self._CONFIG, self._system):
for _slot in (1, 2):
if _slot in self.STATUS:
@ -2131,7 +2155,7 @@ class HBPProtocol(DatagramProtocol):
self._obp_send_bcsq(_dst_id, _stream_id)
self._laststrid.append(_stream_id)
return
if _src_srv_len > 5 and not self.validate_id(_source_server):
if _src_srv_len > 5 and not self.validate_obp_source_server_id(_source_server):
if _stream_id not in self._laststrid:
logger.warning("(%s) Source Server 6 or 7 digits but not a valid DMR ID, discarding Src: %s", self._system, _src_srv_int)
self._obp_send_bcsq(_dst_id, _stream_id)

@ -250,8 +250,8 @@ def _cleanup(files: list[str]) -> None:
try:
if os.path.isfile(f):
os.remove(f)
except Exception:
pass
except Exception as e:
logger.warning("(TTS) cleanup remove %s failed: %s", f, e)
def text_to_ambe(

@ -27,6 +27,7 @@ Use the project interpreter, e.g. `/opt/.pyenv/versions/3.11.8/bin/python3`.
| `replay/` | JSONL session replay |
| `schemas/` | Report v2 JSON Schema validation (`jsonschema` dev dep) |
| `application/` | Report payloads, monitor topology, proxy use cases, subscription store/router |
| `fakes/` | Re-exports for application tests (`InMemorySubscriptionStore` shim) — not run as tests |
| `infrastructure/` | Logging reload, ACL router, **HBP REPEAT + proxy fan-in integration**, MQTT |
| `smoke/` | Quick routing smoke |
| `support/` | Shared stacks (`hbp_repeat_stack`, monitor sim) — not run as tests |

@ -22,7 +22,23 @@
from __future__ import annotations
import json
from pathlib import Path
import jsonschema
import pytest
from adn_server.application.report.dashboard_state import build_dashboard_state
from adn_server.domain.value_objects import bytes_4
_SCHEMA_PATH = Path(__file__).resolve().parents[2] / "schemas" / "report-v2.json"
_EXAMPLES_DIR = _SCHEMA_PATH.parent / "examples"
@pytest.fixture(scope="module")
def validator() -> jsonschema.Draft202012Validator:
with _SCHEMA_PATH.open(encoding="utf-8") as fh:
return jsonschema.Draft202012Validator(json.load(fh))
def test_dashboard_state_omits_idle_masters():
@ -105,3 +121,56 @@ def test_dashboard_state_includes_connected_upstream_peer():
assert "XLX-730" in state["ctable"]["PEERS"]
assert state["ctable"]["PEERS"]["XLX-730"]["connected"] is True
assert state["ctable"]["MASTERS"] == {}
def test_build_dashboard_state_matches_example(validator: jsonschema.Draft202012Validator) -> None:
with (_EXAMPLES_DIR / "dashboard_state.json").open(encoding="utf-8") as fh:
expected = json.load(fh)
systems = {
"MASTER-A": {
"MODE": "MASTER",
"ENABLED": True,
"IP": "10.0.0.1",
"PORT": 62030,
"SINGLE_MODE": False,
"DEFAULT_UA_TIMER": 10,
"TS1_STATIC": "91,92",
"TS2_STATIC": "7302",
"PEERS": {
bytes_4(3120001): {
"CONNECTION": "YES",
"CONNECTED": 1717555200,
"IP": "10.0.0.50",
"PORT": 62031,
"CALLSIGN": b"CE5RPY ",
},
},
},
"MASTER-B": {
"MODE": "MASTER",
"ENABLED": True,
"IP": "10.0.0.2",
"PORT": 62032,
"PEERS": {
bytes_4(3120002): {
"CONNECTION": "YES",
"CONNECTED": 1717555200,
"IP": "10.0.0.51",
"PORT": 62033,
"CALLSIGN": b"EA5GVK ",
},
},
},
"OBP-CL": {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"IP": "44.31.61.68",
"PORT": 62999,
"NETWORK_ID": 73010,
"ENHANCED_OBP": True,
"PEERS": {},
},
}
doc = build_dashboard_state(systems, server_id="7302", ts=expected["ts"])
assert json.loads(json.dumps(doc)) == expected
validator.validate(doc)

@ -32,7 +32,7 @@ from adn_server.application.subscription.store_sync import replace_store_from_ro
from adn_server.application.subscription.subscription_table_ops import make_static_tg_store
from adn_server.domain import bytes_3
from adn_server.domain.subscription import SubscriptionPhase
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def _store_from_bridges(bridges: dict) -> InMemorySubscriptionStore:

@ -37,6 +37,7 @@ from adn_server.application.report import (
)
from adn_server.application.report.payloads import (
parse_peer_options_static,
peer_options_static_valid,
resolve_peer_single_and_timer,
)
from adn_server.application.routing.helpers import peer_should_receive_group_voice
@ -59,6 +60,20 @@ def test_parse_peer_options_static_ts2():
assert ts2 == ["730444"]
@pytest.mark.parametrize(
("options", "expected"),
[
(b"", True),
(b"TS2=730444;TIMER=15;", True),
(b"TS2=bad;TIMER=15;", False),
(b"PASS=secret;TS2=730;", False),
(b"PASS=secret;", False),
],
)
def test_peer_options_static_valid_table(options: bytes, expected: bool) -> None:
assert peer_options_static_valid(options) is expected
def test_resolve_peer_single_and_timer_yaml_defaults() -> None:
yaml_cfg = {"SINGLE_MODE": False, "DEFAULT_UA_TIMER": 60}
single, timer = resolve_peer_single_and_timer({}, yaml_cfg)

@ -34,7 +34,7 @@ from adn_server.domain.subscription import (
SystemId,
TgId,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def test_subscription_to_legacy_row_echo():

@ -34,7 +34,7 @@ from adn_server.domain.subscription import (
SystemId,
TgId,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def _sample_store() -> InMemorySubscriptionStore:

@ -27,7 +27,7 @@ from tests.harness.deterministic import active_routing_table
from adn_server.application.subscription.rule_timer_ops import apply_rule_timer_store
from adn_server.application.subscription.store_sync import replace_store_from_routing_table
from adn_server.domain.subscription import SubscriptionPhase
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def test_apply_rule_timer_store_deactivates_expired_on() -> None:

@ -1232,7 +1232,7 @@ def test_idle_static_bridges_do_not_block_options_fanout() -> None:
from adn_server.application.subscription.routing_table_import import (
subscriptions_from_routing_table,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def _row(*, system: str, ts: int, tgid: int) -> dict:
return {
@ -1343,7 +1343,7 @@ def test_single0_ingress_vterm_bridge_hold_blocks_within_hangtime() -> None:
from adn_server.application.subscription.routing_table_import import (
subscriptions_from_routing_table,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def _row(*, system: str, ts: int, tgid: int) -> dict:
return {
@ -1416,7 +1416,7 @@ def test_single0_listener_bridge_hold_expires_after_hangtime_allows_fresh_ptt()
from adn_server.application.subscription.routing_table_import import (
subscriptions_from_routing_table,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def _row(*, system: str, ts: int, tgid: int) -> dict:
return {
@ -1492,7 +1492,7 @@ def test_single0_listener_fresh_ptt_after_blocked_second_stream() -> None:
from adn_server.application.subscription.routing_table_import import (
subscriptions_from_routing_table,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def _row(*, system: str, ts: int, tgid: int) -> dict:
return {
@ -1646,7 +1646,7 @@ def test_single1_duplex_listen_lock_does_not_block_other_rf_slot() -> None:
from adn_server.application.subscription.routing_table_import import (
subscriptions_from_routing_table,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def _row(*, system: str, ts: int, tgid: int) -> dict:
return {

@ -26,7 +26,7 @@ from adn_server.application.subscription.stat_trimmer_ops import apply_stat_trim
from adn_server.application.subscription.subscription_queries import store_has_table
from adn_server.application.subscription.subscription_table_ops import ensure_stat_relay_store
from adn_server.domain import bytes_3
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def test_stat_trimmer_keeps_obp_stat_bridge_with_inactive_system_legs() -> None:

@ -24,11 +24,11 @@ from __future__ import annotations
from typing import Any
from adn_server.application.subscription.echo_seed import seed_echo_routing_table
from adn_server.application.subscription.routing_table_export import export_routing_table
from adn_server.application.subscription.store_sync import replace_store_from_routing_table
from adn_server.domain import bytes_3, int_id
from adn_server.infrastructure.bootstrap.peer_server import _seed_echo_routing_table
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def _row_fingerprint(row: dict[str, Any]) -> tuple[Any, ...]:
@ -59,7 +59,7 @@ def test_replace_store_from_echo_bridges_round_trip():
"MASTER-B": {"MODE": "MASTER", "DEFAULT_UA_TIMER": 15},
}
}
bridges = _seed_echo_routing_table(config)
bridges = seed_echo_routing_table(config)
store = InMemorySubscriptionStore()
replace_store_from_routing_table(store, bridges)

@ -29,7 +29,7 @@ from adn_server.application.subscription.routing_table_import import subscriptio
from adn_server.domain import bytes_3, int_id
from adn_server.domain.subscription import TgId
from adn_server.domain.voice_routing import ForwardLeg, VoiceIngress
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def _legacy_relay_tables(

@ -31,7 +31,7 @@ from adn_server.application.subscription.subscription_table_ops import (
make_static_tg_store,
)
from adn_server.domain import bytes_3
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def _store(bridges: dict | None = None) -> InMemorySubscriptionStore:

@ -33,7 +33,7 @@ from adn_server.application.subscription.subscription_reset_ops import (
)
from adn_server.domain import bytes_3
from adn_server.domain.subscription import SubscriptionPhase
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def _store(bridges: dict) -> InMemorySubscriptionStore:

@ -35,7 +35,7 @@ from adn_server.domain.subscription import (
SystemId,
TgId,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
def test_build_voice_ingress_hbp_group():

@ -32,7 +32,7 @@ from adn_server.domain.dmr.bptc import encode_emblc
from adn_server.domain.subscription import TgId
from adn_server.domain.voice_routing import ForwardLeg
from adn_server.infrastructure.acl_router import InMemoryAclRouter
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from fakes.subscription_store import InMemorySubscriptionStore
from adn_server.infrastructure.talker_alias_emblc import default_ta_emblc_encoder

@ -26,7 +26,12 @@ import logging
from unittest.mock import patch
from tests.harness.deterministic import DeterministicScenario, PacketSpec
from tests.harness.playback_helpers import FakePlaybackProtocol, install_reactor_capture, send_playback
from tests.harness.playback_helpers import (
FakePlaybackProtocol,
make_capture_call_later,
noop_call_later,
send_playback,
)
from adn_server.application.playback_use_cases import _PLAYBACK_DELAY_S, _RECORD_IDLE_S, PlaybackUseCases
@ -34,7 +39,7 @@ from adn_server.application.playback_use_cases import _PLAYBACK_DELAY_S, _RECORD
def test_dmrd_received_accepts_ingress_pkt_time_kwarg() -> None:
"""Regression: PEER udp_hbp passes ingress_pkt_time (echo must not TypeError)."""
proto = FakePlaybackProtocol()
pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto)
pb = PlaybackUseCases("ECHO", call_later=noop_call_later, get_protocol=lambda: proto)
base = PacketSpec(dst_id=9990, stream_id=0x12121212, slot=2)
args = DeterministicScenario.voice_head_spec(base).decoded_hbp_args()
@ -60,23 +65,22 @@ def test_dmrd_received_accepts_ingress_pkt_time_kwarg() -> None:
def test_ingress_pkt_time_enables_record_to_playback(caplog) -> None:
"""Regression: PEER path (ingress_pkt_time) records voice and schedules playback."""
proto = FakePlaybackProtocol()
pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto)
call_later, scheduled = make_capture_call_later()
pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto)
base = PacketSpec(dst_id=9990, stream_id=0x34343434, slot=2)
mock_reactor, scheduled = install_reactor_capture()
t0 = 1_700_000_000.0
with patch("adn_server.application.playback_use_cases.reactor", mock_reactor):
with patch("adn_server.application.playback_use_cases.time", return_value=t0):
send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base))
send_playback(
pb,
"ECHO",
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1),
ingress_pkt_time=t0 + 2.0,
)
with patch("adn_server.application.playback_use_cases.time", return_value=t0 + _RECORD_IDLE_S + 1):
with caplog.at_level(logging.INFO, logger="adn_server.application.playback_use_cases"):
pb._on_record_idle(proto)
with patch("adn_server.application.playback_use_cases.time", return_value=t0):
send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base))
send_playback(
pb,
"ECHO",
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1),
ingress_pkt_time=t0 + 2.0,
)
with patch("adn_server.application.playback_use_cases.time", return_value=t0 + _RECORD_IDLE_S + 1):
with caplog.at_level(logging.INFO, logger="adn_server.application.playback_use_cases"):
pb._on_record_idle(proto)
assert any(item[0] == _PLAYBACK_DELAY_S for item in scheduled)
assert pb._playback_busy is True

@ -26,7 +26,7 @@ import logging
import re
from tests.harness.deterministic import DeterministicScenario, PacketSpec
from tests.harness.playback_helpers import FakePlaybackProtocol
from tests.harness.playback_helpers import FakePlaybackProtocol, noop_call_later
from adn_server.application.playback_use_cases import PlaybackUseCases
from adn_server.domain import bytes_3, bytes_4
@ -35,7 +35,7 @@ from adn_server.domain import bytes_3, bytes_4
def test_start_playback_logs_duration_with_two_decimals(caplog) -> None:
"""PLAYBACK duration matches bridge-style %.2f (no float noise in logs)."""
proto = FakePlaybackProtocol()
pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto)
pb = PlaybackUseCases("ECHO", call_later=noop_call_later, get_protocol=lambda: proto)
base = PacketSpec(dst_id=9990, stream_id=0x88888888, slot=2)
recorded = [DeterministicScenario.voice_head_spec(base).data()]

@ -25,7 +25,7 @@ from __future__ import annotations
from unittest.mock import patch
from tests.harness.deterministic import DeterministicScenario, PacketSpec
from tests.harness.playback_helpers import FakePlaybackProtocol, install_reactor_capture, run_scheduled
from tests.harness.playback_helpers import FakePlaybackProtocol, make_capture_call_later, run_scheduled
from adn_server.application.playback_use_cases import (
_PACKET_INTERVAL_S,
@ -38,17 +38,16 @@ from adn_server.domain import bytes_3, bytes_4
def test_send_next_packet_emits_all_packets_then_finishes() -> None:
proto = FakePlaybackProtocol()
pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto)
call_later, scheduled = make_capture_call_later()
pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto)
pb._playback_busy = True
pb._playback_stream_id = bytes_4(0x77777777)
pb._playback_packets = [bytes([i]) * 55 for i in range(1, 5)]
pb._playback_index = 0
mock_reactor, scheduled = install_reactor_capture()
with patch("adn_server.application.playback_use_cases.reactor", mock_reactor):
pb._send_next_packet(proto)
while scheduled:
run_scheduled(scheduled)
pb._send_next_packet(proto)
while scheduled:
run_scheduled(scheduled)
assert len(proto.sent) == 4
assert pb._playback_busy is False
@ -58,31 +57,29 @@ def test_send_next_packet_emits_all_packets_then_finishes() -> None:
def test_start_playback_sends_first_packet_and_schedules_rest() -> None:
proto = FakePlaybackProtocol()
pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto)
call_later, scheduled = make_capture_call_later()
pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto)
base = PacketSpec(dst_id=9990, stream_id=0x88888888, slot=2)
recorded = [
DeterministicScenario.voice_head_spec(base).data(),
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1).data(),
]
mock_reactor, scheduled = install_reactor_capture()
with patch("adn_server.application.playback_use_cases.reactor", mock_reactor):
pb._start_playback(
proto,
recorded,
bytes_3(base.rf_src),
bytes_4(base.peer_id),
bytes_3(base.dst_id),
2,
1.5,
)
pb._start_playback(
proto,
recorded,
bytes_3(base.rf_src),
bytes_4(base.peer_id),
bytes_3(base.dst_id),
2,
1.5,
)
assert len(proto.sent) == 1
assert pb._playback_index == 1
assert any(item[0] == _PACKET_INTERVAL_S for item in scheduled)
with patch("adn_server.application.playback_use_cases.reactor", mock_reactor):
while scheduled:
run_scheduled(scheduled)
while scheduled:
run_scheduled(scheduled)
assert len(proto.sent) == 2
assert pb._playback_busy is False
@ -90,7 +87,8 @@ def test_start_playback_sends_first_packet_and_schedules_rest() -> None:
def test_max_duration_commits_recording_when_no_vterm() -> None:
proto = FakePlaybackProtocol()
pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto)
call_later, scheduled = make_capture_call_later()
pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto)
base = PacketSpec(dst_id=9990, stream_id=0x99999999, slot=2)
pb._recording_active = True
pb.CALL_DATA = [DeterministicScenario.voice_head_spec(base).data()]
@ -102,11 +100,9 @@ def test_max_duration_commits_recording_when_no_vterm() -> None:
"peer_id": bytes_4(base.peer_id),
"dst_id": bytes_3(base.dst_id),
}
mock_reactor, scheduled = install_reactor_capture()
with patch("adn_server.application.playback_use_cases.reactor", mock_reactor):
with patch("adn_server.application.playback_use_cases.time", return_value=100.0 + _SOURCE_MAX_S):
pb._on_record_max_duration(proto)
with patch("adn_server.application.playback_use_cases.time", return_value=100.0 + _SOURCE_MAX_S):
pb._on_record_max_duration(proto)
assert not pb._recording_active
assert pb._playback_busy is True
@ -115,14 +111,13 @@ def test_max_duration_commits_recording_when_no_vterm() -> None:
def test_packet_interval_matches_expected() -> None:
proto = FakePlaybackProtocol()
pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto)
call_later, scheduled = make_capture_call_later()
pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto)
pb._playback_busy = True
pb._playback_stream_id = bytes_4(0xAAAAAAAA)
pb._playback_packets = [b"\x01" * 55, b"\x02" * 55]
pb._playback_index = 0
mock_reactor, scheduled = install_reactor_capture()
with patch("adn_server.application.playback_use_cases.reactor", mock_reactor):
pb._send_next_packet(proto)
pb._send_next_packet(proto)
assert scheduled[0][0] == _PACKET_INTERVAL_S

@ -27,7 +27,8 @@ from unittest.mock import patch
from tests.harness.deterministic import DeterministicScenario, PacketSpec
from tests.harness.playback_helpers import (
FakePlaybackProtocol,
install_reactor_capture,
make_capture_call_later,
noop_call_later,
run_scheduled,
send_playback,
)
@ -37,22 +38,21 @@ from adn_server.application.playback_use_cases import _RECORD_IDLE_S, PlaybackUs
def test_idle_timeout_appends_synthetic_vterm_and_schedules_playback() -> None:
proto = FakePlaybackProtocol()
pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto)
call_later, scheduled = make_capture_call_later()
pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto)
base = PacketSpec(dst_id=9990, stream_id=0x55555555, slot=2)
mock_reactor, scheduled = install_reactor_capture()
with patch("adn_server.application.playback_use_cases.reactor", mock_reactor):
with patch("adn_server.application.playback_use_cases.time", return_value=200.0):
send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base))
send_playback(
pb, "ECHO", DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1),
)
with patch("adn_server.application.playback_use_cases.time", return_value=200.0):
send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base))
send_playback(
pb, "ECHO", DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1),
)
idle_calls = [item for item in scheduled if item[0] == _RECORD_IDLE_S]
assert idle_calls
idle_calls = [item for item in scheduled if item[0] == _RECORD_IDLE_S]
assert idle_calls
with patch("adn_server.application.playback_use_cases.time", return_value=200.0 + _RECORD_IDLE_S):
run_scheduled(scheduled, delay=_RECORD_IDLE_S)
with patch("adn_server.application.playback_use_cases.time", return_value=200.0 + _RECORD_IDLE_S):
run_scheduled(scheduled, delay=_RECORD_IDLE_S)
assert not pb._recording_active
assert pb._playback_busy is True
@ -60,7 +60,7 @@ def test_idle_timeout_appends_synthetic_vterm_and_schedules_playback() -> None:
def test_vterm_commit_does_not_require_synthetic_vterm() -> None:
pb = PlaybackUseCases("ECHO")
pb = PlaybackUseCases("ECHO", call_later=noop_call_later)
base = PacketSpec(dst_id=9990, stream_id=0x66666666, slot=2)
recorded = [
DeterministicScenario.voice_head_spec(base).data(),

@ -26,7 +26,12 @@ from unittest.mock import patch
import pytest
from tests.harness.deterministic import DeterministicScenario, PacketSpec
from tests.harness.playback_helpers import FakePlaybackProtocol, install_reactor_capture, send_playback
from tests.harness.playback_helpers import (
FakePlaybackProtocol,
make_capture_call_later,
noop_call_later,
send_playback,
)
from adn_server.application.playback_use_cases import PlaybackUseCases
from adn_server.domain import bytes_4
@ -54,7 +59,7 @@ def _long_voice_recording(
@pytest.mark.behavior
def test_prepare_playback_preserves_source_seq_past_255_packets() -> None:
"""Regression: long QSO keeps source seq; new stream segment when seq byte wraps."""
pb = PlaybackUseCases("ECHO")
pb = PlaybackUseCases("ECHO", call_later=noop_call_later)
base = PacketSpec(dst_id=9990, stream_id=0x55555555, slot=2)
recorded = _long_voice_recording(base, burst_count=500)
pb._playback_stream_id = bytes_4(0x77777777)
@ -83,45 +88,47 @@ def test_prepare_playback_preserves_source_seq_past_255_packets() -> None:
def test_start_playback_sends_preserved_seq_over_30s_recording() -> None:
"""End-to-end: ~500 bursts @ 60ms ≈ 30s; replay seq on wire must match recording."""
proto = FakePlaybackProtocol()
pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto)
call_later, scheduled = make_capture_call_later()
pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto)
base = PacketSpec(dst_id=9990, stream_id=0x66666666, slot=2)
burst_count = 500
recorded = _long_voice_recording(base, burst_count=burst_count)
mock_reactor, scheduled = install_reactor_capture()
pb._playback_stream_id = bytes_4(0x77777777)
out = pb._prepare_playback_packets(recorded)
assert len(out) == len(recorded) + 1
with patch("adn_server.application.playback_use_cases.reactor", mock_reactor):
pb._playback_packets = out
pb._playback_index = 0
pb._playback_busy = True
pb._send_next_packet(proto)
while scheduled:
_delay, fn, args = scheduled.pop(0)
fn(*args)
pb._playback_packets = out
pb._playback_index = 0
pb._playback_busy = True
pb._send_next_packet(proto)
while scheduled:
_delay, fn, args = scheduled.pop(0)
fn(*args)
assert len(proto.sent) == len(out)
def test_recording_rekey_does_not_store_second_vhead() -> None:
pb = PlaybackUseCases("ECHO", get_protocol=lambda: FakePlaybackProtocol())
call_later, _scheduled = make_capture_call_later()
pb = PlaybackUseCases(
"ECHO",
call_later=call_later,
get_protocol=lambda: FakePlaybackProtocol(),
)
base = PacketSpec(dst_id=9990, stream_id=0x11111111, slot=2)
rekey = PacketSpec(dst_id=9990, stream_id=0x22222222, slot=2)
mock_reactor, _scheduled = install_reactor_capture()
with patch("adn_server.application.playback_use_cases.reactor", mock_reactor):
with patch("adn_server.application.playback_use_cases.time") as mock_time:
mock_time.side_effect = [100.0, 100.1, 100.2, 100.3]
send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base))
send_playback(
pb, "ECHO", DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1),
)
send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(rekey))
send_playback(
pb, "ECHO", DeterministicScenario.voice_burst_spec(rekey, seq=2, dtype_vseq=2),
)
with patch("adn_server.application.playback_use_cases.time") as mock_time:
mock_time.side_effect = [100.0, 100.1, 100.2, 100.3]
send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base))
send_playback(
pb, "ECHO", DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1),
)
send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(rekey))
send_playback(
pb, "ECHO", DeterministicScenario.voice_burst_spec(rekey, seq=2, dtype_vseq=2),
)
vheads = sum(
1
@ -133,7 +140,7 @@ def test_recording_rekey_does_not_store_second_vhead() -> None:
def test_prepare_playback_skips_mid_call_vhead_and_preserves_seq() -> None:
pb = PlaybackUseCases("ECHO")
pb = PlaybackUseCases("ECHO", call_later=noop_call_later)
base = PacketSpec(dst_id=9990, stream_id=0x33333333, slot=2)
rekey = PacketSpec(dst_id=9990, stream_id=0x44444444, slot=2)
recorded = [

@ -0,0 +1 @@
"""Test doubles — thin re-exports so application tests avoid infrastructure imports."""

@ -0,0 +1,7 @@
"""In-memory subscription store for application-layer unit tests."""
from __future__ import annotations
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
__all__ = ["InMemorySubscriptionStore"]

@ -22,6 +22,7 @@
from __future__ import annotations
from collections.abc import Callable
from typing import Any
from unittest.mock import MagicMock
@ -41,10 +42,17 @@ def packet_bytes(spec: PacketSpec) -> bytes:
return spec.data()
def install_reactor_capture() -> tuple[MagicMock, list[tuple[float, Any, tuple[Any, ...]]]]:
"""Patch reactor.callLater; returns mock and scheduled (delay, fn, args) list."""
def noop_call_later(*_args: Any, **_kwargs: Any) -> MagicMock:
"""No-op scheduler for tests that never fire timers."""
handle = MagicMock()
handle.active.return_value = False
handle.cancel = MagicMock()
return handle
def make_capture_call_later() -> tuple[Callable[..., Any], list[tuple[float, Any, tuple[Any, ...]]]]:
"""Build a call_later stub; returns (callable, scheduled (delay, fn, args) list)."""
scheduled: list[tuple[float, Any, tuple[Any, ...]]] = []
mock_reactor = MagicMock()
def call_later(delay: float, fn: Any, *args: Any) -> MagicMock:
scheduled.append((delay, fn, args))
@ -53,6 +61,13 @@ def install_reactor_capture() -> tuple[MagicMock, list[tuple[float, Any, tuple[A
handle.cancel = MagicMock()
return handle
return call_later, scheduled
def install_reactor_capture() -> tuple[MagicMock, list[tuple[float, Any, tuple[Any, ...]]]]:
"""Backward-compatible alias; prefer make_capture_call_later + call_later= kwarg."""
call_later, scheduled = make_capture_call_later()
mock_reactor = MagicMock()
mock_reactor.callLater = call_later
return mock_reactor, scheduled

@ -0,0 +1,120 @@
# ADN DMR Peer Server - tests infrastructure hbp auth handshake
#
# 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
###############################################################################
"""Full HBP login exchange on MASTER (RPTL → RPTK → RPTC)."""
from __future__ import annotations
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 MSTNAK, RPTACK, RPTC, RPTK, RPTL
from adn_server.infrastructure.twisted_adapters.udp_hbp import (
HBPProtocol,
_calc_hash,
_get_passphrase_bytes,
)
_PEER = bytes_4(1234567)
_CLIENT_ADDR = ("192.168.1.50", 62031)
_PASSPHRASE = b"test-passphrase"
class _RecordingTransport:
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))
class _AclRouter:
def acl_check(self, peer_id: bytes, acl: object) -> bool:
return True
def _build_rptc(peer: bytes) -> bytes:
return RPTC + peer + b"CE1TEST " + b"\x00" * 85 + b"4"
def _master_protocol(*, passphrase: bytes = _PASSPHRASE) -> tuple[HBPProtocol, _RecordingTransport]:
transport = _RecordingTransport()
config = {
"GLOBAL": {"PING_TIME": 10, "MAX_MISSED": 3, "USE_ACL": False},
"SYSTEMS": {
"HOTSPOT": {
"MODE": "MASTER",
"ENABLED": True,
"MAX_PEERS": 8,
"PASSPHRASE": passphrase,
"OPTIONS": "TS2=9990;",
}
},
}
ensure_system_runtime_config(config)
hbp = HBPProtocol("HOTSPOT", config, router=_AclRouter()) # type: ignore[arg-type]
hbp.transport = transport # type: ignore[assignment]
return hbp, transport
def test_hbp_auth_handshake_rptl_rptk_rptc() -> None:
hbp, transport = _master_protocol()
hbp.datagramReceived(RPTL + _PEER, _CLIENT_ADDR)
assert len(transport.sent) == 1
challenge, addr = transport.sent[0]
assert challenge.startswith(RPTACK)
assert addr == _CLIENT_ADDR
assert hbp._peers[_PEER]["CONNECTION"] == "CHALLENGE_SENT"
salt_str = bytes_4(hbp._peers[_PEER]["SALT"])
sys_cfg = hbp._config
auth_hash = _calc_hash(salt_str, _get_passphrase_bytes(sys_cfg))
transport.sent.clear()
hbp.datagramReceived(RPTK + _PEER + auth_hash, _CLIENT_ADDR)
assert len(transport.sent) == 1
rptk_ack, addr = transport.sent[0]
assert rptk_ack.startswith(RPTACK)
assert addr == _CLIENT_ADDR
assert hbp._peers[_PEER]["CONNECTION"] == "WAITING_CONFIG"
transport.sent.clear()
hbp.datagramReceived(_build_rptc(_PEER), _CLIENT_ADDR)
assert len(transport.sent) == 1
rptc_ack, addr = transport.sent[0]
assert rptc_ack.startswith(RPTACK)
assert addr == _CLIENT_ADDR
assert hbp._peers[_PEER]["CONNECTION"] == "YES"
def test_hbp_auth_wrong_password_mstnak() -> None:
hbp, transport = _master_protocol()
hbp.datagramReceived(RPTL + _PEER, _CLIENT_ADDR)
salt_str = bytes_4(hbp._peers[_PEER]["SALT"])
bad_hash = _calc_hash(salt_str, b"wrong-password")
transport.sent.clear()
hbp.datagramReceived(RPTK + _PEER + bad_hash, _CLIENT_ADDR)
assert len(transport.sent) == 1
nak, addr = transport.sent[0]
assert nak.startswith(MSTNAK)
assert addr == _CLIENT_ADDR
assert _PEER not in hbp._peers

@ -0,0 +1,65 @@
# ADN DMR Peer Server - tests infrastructure obp validate source server
#
# 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 DMRE source-server validation must not use HBP ALLOW_UNREG_ID bypass."""
from __future__ import annotations
from adn_server.domain import bytes_4
from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol
_UNKNOWN = bytes_4(3120999)
_KNOWN = bytes_4(3120001)
def _obp_protocol(*, allow_unreg: bool | None = None) -> HBPProtocol:
system = {
"MODE": "OPENBRIDGE",
"PASSPHRASE": b"test-passphrase\x00\x00\x00\x00\x00\x00",
"VER": 5,
"TARGET_IP": "127.0.0.1",
"TARGET_PORT": 62030,
"NETWORK_ID": bytes_4(73010),
}
if allow_unreg is not None:
system["ALLOW_UNREG_ID"] = allow_unreg
config = {
"GLOBAL": {"SERVER_ID": bytes_4(73010)},
"_SUB_IDS": {3120001: "CE1TST"},
"SYSTEMS": {"OBP-A": system},
}
return HBPProtocol("OBP-A", config)
def test_obp_source_server_rejects_unknown_even_without_allow_unreg_id() -> None:
proto = _obp_protocol()
assert proto.validate_id(_UNKNOWN) is True
assert proto.validate_obp_source_server_id(_UNKNOWN) is False
def test_obp_source_server_accepts_known_subscriber_id() -> None:
proto = _obp_protocol()
assert proto.validate_obp_source_server_id(_KNOWN) == "CE1TST"
def test_obp_source_server_still_rejects_when_allow_unreg_disabled() -> None:
proto = _obp_protocol(allow_unreg=False)
assert proto.validate_id(_UNKNOWN) is False
assert proto.validate_obp_source_server_id(_UNKNOWN) is False

@ -0,0 +1,47 @@
# ADN DMR Peer Server - tests infrastructure security downloader
#
# 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
###############################################################################
"""SecurityDownloader stub and password_crypto round-trip."""
from __future__ import annotations
import pytest
from adn_server.infrastructure.security.password_crypto import decrypt_password
from adn_server.infrastructure.security.password_download import StubSecurityDownloader
pytest.importorskip("cryptography")
from cryptography.fernet import Fernet # noqa: E402
def test_password_crypto_roundtrip(tmp_path) -> None:
key_path = tmp_path / "encryption_key.secret"
key = Fernet.generate_key()
key_path.write_bytes(key)
plaintext = "peer-secret-42"
encrypted = Fernet(key).encrypt(plaintext.encode("utf-8")).decode("utf-8")
assert decrypt_password(encrypted, str(key_path)) == plaintext
def test_stub_security_downloader_is_noop() -> None:
stub = StubSecurityDownloader()
config = {"GLOBAL": {"URL_SECURITY": "", "PORT_SECURITY": "", "PASS_SECURITY": ""}}
stub.init_downloads(config)
stub.periodic_download(config)

@ -22,6 +22,8 @@
from __future__ import annotations
import logging
from adn_server.domain import bytes_4
from adn_server.infrastructure.hbp_constants import DMRD
from adn_server.infrastructure.mesh.obp_v1 import build_dmrd_v1
@ -29,6 +31,15 @@ from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol
_PASS = b"test-passphrase\x00\x00\x00\x00\x00\x00"
_SERVER = bytes_4(73010)
_OBP_ADDR = ("127.0.0.1", 62030)
class _RecordingTransport:
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 _sample_dmr_voice() -> bytes:
@ -87,3 +98,20 @@ def test_hbp_encode_mesh_egress_dmre_when_ver_ge_4() -> None:
assert wire is not None
assert len(wire) == 89
assert wire[:4] == b"DMRE"
def test_obp_hmac_reject_discards_and_warns(caplog) -> None:
proto = _obp_protocol()
proto._config["VER"] = 1
transport = _RecordingTransport()
proto.transport = transport # type: ignore[assignment]
proto._config["TARGET_SOCK"] = _OBP_ADDR
inner = _sample_dmr_voice()
wire = bytearray(build_dmrd_v1(inner, _SERVER, _PASS))
wire[-1] ^= 0xFF
with caplog.at_level(logging.WARNING):
proto.datagramReceived(bytes(wire), _OBP_ADDR)
assert transport.sent == []
assert any("HMAC failed" in record.message for record in caplog.records)

@ -26,7 +26,7 @@ from tests.harness.assertions import assert_forwarded
from tests.harness.deterministic import DeterministicScenario, PacketSpec
from tests.routing.test_echo_subscription_reset import _echo_scenario_config
from adn_server.infrastructure.bootstrap.peer_server import _seed_echo_routing_table
from adn_server.application.subscription.echo_seed import seed_echo_routing_table
def _prod_like_config() -> dict:
@ -49,7 +49,7 @@ def test_echo_routing_missing_without_echo_store_seed() -> None:
def test_echo_routing_works_after_echo_seed_on_vhead() -> None:
config = _prod_like_config()
scenario = DeterministicScenario(config=config, routing_table=_seed_echo_routing_table(config))
scenario = DeterministicScenario(config=config, routing_table=seed_echo_routing_table(config))
scenario.routing.apply_startup_subscriptions()
base = PacketSpec(dst_id=9990, stream_id=0xABCDEF02, slot=2)

@ -26,8 +26,8 @@ import pytest
from tests.harness.assertions import assert_forwarded
from tests.harness.deterministic import DeterministicScenario, PacketSpec, minimal_config
from adn_server.application.subscription.echo_seed import seed_echo_routing_table
from adn_server.domain import ID_MAX, PEER_MAX
from adn_server.infrastructure.bootstrap.peer_server import _seed_echo_routing_table
from adn_server.infrastructure.config_loader import acl_build
@ -67,7 +67,7 @@ def _echo_scenario_config() -> dict:
def _bridges_after_bridgereset(config: dict) -> dict:
bridges = _seed_echo_routing_table(config)
bridges = seed_echo_routing_table(config)
for entry in bridges["9990"]:
if entry["SYSTEM"] == "ECHO":
entry["ACTIVE"] = False

@ -0,0 +1,719 @@
# ADN DMR Peer Server - concurrent OBP streams downlink reproduction
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
"""Reproduce production bug: two concurrent OBP streams on one MASTER TS2 with a
dual-TG hotspot (SINGLE=0) intermittently drops the active listen stream.
Diagnosis-only: no production code changes. Uses DeterministicScenario for OBP
ingress + real HBPProtocol send_peer downlink gate.
"""
from __future__ import annotations
import copy
import logging
from contextlib import contextmanager
from typing import Any
import pytest
from tests.harness.deterministic import (
DeterministicScenario,
FakeClock,
PacketSpec,
add_openbridge_system,
patch_routing_wall_time,
)
from tests.support.hbp_repeat_stack import RecordingTransport
from adn_server.application.routing.downlink import (
DownlinkContext,
peer_slot_blocks_downlink,
)
from adn_server.application.routing.helpers import (
_peer_status_rx_hangtime_blocks,
_peer_transmit_hangtime_blocks,
hbp_slot_blocks_group_voice_for_peer,
peer_hotspot_voice_slot_busy,
peer_single_blocks_foreign_same_tg_downlink,
peer_single_same_tg_foreign_tx_blocks,
slot_has_active_voice,
)
from adn_server.domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3, bytes_4, int_id
from adn_server.domain.hbp_protocol import HBPF_VOICE, STREAM_TO
from adn_server.infrastructure.acl_router import InMemoryAclRouter
from adn_server.infrastructure.config_normalizer import ensure_system_runtime_config
from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol
_TG_A = 3109050
_TG_B = 52090
_STREAM_A = 0xAABBCCDD
_STREAM_B = 0x11223344
_HS_PEER = 730039101
_OBP_PEER = 73010
_RF_A = 3340001
_RF_B = 3340002
_INTERVAL_S = 0.06
_VOICE_BURSTS = 8
def _bridge_row(*, system: str, ts: int, tgid: int) -> dict[str, Any]:
tg_b = bytes_3(tgid)
return {
"SYSTEM": system,
"TS": ts,
"TGID": tg_b,
"ACTIVE": True,
"TIMEOUT": 3600.0,
"TO_TYPE": "ON",
"ON": [tg_b],
"OFF": [],
"RESET": [],
"TIMER": 0.0,
}
def _dual_tg_routing_table() -> dict[str, list[dict[str, Any]]]:
return {
str(_TG_A): [
_bridge_row(system="OBP-CL", ts=1, tgid=_TG_A),
_bridge_row(system="MASTER-A", ts=2, tgid=_TG_A),
],
str(_TG_B): [
_bridge_row(system="OBP-CL", ts=1, tgid=_TG_B),
_bridge_row(system="MASTER-A", ts=2, tgid=_TG_B),
],
}
@contextmanager
def _patch_harness_time(clock: FakeClock):
"""Patch wall clock on routing ingress and downlink gate paths."""
import adn_server.application.routing.downlink as downlink_mod
import adn_server.application.routing_use_cases as routing_mod
orig_routing = routing_mod.time.time
orig_downlink = downlink_mod.time.time
routing_mod.time.time = clock.time
downlink_mod.time.time = clock.time
try:
yield
finally:
routing_mod.time.time = orig_routing
downlink_mod.time.time = orig_downlink
def _build_stack() -> tuple[DeterministicScenario, HBPProtocol, RecordingTransport]:
routing_table = _dual_tg_routing_table()
scenario = DeterministicScenario(
routing_table=routing_table,
enable_reporting=True,
)
config = scenario.config
add_openbridge_system(config, "OBP-CL")
master = config["SYSTEMS"]["MASTER-A"]
master.update(
{
"MAX_PEERS": 8,
"GROUP_HANGTIME": 5.0,
"SINGLE_MODE": False,
"TS2_STATIC": f"{_TG_A},{_TG_B}",
}
)
ensure_system_runtime_config(config)
transport = RecordingTransport()
hbp = HBPProtocol("MASTER-A", config, router=InMemoryAclRouter())
hbp.transport = transport # type: ignore[assignment]
scenario.protocols["MASTER-A"] = hbp
def _send_to_system(target: str, packet: bytes, **kwargs: Any) -> None:
scenario.capture.recorder(target)(
packet,
hops=kwargs.get("_hops", kwargs.get("hops", b"")),
ber=kwargs.get("_ber", kwargs.get("ber", b"\x00")),
rssi=kwargs.get("_rssi", kwargs.get("rssi", b"\x00")),
source_server=kwargs.get(
"_source_server", kwargs.get("source_server", b"\x00\x00\x00\x00"),
),
source_rptr=kwargs.get(
"_source_rptr", kwargs.get("source_rptr", b"\x00\x00\x00\x00"),
),
)
proto = scenario.protocols.get(target)
if proto is not None and hasattr(proto, "send_system"):
proto.send_system(
packet,
_hops=kwargs.get("_hops", b""),
_ber=kwargs.get("_ber", b"\x00"),
_rssi=kwargs.get("_rssi", b"\x00"),
_source_server=kwargs.get("_source_server", b"\x00\x00\x00\x00"),
_source_rptr=kwargs.get("_source_rptr", b"\x00\x00\x00\x00"),
)
scenario.routing._send_to_system = _send_to_system
hs = bytes_4(_HS_PEER)
hbp._peers[hs] = {
"CONNECTION": "YES",
"CONNECTED": scenario.clock.time(),
"LAST_PING": scenario.clock.time(),
"SOCKADDR": ("127.0.0.1", 62031),
"OPTIONS": f"TS2={_TG_A},{_TG_B};SINGLE=0;".encode(),
}
master.setdefault("PEERS", {})[hs] = hbp._peers[hs]
hbp._refresh_connected_peer_count()
hbp._mark_downlink_index_dirty()
return scenario, hbp, transport
def _base_spec(*, tgid: int, stream_id: int, rf_src: int) -> PacketSpec:
return PacketSpec(
peer_id=_OBP_PEER,
rf_src=rf_src,
dst_id=tgid,
slot=1,
stream_id=stream_id,
)
def _dmrd_packet(
*,
tgid: int,
stream_id: int,
rf_src: int,
frame_type: int,
dtype_vseq: int,
seq: int = 0,
) -> bytes:
base = _base_spec(tgid=tgid, stream_id=stream_id, rf_src=rf_src)
spec = PacketSpec(
peer_id=base.peer_id,
rf_src=base.rf_src,
dst_id=base.dst_id,
slot=base.slot,
stream_id=base.stream_id,
seq=seq,
frame_type=frame_type,
dtype_vseq=dtype_vseq,
payload=base.payload,
)
return spec.data()
def _is_voice_burst(packet: bytes) -> bool:
return ((packet[15] & 0x30) >> 4) == HBPF_VOICE
def _parse_dmrd_args(packet: bytes) -> dict[str, Any]:
bits = packet[15]
return {
"peer_id": packet[11:15],
"rf_src": packet[5:8],
"dst_id": packet[8:11],
"seq": packet[4],
"slot": 2 if bits & 0x80 else 1,
"call_type": "group",
"frame_type": (bits & 0x30) >> 4,
"dtype_vseq": bits & 0xF,
"stream_id": packet[16:20],
"data": packet,
}
def _wrap_send_peer_trace(hbp: HBPProtocol) -> dict[str, Any]:
"""Record per-packet send_peer accept/reject without modifying production code."""
trace: dict[str, Any] = {"accepted": [], "rejected": []}
orig = hbp.send_peer
def _traced(peer_id: bytes, packet: bytes, *, _skip_dual_expand: bool = False) -> None:
route_pkt = packet
peer = hbp._peers.get(peer_id)
if peer is not None and packet[:4] == b"DMRD":
from adn_server.application.routing.downlink import remap_dmrd_for_peer
route_pkt = remap_dmrd_for_peer(packet, peer, hbp._config, peer_id=peer_id)
opts_ok = hbp._peer_should_receive_dmrd(peer_id, packet)
gate_ok = hbp._peer_would_accept_group_dmrd(
peer_id, packet if not _skip_dual_expand else route_pkt, routed=_skip_dual_expand,
)
row = {
"tgid": int_id(route_pkt[8:11]),
"stream": route_pkt[16:20],
"opts_ok": opts_ok,
"gate_ok": gate_ok,
"blocked": not (opts_ok and (gate_ok if packet[:4] == b"DMRD" else True)),
}
if row["blocked"]:
trace["rejected"].append(row)
else:
trace["accepted"].append(row)
orig(peer_id, packet, _skip_dual_expand=_skip_dual_expand)
hbp.send_peer = _traced # type: ignore[method-assign]
return trace
def _inject_obp_at(
scenario: DeterministicScenario,
packet: bytes,
*,
pkt_time: float,
) -> None:
scenario.clock.now = pkt_time
with _patch_harness_time(scenario.clock), patch_routing_wall_time(scenario.clock):
scenario.routing.dmrd_received(
"OBP-CL",
ingress_pkt_time=pkt_time,
obp_use_parsed=True,
obp_hops=b"\x00\x00\x00\x00",
obp_source_server=bytes_4(9990),
**_parse_dmrd_args(packet),
)
def _run_interleaved_qsos(
scenario: DeterministicScenario,
*,
voice_bursts: int = _VOICE_BURSTS,
interval_s: float = _INTERVAL_S,
) -> dict[str, Any]:
t0 = scenario.clock.time()
step = 0
tx_stamp_changes = 0
prev_tx_stream: bytes | None = None
hs_addr = ("127.0.0.1", 62031)
delivered_a: list[bytes] = []
delivered_b: list[bytes] = []
blocked_a_checks: list[dict[str, Any]] = []
def _slot_st() -> dict[str, Any]:
proto = scenario.protocols["MASTER-A"]
assert isinstance(proto, HBPProtocol)
return proto.STATUS.setdefault(2, {})
def _record_delivery(pkt: bytes) -> None:
tgid = int_id(pkt[8:11])
sid = pkt[16:20]
if tgid == _TG_A and sid == bytes_4(_STREAM_A):
delivered_a.append(pkt)
elif tgid == _TG_B and sid == bytes_4(_STREAM_B):
delivered_b.append(pkt)
packets_plan: list[tuple[str, bytes]] = []
for label, tgid, stream, rf in (
("A", _TG_A, _STREAM_A, _RF_A),
("B", _TG_B, _STREAM_B, _RF_B),
):
packets_plan.append(
(
label,
_dmrd_packet(
tgid=tgid,
stream_id=stream,
rf_src=rf,
frame_type=HBPF_DATA_SYNC,
dtype_vseq=HBPF_SLT_VHEAD,
),
)
)
for seq in range(1, voice_bursts + 1):
for label, tgid, stream, rf in (
("A", _TG_A, _STREAM_A, _RF_A),
("B", _TG_B, _STREAM_B, _RF_B),
):
packets_plan.append(
(
label,
_dmrd_packet(
tgid=tgid,
stream_id=stream,
rf_src=rf,
frame_type=HBPF_VOICE,
dtype_vseq=(seq % 4) or 1,
seq=seq,
),
)
)
for label, tgid, stream, rf in (
("A", _TG_A, _STREAM_A, _RF_A),
("B", _TG_B, _STREAM_B, _RF_B),
):
packets_plan.append(
(
label,
_dmrd_packet(
tgid=tgid,
stream_id=stream,
rf_src=rf,
frame_type=HBPF_DATA_SYNC,
dtype_vseq=HBPF_SLT_VTERM,
seq=99,
),
)
)
scenario.capture.packets.clear()
transport = scenario.protocols["MASTER-A"]
assert isinstance(transport, HBPProtocol)
rec = transport.transport
assert isinstance(rec, RecordingTransport)
rec.clear()
hbp = transport
hs = bytes_4(_HS_PEER)
peer = hbp._peers[hs]
for label, pkt in packets_plan:
pkt_time = t0 + step * interval_s
step += 1
_inject_obp_at(scenario, pkt, pkt_time=pkt_time)
cur = _slot_st().get("TX_STREAM_ID")
if cur != prev_tx_stream:
tx_stamp_changes += 1
prev_tx_stream = cur
bits = pkt[15] | 0x80
route_pkt = pkt[:15] + bytes([bits]) + pkt[16:]
ctx = hbp._downlink_ctx()
if (
label == "A"
and int_id(pkt[8:11]) == _TG_A
and pkt[16:20] == bytes_4(_STREAM_A)
and (pkt[15] & 0xF) != HBPF_SLT_VTERM
):
blocked = peer_slot_blocks_downlink(ctx, hs, peer, route_pkt, pkt_time=pkt_time)
if blocked:
blocked_a_checks.append(
{
"pkt_time": pkt_time,
"step": step - 1,
"slot_st": copy.deepcopy(_slot_st()),
"peer_slots": copy.deepcopy(ctx.peer_voice_slots.get(hs, {})),
"reason": diagnose_downlink_block(ctx, hs, peer, route_pkt, pkt_time),
}
)
for sent_pkt, addr in rec.sent:
if addr == hs_addr:
_record_delivery(sent_pkt)
rec.clear()
return {
"tx_stamp_changes": tx_stamp_changes,
"start_tx_events": len(
[ev for ev in (scenario.report_factory.events if scenario.report_factory else []) if ",START,TX," in ev]
),
"delivered_a": delivered_a,
"delivered_b": delivered_b,
"blocked_a_checks": blocked_a_checks,
"final_slot_st": copy.deepcopy(_slot_st()),
}
def diagnose_downlink_block(
ctx: DownlinkContext,
peer_id: bytes,
peer: dict[str, Any],
route_pkt: bytes,
pkt_time: float,
) -> str:
"""Pin which gate branch would block this downlink (mirrors helpers/downlink)."""
if route_pkt[:4] != b"DMRD":
return "not_dmrd"
stream_id = route_pkt[16:20]
incoming_tgid_b = route_pkt[8:11]
hang = float(ctx.sys_cfg.get("GROUP_HANGTIME", 0) or 0)
peer_slots = ctx.peer_voice_slots.get(bytes_4(int_id(peer_id)))
pk = bytes_4(int_id(peer_id))
for voice_slot in (2,):
hang_row = ctx.peer_voice_hangtime.get(pk, {}).get(voice_slot)
slot_st = ctx.status.get(voice_slot, {})
if _peer_transmit_hangtime_blocks(hang_row, incoming_tgid_b, pkt_time, hang):
return f"helpers.py:_peer_transmit_hangtime_blocks voice_slot={voice_slot}"
active = (peer_slots or {}).get(voice_slot)
if isinstance(active, dict):
incoming_tgid = int_id(incoming_tgid_b)
active_tgid = int(active.get("tgid", 0) or 0)
active_stream = active.get("stream_id")
active_time = float(active.get("time", 0) or 0)
age = pkt_time - active_time
if active.get("ingress"):
return f"helpers.py:peer_hotspot_voice_slot_busy ingress voice_slot={voice_slot}"
if active.get("bridge_hold") and active_tgid and incoming_tgid != active_tgid:
if age <= hang:
return (
f"helpers.py:peer_hotspot_voice_slot_busy bridge_hold "
f"active_tg={active_tgid} incoming={incoming_tgid}"
)
if active_stream and stream_id:
if active_stream == stream_id:
pass
elif active_tgid and active_tgid == incoming_tgid:
if age < STREAM_TO:
return (
f"helpers.py:peer_hotspot_voice_slot_busy "
f"same_tg_different_stream active={active_stream!r} incoming={stream_id!r}"
)
elif active_tgid and incoming_tgid != active_tgid:
return (
f"helpers.py:peer_hotspot_voice_slot_busy "
f"elif active_tgid={active_tgid} incoming_tgid={incoming_tgid}"
)
else:
return "helpers.py:peer_hotspot_voice_slot_busy else branch (stream/tg mismatch)"
elif isinstance(active, dict):
if not (active_tgid and active_tgid == incoming_tgid):
return (
f"helpers.py:peer_hotspot_voice_slot_busy "
f"no_active_stream active_tg={active_tgid} incoming={incoming_tgid}"
)
if peer_single_blocks_foreign_same_tg_downlink(
peer, pk, voice_slot, incoming_tgid_b, peer_slots, ctx.sys_cfg, now=pkt_time,
):
return f"helpers.py:peer_single_blocks_foreign_same_tg_downlink voice_slot={voice_slot}"
if peer_single_same_tg_foreign_tx_blocks(
peer, pk, incoming_tgid_b, stream_id, slot_st, ctx.sys_cfg, pkt_time=pkt_time,
):
return (
f"helpers.py:peer_single_same_tg_foreign_tx_blocks "
f"TX_STREAM_ID={slot_st.get('TX_STREAM_ID')!r} slot_active={slot_has_active_voice(slot_st, pkt_time)}"
)
if bytes_4(int_id(slot_st.get("RX_PEER", b""))) == pk:
rx_active = (
slot_st.get("RX_TYPE") is not None
and slot_st.get("RX_TYPE") != HBPF_SLT_VTERM
and (pkt_time - float(slot_st.get("RX_TIME", 0))) < STREAM_TO
)
if rx_active and stream_id != slot_st.get("RX_STREAM_ID"):
return (
f"helpers.py:peer_hotspot_voice_slot_busy "
f"rx_active stream_mismatch RX_STREAM_ID={slot_st.get('RX_STREAM_ID')!r}"
)
if _peer_status_rx_hangtime_blocks(
peer_id, slot_st, incoming_tgid_b, pkt_time, hang,
):
return f"helpers.py:_peer_status_rx_hangtime_blocks voice_slot={voice_slot}"
if hbp_slot_blocks_group_voice_for_peer(
slot_st,
peer_id,
incoming_tgid_b,
stream_id,
pkt_time,
hang,
per_peer=True,
peers=ctx.peers,
peer_slots=peer_slots,
peer_hang_row=hang_row,
voice_slot=voice_slot,
sys_cfg=ctx.sys_cfg,
):
return f"helpers.py:hbp_slot_blocks_group_voice_for_peer voice_slot={voice_slot}"
if peer_slot_blocks_downlink(ctx, peer_id, peer, route_pkt, pkt_time=pkt_time):
return "downlink.py:peer_slot_blocks_downlink (composite)"
return "not_blocked"
def test_flip_flop_slot_st_blocks_active_listen_on_cross_tg() -> None:
"""Pure unit: peer listening on TG A must not RX TG B while session is open."""
now = 1_000_000.0
hs = bytes_4(_HS_PEER)
peer = {"OPTIONS": f"TS2={_TG_A},{_TG_B};SINGLE=0;".encode()}
stream_a = bytes_4(_STREAM_A)
stream_b = bytes_4(_STREAM_B)
peer_slots = {
2: {"stream_id": stream_a, "tgid": _TG_A, "time": now, "ingress": False},
}
slot_st = {
"TX_PEER": bytes_4(_OBP_PEER),
"TX_STREAM_ID": stream_b,
"TX_TYPE": HBPF_SLT_VHEAD,
"TX_TIME": now,
"TX_TGID": bytes_3(_TG_B),
"RX_TYPE": HBPF_SLT_VTERM,
}
assert peer_hotspot_voice_slot_busy(
hs, 2, stream_b, bytes_3(_TG_B), slot_st, peer_slots, None, now + 0.06, 5.0,
peer=peer, sys_cfg={"GROUP_HANGTIME": 5.0, "SINGLE_MODE": False},
)
route_b = _dmrd_packet(
tgid=_TG_B,
stream_id=_STREAM_B,
rf_src=_RF_B,
frame_type=HBPF_VOICE,
dtype_vseq=1,
seq=1,
)
route_b = route_b[:15] + bytes([route_b[15] | 0x80]) + route_b[16:]
ctx = DownlinkContext(
config={"SYSTEMS": {"MASTER-A": {"MODE": "MASTER"}}},
system_name="MASTER-A",
sys_cfg={"GROUP_HANGTIME": 5.0, "SINGLE_MODE": False, "MODE": "MASTER"},
peers={hs: peer},
status={2: slot_st},
peer_voice_slots={hs: copy.deepcopy(peer_slots)},
connected_count=2,
)
reason = diagnose_downlink_block(ctx, hs, peer, route_b, now + 0.06)
assert "active_tgid=3109050 incoming_tgid=52090" in reason
def test_flip_flop_slot_st_same_stream_not_blocked_despite_foreign_tx_stamp() -> None:
"""Same stream id on active listen must pass even when flat TX row shows other stream."""
now = 1_000_000.0
hs = bytes_4(_HS_PEER)
peer = {"OPTIONS": f"TS2={_TG_A},{_TG_B};SINGLE=0;".encode()}
stream_a = bytes_4(_STREAM_A)
stream_b = bytes_4(_STREAM_B)
peer_slots = {
2: {"stream_id": stream_a, "tgid": _TG_A, "time": now, "ingress": False},
}
slot_st = {
"TX_PEER": bytes_4(_OBP_PEER),
"TX_STREAM_ID": stream_b,
"TX_TYPE": HBPF_SLT_VHEAD,
"TX_TIME": now,
"TX_TGID": bytes_3(_TG_B),
"RX_TYPE": HBPF_SLT_VTERM,
}
assert not peer_hotspot_voice_slot_busy(
hs, 2, stream_a, bytes_3(_TG_A), slot_st, peer_slots, None, now + 0.06, 5.0,
peer=peer, sys_cfg={"GROUP_HANGTIME": 5.0, "SINGLE_MODE": False},
)
def test_rx_active_own_peer_foreign_stream_blocks_active_downlink() -> None:
"""When STATUS RX shows the hotspot mid-TX on another stream, stream-A downlink blocks.
helpers.py peer_hotspot_voice_slot_busy lines 411-418:
``rx_active and stream_id != slot_st.get('RX_STREAM_ID')`` → True.
"""
now = 1_000_000.0
hs = bytes_4(_HS_PEER)
peer = {"OPTIONS": f"TS2={_TG_A},{_TG_B};SINGLE=0;".encode()}
stream_a = bytes_4(_STREAM_A)
stream_b = bytes_4(_STREAM_B)
peer_slots = {
2: {"stream_id": stream_a, "tgid": _TG_A, "time": now, "ingress": False},
}
slot_st = {
"RX_PEER": hs,
"RX_STREAM_ID": stream_b,
"RX_TYPE": HBPF_SLT_VHEAD,
"RX_TIME": now,
"RX_TGID": bytes_3(_TG_B),
"TX_PEER": bytes_4(_OBP_PEER),
"TX_STREAM_ID": stream_a,
"TX_TYPE": HBPF_SLT_VHEAD,
"TX_TIME": now,
"TX_TGID": bytes_3(_TG_A),
}
route_a = _dmrd_packet(
tgid=_TG_A,
stream_id=_STREAM_A,
rf_src=_RF_A,
frame_type=HBPF_VOICE,
dtype_vseq=1,
seq=1,
)
route_a = route_a[:15] + bytes([route_a[15] | 0x80]) + route_a[16:]
ctx = DownlinkContext(
config={"SYSTEMS": {"MASTER-A": {"MODE": "MASTER"}}},
system_name="MASTER-A",
sys_cfg={"GROUP_HANGTIME": 5.0, "SINGLE_MODE": False, "MODE": "MASTER"},
peers={hs: peer},
status={2: slot_st},
peer_voice_slots={hs: copy.deepcopy(peer_slots)},
connected_count=2,
)
reason = diagnose_downlink_block(ctx, hs, peer, route_a, now + 0.06)
assert "rx_active stream_mismatch" in reason
assert peer_slot_blocks_downlink(ctx, hs, peer, route_a, pkt_time=now + 0.06)
def test_concurrent_obp_tx_row_flip_flop(caplog: pytest.LogCaptureFixture) -> None:
"""BUG: shared STATUS[2] TX_STREAM_ID must change at most twice per stream (VHEAD), not per packet."""
caplog.set_level(logging.INFO)
scenario, _hbp, _transport = _build_stack()
result = _run_interleaved_qsos(scenario, voice_bursts=6)
# Expect: 2 streams × 1 VHEAD stamp each (+ maybe VTERM) — not ~2×packets.
max_expected_stamps = 6
assert result["tx_stamp_changes"] <= max_expected_stamps, (
f"TX_STREAM_ID flip-flop: {result['tx_stamp_changes']} changes "
f"(expected <= {max_expected_stamps}); final STATUS[2]={result['final_slot_st']}"
)
assert result["start_tx_events"] <= max_expected_stamps, (
f"START,TX report spam: {result['start_tx_events']} events"
)
def test_concurrent_obp_active_stream_delivered_without_gaps() -> None:
"""After VHEAD, stream-A voice should not be dropped mid-QSO (minimal harness)."""
scenario, hbp, transport = _build_stack()
send_trace = _wrap_send_peer_trace(hbp)
result = _run_interleaved_qsos(scenario, voice_bursts=_VOICE_BURSTS)
voice_a = [p for p in result["delivered_a"] if _is_voice_burst(p)]
rejected_a = [r for r in send_trace["rejected"] if r["tgid"] == _TG_A]
expected_voice = _VOICE_BURSTS
# Minimal OBP-only harness delivers all A packets; production drops use RX/status
# corruption paths (see test_rx_active_own_peer_foreign_stream_blocks_active_downlink).
assert not rejected_a, f"unexpected stream-A send_peer rejects: {rejected_a}"
assert len(voice_a) == expected_voice, (
f"stream A voice delivered {len(voice_a)}/{expected_voice}"
)
# Stream B may be dropped for the listening hotspot — not asserted.
def test_concurrent_obp_stream_b_blocked_on_listen_session() -> None:
"""Stream B must be dropped while hotspot listens to stream A (one QSO per slot)."""
scenario, hbp, transport = _build_stack()
send_trace = _wrap_send_peer_trace(hbp)
_run_interleaved_qsos(scenario, voice_bursts=4)
rejected_b = [r for r in send_trace["rejected"] if r["tgid"] == _TG_B]
assert rejected_b, "stream B should be blocked for listening hotspot"
ctx = hbp._downlink_ctx()
hs = bytes_4(_HS_PEER)
peer = hbp._peers[hs]
sample = _dmrd_packet(
tgid=_TG_B, stream_id=_STREAM_B, rf_src=_RF_B,
frame_type=HBPF_VOICE, dtype_vseq=1, seq=1,
)
sample = sample[:15] + bytes([sample[15] | 0x80]) + sample[16:]
reason = diagnose_downlink_block(ctx, hs, peer, sample, scenario.clock.time())
assert (
"active_tgid=3109050 incoming_tgid=52090" in reason
or "_peer_transmit_hangtime_blocks" in reason
or "peer_hotspot" in reason
or "peer_slot" in reason
), f"stream B blocked but unexpected gate: {reason}"
def test_concurrent_obp_diagnose_active_stream_block_branch() -> None:
"""Identify the exact branch that drops active stream-A packets mid-QSO."""
scenario, hbp, transport = _build_stack()
send_trace = _wrap_send_peer_trace(hbp)
result = _run_interleaved_qsos(scenario, voice_bursts=_VOICE_BURSTS)
rejected_a = [
r for r in send_trace["rejected"]
if r["tgid"] == _TG_A and r["stream"] == bytes_4(_STREAM_A)
]
if not rejected_a and not result["blocked_a_checks"]:
# After per-leg OBP bridge TX wiring, flat TX_STREAM_ID no longer flip-flops per packet.
if result["tx_stamp_changes"] <= 4:
pytest.skip("no stream-A rejects; per-leg bridge TX fix removed flat-row flip-flop")
assert result["tx_stamp_changes"] > 4, "expected TX row flip-flop in concurrent OBP harness"
pytest.skip("no stream-A send_peer rejects — inspect flip-flop collateral (START,TX spam)")
reasons: dict[str, int] = {}
for row in result["blocked_a_checks"]:
reasons[row["reason"]] = reasons.get(row["reason"], 0) + 1
for rej in rejected_a:
if not rej["opts_ok"]:
key = "udp_hbp.py:_peer_should_receive_dmrd OPTIONS/eligibility"
elif not rej["gate_ok"]:
key = "udp_hbp.py:_peer_would_accept_group_dmrd -> peer_slot_blocks_downlink"
else:
key = "unknown"
reasons[key] = reasons.get(key, 0) + 1
top = max(reasons, key=reasons.get)
assert "peer_hotspot" in top or "peer_slot" in top or "send_peer" in top or "udp_hbp" in top, (
f"unexpected block reasons: {reasons}"
)

@ -0,0 +1,131 @@
# ADN DMR Peer Server - tests schemas report wire contract
#
# 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
###############################################################################
"""ReportWire output must validate against report-v2.json and match committed examples."""
from __future__ import annotations
import json
from pathlib import Path
from unittest.mock import patch
import jsonschema
import pytest
from adn_server.domain.value_objects import bytes_4
from adn_server.infrastructure.twisted_adapters.report.opcodes import REPORT_OPCODES
from adn_server.infrastructure.twisted_adapters.report.wire import ReportWire
_SCHEMA_PATH = Path(__file__).resolve().parents[2] / "schemas" / "report-v2.json"
_EXAMPLES_DIR = _SCHEMA_PATH.parent / "examples"
_VOICE_CSV = "GROUP VOICE,START,RX,MASTER-A,2155905152,1001,3120001,2,52090"
@pytest.fixture(scope="module")
def validator() -> jsonschema.Draft202012Validator:
with _SCHEMA_PATH.open(encoding="utf-8") as fh:
return jsonschema.Draft202012Validator(json.load(fh))
def _dashboard_systems() -> dict:
return {
"MASTER-A": {
"MODE": "MASTER",
"ENABLED": True,
"IP": "10.0.0.1",
"PORT": 62030,
"SINGLE_MODE": False,
"DEFAULT_UA_TIMER": 10,
"TS1_STATIC": "91,92",
"TS2_STATIC": "7302",
"PEERS": {
bytes_4(3120001): {
"CONNECTION": "YES",
"CONNECTED": 1717555200,
"IP": "10.0.0.50",
"PORT": 62031,
"CALLSIGN": b"CE5RPY ",
},
},
},
"MASTER-B": {
"MODE": "MASTER",
"ENABLED": True,
"IP": "10.0.0.2",
"PORT": 62032,
"PEERS": {
bytes_4(3120002): {
"CONNECTION": "YES",
"CONNECTED": 1717555200,
"IP": "10.0.0.51",
"PORT": 62033,
"CALLSIGN": b"EA5GVK ",
},
},
},
"OBP-CL": {
"MODE": "OPENBRIDGE",
"ENABLED": True,
"IP": "44.31.61.68",
"PORT": 62999,
"NETWORK_ID": 73010,
"ENHANCED_OBP": True,
"PEERS": {},
},
}
def _wire_json(frame: bytes) -> dict:
return json.loads(frame[1:].decode("utf-8"))
def test_report_wire_state_frames_match_dashboard_state_example(
validator: jsonschema.Draft202012Validator,
) -> None:
with (_EXAMPLES_DIR / "dashboard_state.json").open(encoding="utf-8") as fh:
expected = json.load(fh)
expected_wire = {key: value for key, value in expected.items() if key != "server_id"}
wire = ReportWire()
with patch("adn_server.infrastructure.twisted_adapters.report.wire.time") as mock_time:
mock_time.time.return_value = expected["ts"]
frames = wire.state_frames(_dashboard_systems(), force=True)
assert len(frames) == 1
assert frames[0][:1] == REPORT_OPCODES["STATE_SND"]
payload = _wire_json(frames[0])
assert json.loads(json.dumps(payload)) == expected_wire
validator.validate({**payload, "server_id": expected.get("server_id", "7302")})
def test_report_wire_bridge_event_frames_match_voice_event_example(
validator: jsonschema.Draft202012Validator,
) -> None:
with (_EXAMPLES_DIR / "voice_event.json").open(encoding="utf-8") as fh:
expected = json.load(fh)
wire = ReportWire()
frames = wire.bridge_event_frames(_VOICE_CSV)
assert len(frames) == 1
assert frames[0][:1] == REPORT_OPCODES["VOICE_EVENT_SND"]
payload = _wire_json(frames[0])
assert payload["type"] == expected["type"]
for key, value in expected.items():
if key == "ts":
continue
assert payload[key] == value
validator.validate({**payload, "ts": expected["ts"]})

@ -38,7 +38,7 @@ def _reload_uc(scenario, *, start_looping_call) -> VoiceUseCases:
)
def test_check_voice_config_reload_starts_enabled_announcement_loop() -> None:
def test_apply_voice_config_starts_enabled_announcement_loop() -> None:
scenario, _ = voice_master_scenario()
scenario.config["VOICE"] = {
"ANNOUNCEMENTS": [
@ -61,14 +61,14 @@ def test_check_voice_config_reload_starts_enabled_announcement_loop() -> None:
return handle
uc = _reload_uc(scenario, start_looping_call=start_looping_call)
uc.check_voice_config_reload()
uc.apply_voice_config()
assert len(started) == 1
assert started[0][0] == 120.0
assert 0 in uc._ann_tasks
def test_check_voice_config_reload_stops_removed_announcement() -> None:
def test_apply_voice_config_stops_removed_announcement() -> None:
scenario, _ = voice_master_scenario()
scenario.config["VOICE"] = {"ANNOUNCEMENTS": [{"ENABLED": False, "TG": 91, "FILE": "x"}]}
stop_mock = MagicMock()
@ -80,13 +80,13 @@ def test_check_voice_config_reload_stops_removed_announcement() -> None:
)
uc._ann_tasks[0] = MagicMock(running=True, stop=stop_mock)
uc.check_voice_config_reload()
uc.apply_voice_config()
assert 0 not in uc._ann_tasks
stop_mock.assert_called_once()
def test_check_voice_config_reload_starts_tts_loop() -> None:
def test_apply_voice_config_starts_tts_loop() -> None:
scenario, _ = voice_master_scenario()
scenario.config["VOICE"] = {
"TTS_ANNOUNCEMENTS": [
@ -106,7 +106,7 @@ def test_check_voice_config_reload_starts_tts_loop() -> None:
return MagicMock(running=True, stop=MagicMock())
uc = _reload_uc(scenario, start_looping_call=start_looping_call)
uc.check_voice_config_reload()
uc.apply_voice_config()
assert started == [30.0]
assert 0 in uc._tts_tasks

Loading…
Cancel
Save

Powered by TurnKey Linux.