Merge pull request #59 from ce5rpy/fix/silent-tg-log-once

fix: log silent dynamic TG activation once per peer and TG
pull/60/head
ce5rpy 2 months ago committed by GitHub
commit f807709f2d
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -102,6 +102,36 @@ class HbpForwardMixin:
cache.add(key)
logger.warning(msg, *args)
def _silent_activation_log_key(
self,
system_name: str,
peer_id: bytes,
dst_id: bytes,
) -> tuple:
return ("silent_activation", system_name, bytes_4(int_id(peer_id)), dst_id)
def _log_ingress_info_once(
self,
key: tuple,
msg: str,
*args: object,
) -> None:
cache = self._ingress_drop_log_cache()
if key in cache:
return
cache.add(key)
logger.info(msg, *args)
def _clear_silent_activation_log(
self,
system_name: str,
peer_id: bytes,
dst_id: bytes,
) -> None:
self._ingress_drop_log_cache().discard(
self._silent_activation_log_key(system_name, peer_id, dst_id),
)
def _clear_ingress_drop_log(
self,
system_name: str,
@ -174,12 +204,16 @@ class HbpForwardMixin:
if tg_has_active_conversation(
protocols, systems_cfg, dst_id, stream_id, rf_src, pkt_time,
):
logger.info(
self._log_ingress_info_once(
self._silent_activation_log_key(system_name, peer_id, dst_id),
"(%s) TG %s has active QSO — activating dynamic TG silently for peer %s (uplink suppressed)",
system_name, int_id(dst_id), int_id(peer_id),
)
_slot_st["_suppress_uplink"] = True
_slot_st["_silent_activation_tg"] = int_id(dst_id)
# Bind stream early so voice frames do not re-enter the new-stream
# path and spam silent-activation logs (udp_hbp parity).
_slot_st["RX_STREAM_ID"] = stream_id
else:
self._log_ingress_warning_once(
self._ingress_drop_key(

@ -328,6 +328,8 @@ class RoutingUseCases(
slot_st_vterm = st.get(slot, {})
if isinstance(slot_st_vterm, dict):
_suppress_uplink_vterm = bool(slot_st_vterm.get("_suppress_uplink"))
if _suppress_uplink_vterm:
self._clear_silent_activation_log(system_name, peer_id, dst_id)
if not _suppress_uplink_vterm:
_rx_report_peer = peer_id
if not source_is_obp:

@ -22,6 +22,8 @@
from __future__ import annotations
import logging
import pytest
from tests.harness.assertions import assert_forwarded, assert_inject_ok
from tests.harness.deterministic import DeterministicScenario, PacketSpec, active_routing_table
@ -87,6 +89,47 @@ def test_hbp_stream_collision_silent_activation_on_busy_tg() -> None:
assert scenario.protocols["MASTER-A"].STATUS[2].get("_suppress_uplink") is True
def test_hbp_silent_activation_logs_once_per_peer_tg(caplog: pytest.LogCaptureFixture) -> None:
"""Connect-PTT style bursts may open many stream IDs; log silent activation once."""
bridges = active_routing_table(91, (("MASTER-A", 2), ("MASTER-B", 2)))
scenario = DeterministicScenario(routing_table=bridges)
t0 = scenario.clock.time()
slot = scenario.protocols["MASTER-A"].STATUS[2]
slot.update(
{
"RX_STREAM_ID": bytes_4(0x80808080),
"RX_TYPE": HBPF_SLT_VHEAD,
"RX_RFS": bytes_3(1111111),
"RX_TGID": bytes_3(91),
"RX_TIME": t0,
"RX_START": t0,
}
)
peer = 2222222
streams = (0x70707070, 0x71717171, 0x72727272)
with caplog.at_level(logging.INFO):
for sid in streams:
base = PacketSpec(dst_id=91, stream_id=sid, rf_src=peer)
ok = scenario.inject_hbp(
"MASTER-A",
DeterministicScenario.voice_head_spec(base),
ingress_pkt_time=t0 + 0.1,
)
assert_inject_ok(ok)
burst = DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1)
assert_inject_ok(
scenario.inject_hbp("MASTER-A", burst, ingress_pkt_time=t0 + 0.12)
)
hits = [
r.message
for r in caplog.records
if "activating dynamic TG silently" in r.message
]
assert len(hits) == 1
@pytest.mark.behavior
def test_hbp_collision_allows_same_subscriber_rekey() -> None:
"""Regression: same RF source may start a new stream while prior call is still open."""

Loading…
Cancel
Save

Powered by TurnKey Linux.