From c2fd82698dd8065987a5ad5723421178ed63247f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Tue, 22 Sep 2026 17:33:30 -0300 Subject: [PATCH] fix(obp): let a reused stream_id start a new call instead of dropping it STATUS is keyed by stream_id alone and the trimmer only drops a row after 180s idle, so when a peer reuses an id for another destination the second call lands on the first one's row, fails the TGID check in loop control and is discarded frame by frame. to_target already evicts a stale row on a forward leg; ingress did not. Ingress now evicts it too, once the old row has been idle for 2s so two genuinely interleaved streams keep their own rows. Also drops the bogus TS from the loop-control warning (it printed a slot index from an unrelated loop), gives that branch the once-per-stream guard its siblings have, and says what the condition is. Seen on a production bridge: one talkgroup took 116 frames and completed no calls at all. --- .../application/routing/obp_forward.py | 40 ++++-- tests/routing/test_obp_reused_stream_id.py | 118 ++++++++++++++++++ 2 files changed, 148 insertions(+), 10 deletions(-) create mode 100644 tests/routing/test_obp_reused_stream_id.py diff --git a/src/adn_server/application/routing/obp_forward.py b/src/adn_server/application/routing/obp_forward.py index 526fe08..e959594 100644 --- a/src/adn_server/application/routing/obp_forward.py +++ b/src/adn_server/application/routing/obp_forward.py @@ -56,6 +56,10 @@ from .helpers import group_voice_tg_ingress_collision, obp_is_canonical_ingress, logger = logging.getLogger(__name__) +# Idle time before another destination may take over a stream_id. The trimmer only +# drops a row after 180s; well above one voice frame so interleaved streams keep theirs. +OBP_REUSED_STREAM_IDLE_S = 2.0 + class ObpForwardMixin: """routerOBP group voice, stream tracking, sendDataToOBP.""" @@ -170,6 +174,23 @@ class ObpForwardMixin: if status is None: return True + # A peer reusing a stream_id for another destination is starting a new call, + # not continuing the one the row describes. to_target already evicts on a + # forward leg; without the same here loop control refuses the new call. + _prev = status.get(stream_id) + if isinstance(_prev, dict) and _prev.get("TGID") != dst_id: + _idle = pkt_time - _prev.get("LAST", _prev.get("START", pkt_time)) + if _idle >= OBP_REUSED_STREAM_IDLE_S: + logger.info( + "(%s) stream %s reused for TG %s (was TG %s, idle %.1fs): starting a new call", + system_name, + int_id(stream_id), + int_id(dst_id), + int_id(_prev.get("TGID", b"\x00\x00\x00")), + _idle, + ) + del status[stream_id] + if stream_id not in status: if group_voice_tg_ingress_collision( protocols, systems_cfg, dst_id, stream_id, rf_src, pkt_time, @@ -266,7 +287,6 @@ class ObpForwardMixin: # Legacy routerOBP ~2409: LoopControl only on 2nd+ packet (else branch). hr_times: dict[str, float] = {} - _sysslot_last = 0 for other_name, proto in (protocols or {}).items(): omode = systems_cfg.get(other_name, {}).get("MODE") if other_name != system_name and omode != "OPENBRIDGE": @@ -274,7 +294,6 @@ class ObpForwardMixin: if not ostatus: continue for _sysslot in ostatus: - _sysslot_last = _sysslot if isinstance(_sysslot, int) else _sysslot_last slot_st = ostatus.get(_sysslot) if not isinstance(slot_st, dict): continue @@ -303,15 +322,16 @@ class ObpForwardMixin: hr_times[other_name] = obp_status[stream_id]["1ST"] fi = min(hr_times, key=hr_times.get, default=False) - hr_times.clear() if not fi: - logger.warning( - "(%s) OBP *LoopControl* fi is empty for some reason : STREAM ID: %s, TG: %s, TS: %s", - system_name, - int_id(stream_id), - int_id(dst_id), - _sysslot_last, - ) + if not st.get("LOOPLOG"): + logger.warning( + "(%s) OBP *LoopControl* no bridge holds stream %s for TG %s, dropping", + system_name, + int_id(stream_id), + int_id(dst_id), + ) + st["LOOPLOG"] = True + st["LAST"] = pkt_time return False if system_name != fi: if "LOOPLOG" not in st or not st["LOOPLOG"]: diff --git a/tests/routing/test_obp_reused_stream_id.py b/tests/routing/test_obp_reused_stream_id.py new file mode 100644 index 0000000..21d5540 --- /dev/null +++ b/tests/routing/test_obp_reused_stream_id.py @@ -0,0 +1,118 @@ +# ADN DMR Peer Server - tests routing obp reused stream id +# +# 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 +############################################################################### + +"""A peer that reuses a stream_id for another destination. + +Seen in production: one bridge sent stream 128 to three destinations in a day and +another reused one id across five. STATUS is keyed by stream_id alone and the +trimmer only drops a row after 180s idle, so the second destination landed on the +first one's row, failed the TGID check in loop control and was discarded frame by +frame — 116 frames to one talkgroup that never got a single call through. +to_target already evicts a stale row on a forward leg; ingress did not. +""" + +from __future__ import annotations + +import pytest +from tests.harness.assertions import assert_forwarded, packets_to +from tests.harness.deterministic import ( + DeterministicScenario, + PacketSpec, + add_openbridge_system, + minimal_config, + patch_routing_wall_time, +) + +from adn_server.application.routing.obp_forward import OBP_REUSED_STREAM_IDLE_S + +_STREAM = 128 # as the peer really sends it +_TG_FIRST = 71401 +_TG_SECOND = 71411 + + +def _bridge_row(system: str, tgid: int) -> dict: + from adn_server.domain import bytes_3 + + tg_b = bytes_3(tgid) + return { + "SYSTEM": system, "TS": 1, "TGID": tg_b, "ACTIVE": True, + "TIMEOUT": 3600.0, "TO_TYPE": "ON", "ON": [tg_b], "OFF": [], "RESET": [], "TIMER": 0.0, + } + + +def _scenario() -> DeterministicScenario: + config = minimal_config(()) + add_openbridge_system(config, "OBP-IN") + add_openbridge_system(config, "OBP-OUT") + bridges = { + str(_TG_FIRST): [_bridge_row("OBP-IN", _TG_FIRST), _bridge_row("OBP-OUT", _TG_FIRST)], + str(_TG_SECOND): [_bridge_row("OBP-IN", _TG_SECOND), _bridge_row("OBP-OUT", _TG_SECOND)], + } + scenario = DeterministicScenario(config=config, routing_table=bridges) + scenario.routing._finalize_routing_state() + return scenario + + +@pytest.mark.behavior +def test_a_reused_stream_id_reaches_its_new_destination() -> None: + """The production timeline: TG 71401 goes quiet, and 45s later TG 71411 arrives + on the same stream id. It used to be dropped for the rest of the 180s window. + """ + scenario = _scenario() + with patch_routing_wall_time(scenario.clock): + first = PacketSpec(dst_id=_TG_FIRST, stream_id=_STREAM, slot=1) + scenario.inject_obp("OBP-IN", DeterministicScenario.voice_head_spec(first)) + assert_forwarded(scenario, "OBP-OUT", count=1, dst_id=_TG_FIRST) + + scenario.clock.advance(45.0) + second = PacketSpec(dst_id=_TG_SECOND, stream_id=_STREAM, slot=1) + scenario.inject_obp("OBP-IN", DeterministicScenario.voice_head_spec(second)) + + got = packets_to(scenario, "OBP-OUT") + from adn_server.domain import int_id + + assert [int_id(p.fields["dst_id"]) for p in got] == [_TG_FIRST, _TG_SECOND] + + +@pytest.mark.behavior +def test_an_interleaved_stream_id_does_not_evict_the_live_call() -> None: + """Without the idle rule each frame would drop the other call's row, and both + would lose their packet counters, LC and loss tracking. + """ + scenario = _scenario() + with patch_routing_wall_time(scenario.clock): + first = PacketSpec(dst_id=_TG_FIRST, stream_id=_STREAM, slot=1) + scenario.inject_obp("OBP-IN", DeterministicScenario.voice_head_spec(first)) + + second = PacketSpec(dst_id=_TG_SECOND, stream_id=_STREAM, slot=1) + for _ in range(8): + scenario.clock.advance(0.06) # one voice frame apart + scenario.inject_obp("OBP-IN", DeterministicScenario.voice_head_spec(second)) + + status = scenario.protocols["OBP-IN"].STATUS + from adn_server.domain import bytes_4, int_id + + row = status.get(bytes_4(_STREAM)) + assert row is not None + assert int_id(row["TGID"]) == _TG_FIRST + + +def test_the_idle_threshold_is_well_above_one_voice_frame() -> None: + assert OBP_REUSED_STREAM_IDLE_S > 0.06