Merge pull request #96 from Amateur-Digital-Network/fix/obp-ingress-reused-stream-id

fix(obp): let a reused stream_id start a new call instead of dropping it
pull/100/head
ce5rpy 6 days ago committed by GitHub
commit 5f78545809
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -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"]:

@ -0,0 +1,118 @@
# ADN DMR Peer Server - tests routing obp reused stream id
#
# 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
###############################################################################
"""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
Loading…
Cancel
Save

Powered by TurnKey Linux.