diff --git a/src/adn_server/infrastructure/options_redaction.py b/src/adn_server/infrastructure/options_redaction.py new file mode 100644 index 0000000..3deedf5 --- /dev/null +++ b/src/adn_server/infrastructure/options_redaction.py @@ -0,0 +1,19 @@ +"""Shared OPTIONS text redaction for logging (mask ``PASS=`` secrets).""" + +from __future__ import annotations + +import re +from typing import Any + +_PASS_RE = re.compile(r"(?i)(PASS=)[^;]*") + + +def redact_pass_in_options(options: Any) -> str: + """OPTIONS text for logging with ``PASS=`` secret replaced by ``PASS=*******``.""" + if options is None: + return "" + if isinstance(options, (bytes, bytearray)): + text = bytes(options).decode("utf-8", errors="replace") + else: + text = str(options) + return _PASS_RE.sub(r"\1*******", text) diff --git a/src/adn_server/infrastructure/proxy/persistence/proxy_repository.py b/src/adn_server/infrastructure/proxy/persistence/proxy_repository.py index 7c45c0d..e781beb 100644 --- a/src/adn_server/infrastructure/proxy/persistence/proxy_repository.py +++ b/src/adn_server/infrastructure/proxy/persistence/proxy_repository.py @@ -81,6 +81,13 @@ class ProxySelfServiceRepository(ProxySelfServiceStore): ).addErrback( lambda f: logger.error("(SELF_SERVICE) psswd: %s", f.getTraceback()) ) + elif action == "clear_psswd": + self._pool.runOperation( + "UPDATE Clients SET psswd=NULL WHERE dmr_id=%s", + (peer_id_bytes,), + ).addErrback( + lambda f: logger.error("(SELF_SERVICE) clear_psswd: %s", f.getTraceback()) + ) elif action == "opt_rcvd": self._pool.runOperation( "UPDATE Clients SET opt_rcvd=1 WHERE dmr_id=%s", diff --git a/src/adn_server/infrastructure/proxy/self_service_bridge.py b/src/adn_server/infrastructure/proxy/self_service_bridge.py index d27885a..e40f48f 100644 --- a/src/adn_server/infrastructure/proxy/self_service_bridge.py +++ b/src/adn_server/infrastructure/proxy/self_service_bridge.py @@ -27,6 +27,7 @@ import struct from hashlib import pbkdf2_hmac from typing import Any +from twisted.internet import reactor from twisted.internet.defer import inlineCallbacks from twisted.internet.interfaces import IDelayedCall from twisted.internet.task import LoopingCall @@ -40,6 +41,9 @@ from adn_server.application.proxy import ProxyUseCases from adn_server.domain.proxy import ClientEndpoint from adn_server.domain.value_objects import int_id from adn_server.infrastructure.hbp_constants import RPTACK, RPTC, RPTCL, RPTO +from adn_server.infrastructure.options_redaction import redact_pass_in_options + +_RPTC_FALLBACK_DELAY = 10.0 def _peer_id_from_db(value: Any) -> bytes | None: @@ -114,6 +118,11 @@ class ProxySelfServiceBridge: command = data[:4] host, port = addr if command == RPTO: + self._log.info( + "(SELF_SERVICE) RPTO from %s:%s peer=%s len=%d payload=%s", + host, port, int_id(peer_id), len(data), + redact_pass_in_options(data[8:8 + 40]), + ) return self._handle_rpto(data, peer_id, host, port) if command == RPTC and len(data) >= 5 and data[:5] != RPTCL: self._handle_rptc(data, peer_id, host) @@ -125,13 +134,22 @@ class ProxySelfServiceBridge: self._store.updt_tbl("log_out", peer_id) def _handle_rptc(self, data: bytes, peer_id: bytes, host: str) -> None: - if self._use_cases.resolve_client(peer_id) is None: + client = self._use_cases.resolve_client(peer_id) + if client is None: + self._log.debug( + "(SELF_SERVICE) RPTC from %s peer=%s but peer NOT in proxy slots", + host, int_id(peer_id), + ) return mode = data[97:98].decode("utf-8", errors="replace") if len(data) >= 98 else "4" callsign = data[8:16].rstrip().decode("utf-8", errors="replace") self._store.ins_conf(int_id(peer_id), peer_id, callsign, host, mode) if peer_id in self._mysql_option_peers: self._fetch_options_now(peer_id) + return + # Schedule a fallback fetch: if the hotspot does not send an RPTO within the + # delay window, treat it the same as an empty RPTO (BD is the authority). + self._schedule_rptc_fallback(peer_id) def _handle_rpto( self, @@ -140,37 +158,101 @@ class ProxySelfServiceBridge: host: str, port: int, ) -> bool: - if self._use_cases.resolve_client(peer_id) is None: + client = self._use_cases.resolve_client(peer_id) + if client is None: + self._log.warning( + "(SELF_SERVICE) RPTO from %s:%s peer=%s but peer NOT in proxy slots — " + "cannot process PASS/OPTIONS", + host, port, int_id(peer_id), + ) return False - if data[8:].upper().startswith(b"PASS=") and len(data) >= 13: - psswd_raw = data[13:] - if len(psswd_raw) >= 6: - dk = pbkdf2_hmac( - "sha256", - psswd_raw, - self._pbkdf2_salt.encode("utf-8"), - self._pbkdf2_iterations, - ).hex() - self._store.updt_tbl("psswd", peer_id, psswd=dk) - self._client_sender.send_to_client( - RPTACK + peer_id, - ClientEndpoint(host=host, port=port), + payload = data[8:].rstrip(b"\x00") + if payload.upper().startswith(b"PASS="): + psswd_raw = payload[5:] + if len(psswd_raw) < 6: + self._log.warning( + "(SELF_SERVICE) RPTO PASS= from %s (peer=%s) too short (%d bytes) — " + "skipped, packet NOT injected to master", + host, int_id(peer_id), len(psswd_raw), ) - self._log.info("(SELF_SERVICE) Password stored for: %s", int_id(peer_id)) - self._mysql_option_peers.add(peer_id) - self._fetch_options_now(peer_id) + return True + dk = pbkdf2_hmac( + "sha256", + psswd_raw, + self._pbkdf2_salt.encode("utf-8"), + self._pbkdf2_iterations, + ).hex() + self._store.updt_tbl("psswd", peer_id, psswd=dk) + self._client_sender.send_to_client( + RPTACK + peer_id, + ClientEndpoint(host=host, port=port), + ) + self._log.info("(SELF_SERVICE) Password stored for: %s", int_id(peer_id)) + self._mysql_option_peers.add(peer_id) + self._fetch_options_now(peer_id) return True - self._mysql_option_peers.discard(peer_id) - self._store.updt_tbl("opt_rcvd", peer_id) - self._cancel_opt_timer(peer_id) - self._log.info("(SELF_SERVICE) Options received from: %s", int_id(peer_id)) - return False + if payload: + # Hotspot sent OPTIONS directly (e.g. TS2=730;SINGLE=0;): it is the authority. + # Clear any stored password so only IP auto-login works (no password login). + self._clear_password(peer_id, host, port) + self._mysql_option_peers.discard(peer_id) + self._store.updt_tbl("opt_rcvd", peer_id) + self._cancel_opt_timer(peer_id) + self._log.info("(SELF_SERVICE) Options received from: %s", int_id(peer_id)) + return False + # Empty payload: BD is the authority. Same treatment as PASS (password cleared). + self._clear_password(peer_id, host, port) + self._mysql_option_peers.add(peer_id) + self._fetch_options_now(peer_id) + self._log.info( + "(SELF_SERVICE) Empty OPTIONS from %s (peer=%s) — BD is authority", + host, int_id(peer_id), + ) + return True + + def _clear_password( + self, peer_id: bytes, host: str, port: int + ) -> None: + """Clear stored password (NULL) when the hotspot does not send PASS. + IP auto-login still works; password login is impossible with a NULL hash.""" + self._store.updt_tbl("clear_psswd", peer_id) + self._client_sender.send_to_client( + RPTACK + peer_id, + ClientEndpoint(host=host, port=port), + ) def _cancel_opt_timer(self, peer_id: bytes) -> None: timer = self._opt_timers.pop(peer_id, None) if timer is not None and timer.active(): timer.cancel() + def _schedule_rptc_fallback(self, peer_id: bytes) -> None: + """If the hotspot does not send RPTO within the delay, fetch OPTIONS from DB. + + Mirrors the legacy proxy 10s ``login_opt`` timer: when no RPTO arrives + after RPTC, the server treats it as "BD is the authority" and fetches + OPTIONS from MySQL so the peer gets its configured static TGs. + """ + self._cancel_opt_timer(peer_id) + timer = reactor.callLater( + _RPTC_FALLBACK_DELAY, self._rptc_fallback_fire, peer_id + ) + self._opt_timers[peer_id] = timer + + def _rptc_fallback_fire(self, peer_id: bytes) -> None: + """Called by the RPTC fallback timer if no RPTO was received.""" + self._opt_timers.pop(peer_id, None) + if self._use_cases.resolve_client(peer_id) is None: + return + if peer_id in self._mysql_option_peers: + return + self._log.info( + "(SELF_SERVICE) No RPTO from peer=%s after %.0fs — fetching DB options", + int_id(peer_id), _RPTC_FALLBACK_DELAY, + ) + self._mysql_option_peers.add(peer_id) + self._fetch_options_now(peer_id) + def _fetch_options_now(self, peer_id: bytes) -> None: """Load OPTIONS from MySQL and inject RPTO to the server (no legacy 10s delay).""" self._cancel_opt_timer(peer_id) @@ -202,15 +284,33 @@ class ProxySelfServiceBridge: def _login_opt(self, peer_id: bytes) -> None: self._opt_timers.pop(peer_id, None) if peer_id not in self._mysql_option_peers: + self._log.debug( + "(SELF_SERVICE) _login_opt skip: peer=%s not in mysql_option_peers", + int_id(peer_id), + ) return if self._use_cases.resolve_client(peer_id) is None: + self._log.warning( + "(SELF_SERVICE) _login_opt skip: peer=%s no longer in proxy slots", + int_id(peer_id), + ) return try: rows = yield self._store.slct_opt(peer_id) if not rows or not rows[0]: + self._log.info( + "(SELF_SERVICE) _login_opt: peer=%s — no options in DB; " + "server YAML defaults will apply", + int_id(peer_id), + ) return options = rows[0][0] if not options: + self._log.info( + "(SELF_SERVICE) _login_opt: peer=%s — options column empty in DB; " + "server YAML defaults will apply", + int_id(peer_id), + ) return self._inject_rpto(peer_id, options) self._log.info("(SELF_SERVICE) Options sent at login for: %s", int_id(peer_id)) @@ -226,9 +326,19 @@ class ProxySelfServiceBridge: continue pid = _peer_id_from_db(row[0]) options = row[1] - if pid is None or not options: + if pid is None: + continue + if not options: + self._log.debug( + "(SELF_SERVICE) send_opts: peer=%s has modified=1 but OPTIONS is EMPTY in DB", + int_id(pid), + ) continue if self._use_cases.resolve_client(pid) is None: + self._log.debug( + "(SELF_SERVICE) send_opts: peer=%s has modified=1 but is NOT connected", + int_id(pid), + ) continue if pid not in self._mysql_option_peers: continue diff --git a/src/adn_server/infrastructure/proxy/udp_fanin.py b/src/adn_server/infrastructure/proxy/udp_fanin.py index 1c24093..90622ae 100644 --- a/src/adn_server/infrastructure/proxy/udp_fanin.py +++ b/src/adn_server/infrastructure/proxy/udp_fanin.py @@ -31,12 +31,15 @@ from twisted.internet.protocol import DatagramProtocol from adn_server.application.ports import ProxyMasterSink from adn_server.application.proxy import ProxyUseCases, peer_id_from_packet from adn_server.domain.result import is_fail +from adn_server.infrastructure.hbp_constants import RPTC, RPTO if TYPE_CHECKING: from .self_service_bridge import ProxySelfServiceBridge _logger = logging.getLogger(__name__) +_CONTROL_COMMANDS = frozenset({RPTC, RPTO}) + class _DatagramWriter(Protocol): def write(self, data: bytes, addr: tuple[str, int]) -> None: @@ -86,13 +89,14 @@ class ProxyFanInProtocol(DatagramProtocol): new_session = self._proxy.resolve_client(peer_id) is None result = self._proxy.attach_client(peer_id, host, port) if is_fail(result): - if self.debug: - self._log.debug( - "(PROXY) attach rejected peer=%s from %s:%s: %s", + if self.debug or command in _CONTROL_COMMANDS: + self._log.warning( + "(PROXY) attach rejected peer=%s from %s:%s: %s (cmd=%r)", peer_id.hex(), host, port, result.error, + command, ) return if self._on_attached is not None: diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 5c377aa..04d2e7b 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -129,6 +129,7 @@ from ..mesh.obp_v1 import ( verify_bcve, ) from ..mesh.registry import MeshCodecRegistry +from ..options_redaction import redact_pass_in_options logger = logging.getLogger(__name__) @@ -1521,7 +1522,8 @@ class HBPProtocol(DatagramProtocol): invalidate_peer_options_cache(_this_peer) self._mark_downlink_index_dirty() self.send_peer(_peer_id, b"".join([RPTACK, _peer_id])) - logger.info("(%s) Peer %s has sent options %s", self._system, _this_peer["CALLSIGN"], _this_peer["OPTIONS"]) + _opt_log = redact_pass_in_options(_this_peer["OPTIONS"]) + logger.info("(%s) Peer %s has sent options %s", self._system, _this_peer["CALLSIGN"], _opt_log) # Inject-only multi-hotspot: OPTIONS live on each peer; do not let last RPTO # overwrite the shared SYSTEM row (legacy had one peer per virtual master). if not is_proxy_inject_only(self._CONFIG, self._system): @@ -2198,8 +2200,6 @@ class HBPProtocol(DatagramProtocol): self._obp_send_bcsq(_dst_id, _stream_id) self._laststrid.append(_stream_id) return - if _call_type == "group" and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD: - logger.info("(%s) CALL RX (OBP DMRE) src %s -> TG %s slot %s", self._system, int_id(_rf_src), int_id(_dst_id), _slot) self.note_dmrd_stream(_peer_id, _rf_src, _stream_id) if ( _call_type in ("group", "vcsbk") diff --git a/tests/infrastructure/test_options_redaction.py b/tests/infrastructure/test_options_redaction.py new file mode 100644 index 0000000..d5e3981 --- /dev/null +++ b/tests/infrastructure/test_options_redaction.py @@ -0,0 +1,47 @@ +# ADN DMR Peer Server - tests infrastructure options redaction +# +# 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 +############################################################################### + +"""Redaction of PASS= secret in the OPTIONS log line.""" + +from __future__ import annotations + +from adn_server.infrastructure.options_redaction import redact_pass_in_options + + +def test_redact_pass_masks_secret() -> None: + assert redact_pass_in_options(b"TS2=730;PASS=secret123;SINGLE=1;") == ( + "TS2=730;PASS=*******;SINGLE=1;" + ) + + +def test_redact_pass_case_insensitive() -> None: + assert redact_pass_in_options(b"pass=hunter2;TS1=730;") == "pass=*******;TS1=730;" + + +def test_redact_pass_no_pass_returns_as_is() -> None: + assert redact_pass_in_options(b"TS2=730;SINGLE=1;") == "TS2=730;SINGLE=1;" + + +def test_redact_pass_accepts_str() -> None: + assert redact_pass_in_options("TS2=730;PASS=pw;") == "TS2=730;PASS=*******;" + + +def test_redact_pass_none_returns_empty() -> None: + assert redact_pass_in_options(None) == "" diff --git a/tests/infrastructure/test_proxy_self_service.py b/tests/infrastructure/test_proxy_self_service.py index 2ce6763..4ecdd15 100644 --- a/tests/infrastructure/test_proxy_self_service.py +++ b/tests/infrastructure/test_proxy_self_service.py @@ -135,15 +135,44 @@ def test_self_service_settings_reads_database_and_self_service_keys() -> None: assert ss["pbkdf2_iterations"] == 2000 -def test_rptc_login_opt_skips_without_pass() -> None: - bridge, sink, store, _sender = _bridge() +def test_rptc_schedules_fallback_when_no_pass() -> None: + """RPTC without prior PASS schedules a fallback timer to fetch DB options.""" + bridge, _sink, store, _sender = _bridge() peer = bytes_4(7300444) store.options_by_peer[peer] = "TS2=730444;" packet = RPTC + peer + b"CE1ILI " + b"\x00" * 85 + b"4" bridge.before_inject(packet, ("192.168.1.10", 62031), peer) assert ("ins_conf", peer) in store.actions - _run_deferred(bridge._login_opt(peer)) - assert sink.injected == [] + assert peer in bridge._opt_timers + + +def test_rptc_fallback_fetches_db_options() -> None: + """When the fallback timer fires (no RPTO), DB options are fetched and injected.""" + bridge, sink, store, _sender = _bridge() + peer = bytes_4(7300444) + store.options_by_peer[peer] = "TS2=730444;SINGLE=1;" + packet = RPTC + peer + b"CE1ILI " + b"\x00" * 85 + b"4" + bridge.before_inject(packet, ("192.168.1.10", 62031), peer) + # Simulate the timer firing + bridge._rptc_fallback_fire(peer) + assert peer in bridge._mysql_option_peers + assert sink.injected + assert sink.injected[0][0] == RPTO + peer + b"TS2=730444;SINGLE=1;" + + +def test_rpto_cancels_rptc_fallback() -> None: + """RPTO with content cancels the RPTC fallback timer.""" + bridge, _sink, _store, _sender = _bridge() + peer = bytes_4(7300444) + packet = RPTC + peer + b"CE1ILI " + b"\x00" * 85 + b"4" + bridge.before_inject(packet, ("192.168.1.10", 62031), peer) + assert peer in bridge._opt_timers + bridge.before_inject( + RPTO + peer + b"TS2=730;SINGLE=0;", + ("192.168.1.10", 62031), + peer, + ) + assert peer not in bridge._opt_timers def test_rptc_after_pass_reinjects_rpto() -> None: @@ -192,7 +221,58 @@ def test_rpto_pass_fetches_options_immediately() -> None: assert sink.injected[0][0] == RPTO + peer + b"TS2=730444;SINGLE=1;" +def test_rpto_empty_fetches_options_from_db() -> None: + """Empty RPTO payload: BD is the authority (like PASS). Password cleared.""" + bridge, sink, store, sender = _bridge() + peer = bytes_4(7300444) + store.options_by_peer[peer] = "TS2=730444;SINGLE=1;" + skip = bridge.before_inject( + RPTO + peer, + ("192.168.1.10", 62031), + peer, + ) + assert skip is True + assert ("clear_psswd", peer) in store.actions + assert sender.sent and sender.sent[0][0][:6] == b"RPTACK" + assert sink.injected + assert sink.injected[0][0] == RPTO + peer + b"TS2=730444;SINGLE=1;" + + +def test_rpto_empty_no_db_options_skips_inject() -> None: + """Empty RPTO + empty BD: no inject; server YAML defaults apply.""" + bridge, sink, store, sender = _bridge() + peer = bytes_4(7300444) + skip = bridge.before_inject( + RPTO + peer, + ("192.168.1.10", 62031), + peer, + ) + assert skip is True + assert ("clear_psswd", peer) in store.actions + assert sender.sent and sender.sent[0][0][:6] == b"RPTACK" + assert sink.injected == [] + + +def test_rpto_with_content_clears_password_and_passes_through() -> None: + """OPTIONS with content (no PASS): hotspot is authority; password cleared + so only IP auto-login works; user cannot login by password (NULL hash).""" + bridge, sink, store, sender = _bridge() + peer = bytes_4(7300444) + skip = bridge.before_inject( + RPTO + peer + b"TS2=730;SINGLE=0;", + ("192.168.1.10", 62031), + peer, + ) + assert skip is False + assert ("clear_psswd", peer) in store.actions + assert ("opt_rcvd", peer) in store.actions + assert sender.sent and sender.sent[0][0][:6] == b"RPTACK" + assert sink.injected == [] + assert peer not in bridge._mysql_option_peers + + def test_send_opts_skips_without_pass() -> None: + """Peer that sent OPTIONS with content (not in _mysql_option_peers) is skipped by send_opts.""" bridge, sink, store, _sender = _bridge() peer = bytes_4(7300444) store.pending_modified = [(peer, "TS2=730444;")]