diff --git a/src/adn_server/application/routing/hbp_forward.py b/src/adn_server/application/routing/hbp_forward.py index 4c14230..6d36e12 100644 --- a/src/adn_server/application/routing/hbp_forward.py +++ b/src/adn_server/application/routing/hbp_forward.py @@ -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( diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index 8fee588..a9e2391 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -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: diff --git a/tests/hbp/test_timeout_collision.py b/tests/hbp/test_timeout_collision.py index 9a7d718..b961677 100644 --- a/tests/hbp/test_timeout_collision.py +++ b/tests/hbp/test_timeout_collision.py @@ -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."""