fix: echo TG 9990 monitor parity and UA session exclusions (#7)

Exclude 9990–9999 from SINGLE/UA locks, remap inject-only echo TX events,
resolve bridge TX peer_id for monitor display, and fix dynamic TG restore callback.
pull/8/head
ce5rpy 3 months ago committed by GitHub
parent 8574e1e42f
commit 3b8332cfa4
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -113,7 +113,7 @@ class DynamicTgUseCases:
]
tgids = restore_peer_ua_entries_to_memory(sys_cfg, peer_id, active, now=now)
if self._on_restored is not None:
self._on_restored(peer_id, system_name, sys_cfg, active, now)
self._on_restored(peer_id, system_name, sys_cfg, active, now=now)
if tgids:
logger.info(
"(DYNAMIC_TG) Restored %s TG(s) for peer %s on %s: %s",

@ -31,7 +31,7 @@ from __future__ import annotations
import copy
from typing import Any
from adn_server.application.routing.helpers import peer_should_receive_group_voice
from adn_server.application.routing.helpers import is_special_tg, peer_should_receive_group_voice
from adn_server.application.proxy.deployment import is_proxy_inject_only, proxy_target_system
from adn_server.domain.value_objects import bytes_4, int_id
@ -246,7 +246,12 @@ def _peers_receiving_tgid(
def _echo_tx_target_peer(parts: list[str], peers: dict[Any, Any]) -> bytes | None:
"""Echo/static downlink: field 5 is 9990 and field 6 resolves one hotspot."""
"""Echo/service downlink TX: field 5 may be 9990 or hotspot id; field 8 is 9990–9999."""
tgid_slot = _voice_event_tgid_slot(parts)
if tgid_slot is not None:
tgid, _ = tgid_slot
if is_special_tg(str(tgid)):
return _peer_key_from_voice_csv(parts, peers)
if len(parts) <= 5:
return None
try:
@ -320,13 +325,19 @@ def remap_inject_proxy_voice_events(
trx = parts[2].strip() if len(parts) > 2 else ""
if trx == "TX":
tgid_slot = _voice_event_tgid_slot(parts)
echo_peer = _echo_tx_target_peer(parts, peers)
if echo_peer is not None:
slot = slot_map.get(echo_peer)
if slot is not None:
tx_parts = list(parts)
if tgid_slot is not None:
tgid, _ = tgid_slot
if is_special_tg(str(tgid)):
tx_parts[5] = str(tgid)
return [
_remap_voice_event_to_slot(
parts, target=target, slot=slot, peer_key=echo_peer
tx_parts, target=target, slot=slot, peer_key=echo_peer
)
]
try:

@ -104,8 +104,14 @@ def tg4000_reset_on_vhead(int_dst_id: int, frame_type: int, dtype_vseq: int) ->
def is_ua_session_tgid(tgid: int) -> bool:
"""True when a keyed TG may be stored as a user-activated dynamic session."""
return int(tgid) > 0 and int(tgid) != 4000
"""True when a keyed TG may be stored as a user-activated dynamic session.
Excludes TG 4000 (reset command) and service/echo 9990–9999 (no SINGLE lock).
"""
t = int(tgid)
if t <= 0 or t == 4000:
return False
return not is_special_tg(str(t))
def obp_target_bcsq_quenches_stream(

@ -417,6 +417,11 @@ class RoutingUseCases(
bridge_match_slot=bridge_match_slot,
dst_int=dst_int,
)
_tx_report_peer = int_id(peer_id)
if not source_is_obp:
_tx_report_peer = int_id(
resolve_voice_peer_id(peer_id, rf_src, system_name, systems_cfg)
)
forwarded = []
_leg_iter: list[tuple[str, dict[str, Any]]] = [
(
@ -496,7 +501,7 @@ class RoutingUseCases(
)
self._send_routing_event(
"GROUP VOICE,START,TX,{},{},{},{},{},{}".format(
entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry.get("TS", 1), int_id(target_tgid)
entry["SYSTEM"], int_id(stream_id), _tx_report_peer, int_id(rf_src), entry.get("TS", 1), int_id(target_tgid)
)
)
if "EMB_LC" not in _target_status[stream_id]:
@ -533,7 +538,7 @@ class RoutingUseCases(
call_duration = pkt_time - _target_status[stream_id].get("START", pkt_time)
self._send_routing_event(
"GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}".format(
entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry.get("TS", 1), int_id(target_tgid), call_duration
entry["SYSTEM"], int_id(stream_id), _tx_report_peer, int_id(rf_src), entry.get("TS", 1), int_id(target_tgid), call_duration
)
)
elif dtype_vseq in (1, 2, 3, 4):
@ -624,7 +629,7 @@ class RoutingUseCases(
"GROUP VOICE,START,TX,{},{},{},{},{},{}".format(
entry["SYSTEM"],
int_id(stream_id),
int_id(peer_id),
_tx_report_peer,
int_id(rf_src),
entry_ts,
int_id(entry_tgid_b),
@ -648,11 +653,18 @@ class RoutingUseCases(
dmrbits = _ts_st["TX_T_LC"][0:98] + dmrbits[98:166] + _ts_st["TX_T_LC"][98:197]
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(
resolve_voice_peer_id(
_end_peer, rf_src, system_name, systems_cfg
)
)
self._send_routing_event(
"GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}".format(
entry["SYSTEM"],
int_id(stream_id),
int_id(_end_peer),
_end_report_peer,
int_id(rf_src),
entry_ts,
int_id(entry_tgid_b),

@ -26,6 +26,7 @@ from adn_server.application.dynamic_tg_use_cases import DynamicTgUseCases
from adn_server.application.routing.helpers import register_peer_ua_session
from adn_server.domain.dynamic_tg import DynamicTgEntry
from tests.application.test_peer_single_downlink import _peer_id, _sys_cfg
from twisted.internet import defer
class _CaptureStore:
@ -66,3 +67,48 @@ def test_persist_after_register_stores_absolute_expires_from_memory() -> None:
entry = store.replaced[0]
assert entry.expires_at == now + 60.0 * 60.0
assert entry.single_mode is True
def test_restore_peer_passes_now_as_keyword_to_on_restored() -> None:
"""sync_restored_dynamic_tgs declares ``now`` keyword-only; restore must match."""
import time
peer_id = _peer_id()
sys_cfg = _sys_cfg()
now = time.time()
entry = DynamicTgEntry(
int_id=730039101,
system_name="MASTER-A",
slot=2,
tgid=7304,
single_mode=True,
expires_at=now + 3600.0,
updated_at=now,
)
seen: dict[str, float] = {}
def on_restored(
_peer_id: bytes,
_system_name: str,
_sys_cfg: dict,
entries: list[DynamicTgEntry],
*,
now: float,
) -> None:
seen["now"] = now
seen["count"] = float(len(entries))
class _Store:
def load_peer(self, int_id: int, system_name: str):
return defer.succeed([entry])
def purge_expired(self, now: float) -> None:
pass
uc = DynamicTgUseCases(_Store(), on_restored=on_restored)
out: list[list[int]] = []
d = uc.restore_peer(peer_id, "MASTER-A", sys_cfg)
d.addCallback(lambda tgids: out.append(tgids))
assert out == [[7304]]
assert seen["count"] == 1.0
assert isinstance(seen.get("now"), float)

@ -142,6 +142,13 @@ def test_expand_inject_proxy_emits_all_virtual_masters() -> None:
{"parts": {3: "SYSTEM-4", 5: "9990"}},
id="tx_echo_keeps_9990_for_hotspot_rx",
),
pytest.param(
"GROUP VOICE,START,TX,SYSTEM,4100887026,730039101,730039101,2,9990",
[(730039101,)],
{730039101: 4},
{"parts": {3: "SYSTEM-4", 5: "9990"}},
id="tx_echo_peer_id_in_field5_dst_9990",
),
pytest.param(
"GROUP VOICE,START,RX,SYSTEM,4100887026,73003,7300392,2,9990",
[(730039101,)],

@ -388,6 +388,27 @@ def test_register_peer_ua_session_ignores_4000() -> None:
assert peer_single_exclusive_tgid(peer_single, 2, sys_cfg, peer_id=peer_id, now=now) is None
def test_register_peer_ua_session_ignores_echo_9990() -> None:
"""Echo TG 9990 must not create SINGLE=1 exclusive listen lock."""
peer_id = _peer_id()
sys_cfg = _sys_cfg()
now = 1_000_000.0
peer = {"OPTIONS": b"TS2=730444;SINGLE=1;TIMER=5;"}
register_peer_ua_session(peer, peer_id, 2, 730444, sys_cfg, now=now)
register_peer_ua_session(peer, peer_id, 2, 9990, sys_cfg, now=now)
assert peer_single_exclusive_tgid(peer, 2, sys_cfg, peer_id=peer_id, now=now) == 730444
def test_register_peer_ua_session_ignores_service_999x() -> None:
"""On-demand / service 9991–9999 are not UA sessions (same class as echo 9990)."""
peer_id = _peer_id()
sys_cfg = _sys_cfg()
now = 1_000_000.0
peer = {"OPTIONS": b"TS2=730444;SINGLE=1;TIMER=5;"}
register_peer_ua_session(peer, peer_id, 2, 9999, sys_cfg, now=now)
assert peer_single_exclusive_tgid(peer, 2, sys_cfg, peer_id=peer_id, now=now) is None
def test_new_tx_replaces_single_session_tg() -> None:
peer = {"OPTIONS": b"TS2=730,7305;SINGLE=1;TIMER=5;"}
sys_cfg = _sys_cfg()

Loading…
Cancel
Save

Powered by TurnKey Linux.