From cc4031ae4d887cd2dbcfdc29b0ccaa8fd57dd72c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Mon, 13 Jul 2026 01:59:08 -0400 Subject: [PATCH] fix: stop premature OBP monitor voice END events Restrict stream trimmer END,RX/END,TX to ingress rows only so BCSQ-quenched forward legs no longer cascade END to still-active peers; drop stale forward STATUS before START,TX on stream_id reuse. --- src/adn_server/application/routing/timers.py | 29 +- .../application/routing_use_cases.py | 11 + .../application/test_obp_monitor_voice_end.py | 254 ++++++++++++++++++ .../test_voice_subscription_plan.py | 2 +- 4 files changed, 290 insertions(+), 6 deletions(-) create mode 100644 tests/application/test_obp_monitor_voice_end.py diff --git a/src/adn_server/application/routing/timers.py b/src/adn_server/application/routing/timers.py index e9e0141..260df00 100644 --- a/src/adn_server/application/routing/timers.py +++ b/src/adn_server/application/routing/timers.py @@ -53,6 +53,16 @@ from .helpers import obp_clear_deferred_bridge_tx_leg logger = logging.getLogger(__name__) +def _obp_status_is_forward_leg(st: dict[str, Any]) -> bool: + """``to_target`` OPENBRIDGE rows carry rewritten LC (monitor TX leg).""" + return isinstance(st, dict) and "H_LC" in st + + +def _obp_status_is_ingress(st: dict[str, Any]) -> bool: + """Ingress/routerOBP source stream (monitor canonical RX).""" + return isinstance(st, dict) and not _obp_status_is_forward_leg(st) + + class RoutingTimerMixin: """rule_timer, stream_trimmer, bridge_reset, stat_trimmer, bridge_debug loops.""" @@ -182,7 +192,9 @@ class RoutingTimerMixin: return if tst.get("TGID", b"\x00\x00\x00") != tgid: return - self._obp_emit_end_tx_forward_leg(system_name, stream_id, tst, time.time()) + now = time.time() + self._obp_emit_end_tx_forward_leg(system_name, stream_id, tst, now) + tst["_bcsq_quenched"] = True def flush_monitor_events_for_system(self, system_name: str, protocol: Any) -> None: if not self._config.get("REPORTS", {}).get("REPORT", True): @@ -267,8 +279,13 @@ class RoutingTimerMixin: if st.get("_to") and last < now - 180: to_remove.append(stream_id) continue - # Stage 1: 5s idle, not yet timed out → mark _to, emit END + # Stage 1: 5s idle — monitor END only on ingress (canonical RX). + # Forward legs (H_LC): no END,RX and no cross-leg END,TX cascade. if "_to" not in st and "_fin" not in st and last < now - 5: + if _obp_status_is_forward_leg(st): + if st.get("_end_tx_sent") or st.get("_bcsq_quenched"): + st["_to"] = True + continue rfs = st.get("RFS", b"\x00\x00\x00") peer = st.get("RX_PEER", b"\x00\x00\x00\x00") tgid = st.get("TGID", b"\x00\x00\x00") @@ -280,8 +297,6 @@ class RoutingTimerMixin: ) ) st["_to"] = True - # Legacy trimmer emits END,RX here only; forward legs waited ~180s. - # END,TX now (END_TX_FORWARD): same event as legacy VTERM path (~2039). self._obp_emit_end_tx_for_forward_legs(stream_id, system_name, now) continue for stream_id in to_remove: @@ -291,7 +306,11 @@ class RoutingTimerMixin: for _tgid_k, _sid in list(_bmap.items()): if _sid == stream_id: _bmap.pop(_tgid_k, None) - self._obp_emit_end_tx_for_forward_legs(stream_id, system_name, now) + st_rem = obp_status.get(stream_id) + if isinstance(st_rem, dict) and _obp_status_is_forward_leg(st_rem): + self._obp_emit_end_tx_forward_leg(system_name, stream_id, st_rem, now) + elif isinstance(st_rem, dict) and _obp_status_is_ingress(st_rem): + self._obp_emit_end_tx_for_forward_legs(stream_id, system_name, now) obp_status.pop(stream_id, None) continue for slot in (1, 2): diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index 33df780..2d9c5ea 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -509,6 +509,17 @@ class RoutingUseCases( if isinstance(_prev_obp_st, dict) and _prev_obp_st.get("_fin"): del _target_status[stream_id] self.clear_talker_alias_stream(system_name, stream_id) + _prev_obp_st = None + # Reused stream_id on a forward leg: drop stale row before START,TX. + if isinstance(_prev_obp_st, dict) and "H_LC" in _prev_obp_st: + _stale_forward = ( + _prev_obp_st.get("_end_tx_sent") + or _prev_obp_st.get("RX_PEER") != peer_id + or _prev_obp_st.get("RFS") != rf_src + or _prev_obp_st.get("TGID") != dst_id_b + ) + if _stale_forward: + del _target_status[stream_id] if stream_id not in _target_status: _target_status[stream_id] = { "START": pkt_time, diff --git a/tests/application/test_obp_monitor_voice_end.py b/tests/application/test_obp_monitor_voice_end.py new file mode 100644 index 0000000..1a7c26a --- /dev/null +++ b/tests/application/test_obp_monitor_voice_end.py @@ -0,0 +1,254 @@ +# ADN DMR Peer Server - OBP monitor voice END reporting +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""OBP monitor voice END,RX / END,TX report timing (no premature cascade).""" + +from __future__ import annotations + +from contextlib import contextmanager +from typing import Any + +from tests.harness.deterministic import ( + FakeClock, + FakeObpProtocol, + FakeReportFactory, + FakeReportSender, + add_openbridge_system, + minimal_config, +) + +from adn_server.application.reporting_use_cases import ReportingUseCases +from adn_server.application.routing_use_cases import RoutingUseCases +from adn_server.domain import bytes_3, bytes_4 +from adn_server.domain.dmr.bptc import encode_emblc +from adn_server.infrastructure.acl_router import InMemoryAclRouter +from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore +from adn_server.infrastructure.talker_alias_emblc import default_ta_emblc_encoder + + +def _event_system(event: str) -> str: + parts = event.split(",") + return parts[3] if len(parts) > 3 else "" + + +def _end_tx_for(events: list[str], system: str) -> list[str]: + return [e for e in events if ",END,TX," in e and _event_system(e) == system] + + +def _end_rx_for(events: list[str], system: str) -> list[str]: + return [e for e in events if ",END,RX," in e and _event_system(e) == system] + + +@contextmanager +def _patch_wall_time(clock: FakeClock): + import adn_server.application.routing.timers as timers_mod + import adn_server.application.routing_use_cases as buc + + orig_t = timers_mod.time.time + orig_b = buc.time.time + timers_mod.time.time = clock.time + buc.time.time = clock.time + try: + yield + finally: + timers_mod.time.time = orig_t + buc.time.time = orig_b + + +def _obp_stack( + names: tuple[str, ...], +) -> tuple[RoutingUseCases, FakeReportFactory, FakeClock, dict[str, FakeObpProtocol]]: + config = minimal_config() + for name in names: + add_openbridge_system(config, name) + config["REPORTS"]["REPORT"] = True + clock = FakeClock() + report_factory = FakeReportFactory() + protocols = {name: FakeObpProtocol(name) for name in names} + routing = RoutingUseCases( + InMemoryAclRouter(), + config, + InMemorySubscriptionStore(), + get_protocols=lambda: protocols, + reporting=ReportingUseCases(FakeReportSender(report_factory), config), + encode_emblc=encode_emblc, + ta_emblc_encoder=default_ta_emblc_encoder, + ) + return routing, report_factory, clock, protocols + + +def _forward_leg( + *, + peer_id: bytes, + rf_src: bytes, + tgid: bytes, + start: float, + last: float, + **extra: Any, +) -> dict[str, Any]: + return { + "H_LC": b"\x01", + "RX_PEER": peer_id, + "RFS": rf_src, + "TGID": tgid, + "START": start, + "LAST": last, + **extra, + } + + +def _ingress_leg( + *, + peer_id: bytes, + rf_src: bytes, + tgid: bytes, + start: float, + last: float, + **extra: Any, +) -> dict[str, Any]: + return { + "RX_PEER": peer_id, + "RFS": rf_src, + "TGID": tgid, + "START": start, + "LAST": last, + "_monitor_canonical_rx": True, + **extra, + } + + +def test_bcsq_emits_end_tx_only_on_quenched_leg() -> None: + routing, factory, clock, protocols = _obp_stack(("OBP-USA", "OBP-ES", "OBP-CU")) + stream_id = bytes_4(0x1199AABB) + tgid = bytes_3(52090) + now = clock.time() + protocols["OBP-ES"].STATUS[stream_id] = _forward_leg( + peer_id=bytes_4(71411), + rf_src=bytes_3(3120001), + tgid=tgid, + start=now, + last=now, + ) + + with _patch_wall_time(clock): + routing.on_obp_bcsq_received("OBP-ES", tgid, stream_id) + + assert len(_end_tx_for(factory.events, "OBP-ES")) == 1 + assert _end_tx_for(factory.events, "OBP-CU") == [] + assert _end_tx_for(factory.events, "OBP-USA") == [] + assert protocols["OBP-ES"].STATUS[stream_id].get("_bcsq_quenched") is True + assert protocols["OBP-ES"].STATUS[stream_id].get("_end_tx_sent") is True + + +def test_forward_leg_idle_does_not_cascade_end_tx() -> None: + routing, factory, clock, protocols = _obp_stack(("OBP-USA", "OBP-ES", "OBP-CU")) + stream_id = bytes_4(0x297999345) + tgid = bytes_3(52090) + now = clock.time() + protocols["OBP-USA"].STATUS[stream_id] = _ingress_leg( + peer_id=bytes_4(31031), + rf_src=bytes_3(3120001), + tgid=tgid, + start=now, + last=now, + ) + protocols["OBP-ES"].STATUS[stream_id] = _forward_leg( + peer_id=bytes_4(71411), + rf_src=bytes_3(3120001), + tgid=tgid, + start=now - 6, + last=now - 6, + _end_tx_sent=True, + _bcsq_quenched=True, + ) + protocols["OBP-CU"].STATUS[stream_id] = _forward_leg( + peer_id=bytes_4(31031), + rf_src=bytes_3(3120001), + tgid=tgid, + start=now, + last=now, + ) + + with _patch_wall_time(clock): + routing.stream_trimmer_loop() + + assert _end_rx_for(factory.events, "OBP-ES") == [] + assert _end_tx_for(factory.events, "OBP-CU") == [] + assert _end_tx_for(factory.events, "OBP-ES") == [] + assert _end_rx_for(factory.events, "OBP-USA") == [] + assert protocols["OBP-ES"].STATUS[stream_id].get("_to") is True + + +def test_ingress_idle_emits_end_rx_and_active_forward_end_tx_only() -> None: + routing, factory, clock, protocols = _obp_stack(("OBP-USA", "OBP-ES", "OBP-CU")) + stream_id = bytes_4(0x297999345) + tgid = bytes_3(52090) + now = clock.time() + protocols["OBP-USA"].STATUS[stream_id] = _ingress_leg( + peer_id=bytes_4(31031), + rf_src=bytes_3(3120001), + tgid=tgid, + start=now - 10, + last=now - 6, + ) + protocols["OBP-ES"].STATUS[stream_id] = _forward_leg( + peer_id=bytes_4(71411), + rf_src=bytes_3(3120001), + tgid=tgid, + start=now - 10, + last=now - 1, + _end_tx_sent=True, + ) + protocols["OBP-CU"].STATUS[stream_id] = _forward_leg( + peer_id=bytes_4(31031), + rf_src=bytes_3(3120001), + tgid=tgid, + start=now - 10, + last=now - 1, + ) + + with _patch_wall_time(clock): + routing.stream_trimmer_loop() + + assert len(_end_rx_for(factory.events, "OBP-USA")) == 1 + assert _end_tx_for(factory.events, "OBP-ES") == [] + assert len(_end_tx_for(factory.events, "OBP-CU")) == 1 + usa_event = _end_rx_for(factory.events, "OBP-USA")[0] + assert usa_event.endswith("4.00") or usa_event.endswith("4.0") + + +def test_bcsq_end_tx_not_duplicated() -> None: + routing, factory, clock, protocols = _obp_stack(("OBP-ES",)) + stream_id = bytes_4(0xAABBCCDD) + tgid = bytes_3(52090) + now = clock.time() + protocols["OBP-ES"].STATUS[stream_id] = _forward_leg( + peer_id=bytes_4(71411), + rf_src=bytes_3(3120001), + tgid=tgid, + start=now, + last=now, + _end_tx_sent=True, + ) + + with _patch_wall_time(clock): + routing.on_obp_bcsq_received("OBP-ES", tgid, stream_id) + + assert factory.events == [] diff --git a/tests/application/test_voice_subscription_plan.py b/tests/application/test_voice_subscription_plan.py index 98133d2..9406fe9 100644 --- a/tests/application/test_voice_subscription_plan.py +++ b/tests/application/test_voice_subscription_plan.py @@ -32,8 +32,8 @@ 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 fakes.subscription_store import InMemorySubscriptionStore from adn_server.infrastructure.talker_alias_emblc import default_ta_emblc_encoder +from fakes.subscription_store import InMemorySubscriptionStore def _routing(routing_table: dict) -> RoutingUseCases: