Merge pull request #29 from ce5rpy/develop

fix: self-service RPTO handling, PASS redaction, and proxy log cleanup
pull/30/head
ce5rpy 3 months ago committed by GitHub
commit 7f6b57bd26
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -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)

@ -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",

@ -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

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

@ -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")

@ -0,0 +1,47 @@
# ADN DMR Peer Server - tests infrastructure options redaction
#
# 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
###############################################################################
"""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) == ""

@ -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;")]

Loading…
Cancel
Save

Powered by TurnKey Linux.