diff --git a/src/adn_server/infrastructure/proxy/obp_fanin.py b/src/adn_server/infrastructure/proxy/obp_fanin.py index 44b7136..2c3de3e 100644 --- a/src/adn_server/infrastructure/proxy/obp_fanin.py +++ b/src/adn_server/infrastructure/proxy/obp_fanin.py @@ -181,6 +181,7 @@ class ObpFanInDemux: self._registry = registry self.debug = debug self._log = logger or _logger + self._last_stream: dict[str, bytes] = {} def deliver( self, @@ -210,14 +211,28 @@ class ObpFanInDemux: if entry is None: return if self.debug: - self._log.debug( - "(OBP_PROXY) RX %s from %s:%s len=%d -> %s", - data[:4], - host, - port, - len(data), - system_name, - ) + opcode = data[:4] + if opcode in (DMRD, DMRE) and len(data) >= 20: + stream_id = data[16:20] + if self._last_stream.get(system_name) != stream_id: + self._last_stream[system_name] = stream_id + self._log.debug( + "(OBP_PROXY) RX %s from %s:%s stream=%s -> %s", + opcode, + host, + port, + stream_id.hex(), + system_name, + ) + else: + self._log.debug( + "(OBP_PROXY) RX %s from %s:%s len=%d -> %s", + opcode, + host, + port, + len(data), + system_name, + ) entry.reply_transport.note_ingress(transport) entry.sink.inject(data, addr) diff --git a/tests/infrastructure/test_obp_proxy.py b/tests/infrastructure/test_obp_proxy.py index 15cfc0b..f750be9 100644 --- a/tests/infrastructure/test_obp_proxy.py +++ b/tests/infrastructure/test_obp_proxy.py @@ -114,7 +114,7 @@ class _FakeReactor: return _FakeUdpPort(port, protocol) -def _sample_dmr_voice() -> bytes: +def _sample_dmr_voice(stream_id: int = 0xAABBCCDD) -> bytes: return b"".join( [ DMRD, @@ -123,7 +123,7 @@ def _sample_dmr_voice() -> bytes: bytes_4(52090)[1:4], bytes_4(1), bytes([0x10]), - bytes_4(0xAABBCCDD), + bytes_4(stream_id), b"\x00" * 33, ] ) @@ -568,3 +568,33 @@ def test_control_frame_prefers_configured_peer_over_relaxed_target() -> None: demux.deliver(build_bcka(_PASS), peer_pt, local_port=62032, transport=_RecordingTransport()) assert protocols["OBP-PT"].packets assert not protocols["OBP-FR"].packets + + +def test_debug_rx_log_once_per_stream(caplog) -> None: + """A call sends many DMRD/DMRE packets; the RX debug line must fire once per + stream_id, not once per packet, or an active call floods the log.""" + import logging as _logging + + receiver = _RecordingObp() + transport = _RecordingTransport() + registry = ObpBridgeRegistry() + registry.register( + ObpBridgeEntry( + system_name="OBP-CL", + network_id=_NETWORK, + passphrase=_PASS, + sink=InProcessObpSink(receiver), + reply_transport=ObpIngressReplyTransport(transport), + ) + ) + demux = ObpFanInDemux(registry, debug=True) + caplog.set_level(_logging.DEBUG) + + wire_a = build_dmrd_v1(_sample_dmr_voice(stream_id=0x11111111), _NETWORK, _PASS) + for _ in range(5): + demux.deliver(wire_a, _ADDR, local_port=62032, transport=transport) + assert sum("RX" in r.getMessage() for r in caplog.records) == 1 + + wire_b = build_dmrd_v1(_sample_dmr_voice(stream_id=0x22222222), _NETWORK, _PASS) + demux.deliver(wire_b, _ADDR, local_port=62032, transport=transport) + assert sum("RX" in r.getMessage() for r in caplog.records) == 2