From 8d941adf20105c798d9674aa0d01ff80f3fa4a23 Mon Sep 17 00:00:00 2001 From: agessaman Date: Mon, 14 Sep 2026 21:31:11 -0700 Subject: [PATCH] fix(path): authenticate channel RF correlation (#255) --- modules/meshcore_payload_decode.py | 14 +- modules/message_handler.py | 200 +++++++++++++++++ tests/test_message_handler.py | 350 ++++++++++++++++++++++++++++- 3 files changed, 560 insertions(+), 4 deletions(-) diff --git a/modules/meshcore_payload_decode.py b/modules/meshcore_payload_decode.py index f90e90e..d9c32fb 100644 --- a/modules/meshcore_payload_decode.py +++ b/modules/meshcore_payload_decode.py @@ -161,8 +161,10 @@ class ChannelKeyStore: def decrypt_group_text(ciphertext: bytes, cipher_mac: bytes, key16: bytes) -> Optional[dict[str, Any]]: """Verify+decrypt a GRP_TXT ciphertext with a single channel key. - Returns ``{timestamp, flags, sender, text}`` on success, or ``None`` if the - MAC fails or the plaintext is malformed. + Returns ``{timestamp, flags, sender, text, message}`` on success, or + ``None`` if the MAC fails or the plaintext is malformed. ``message`` is + the complete decoded text while ``text`` retains the existing + sender-prefix-stripped value. """ if len(ciphertext) < 16 or len(ciphertext) % 16 != 0: return None @@ -202,7 +204,13 @@ def decrypt_group_text(ciphertext: bytes, cipher_mac: bytes, key16: bytes) -> Op sender = candidate content = text[colon + 2:] - return {"timestamp": timestamp, "flags": flags, "sender": sender, "text": content} + return { + "timestamp": timestamp, + "flags": flags, + "sender": sender, + "text": content, + "message": text, + } def _iso_utc(unix_ts: int) -> Optional[str]: diff --git a/modules/message_handler.py b/modules/message_handler.py index 7699891..e41c42b 100644 --- a/modules/message_handler.py +++ b/modules/message_handler.py @@ -14,6 +14,7 @@ from typing import Any, TypedDict from .enums import AdvertFlags, DeviceRole, PayloadType, PayloadVersion, RouteType from .graph_trace_helper import update_mesh_graph_from_trace_data +from .meshcore_payload_decode import channel_hash_for_key, decrypt_group_text from .models import MeshMessage from .neighbors_discovery import upsert_zero_hop_observed_path_via_manager from .security_utils import sanitize_input, sanitize_name @@ -36,6 +37,7 @@ RF_MATCH_PARTIAL = "partial" # prefix. Channel messages have no prefix to match on, so this is the only positive # correlation available to them (see _rf_data_matches_chan_payload). RF_MATCH_PAYLOAD = "payload" +RF_MATCH_CHANNEL_AUTHENTICATED = "channel_authenticated" RF_MATCH_FALLBACK = "fallback" @@ -76,6 +78,11 @@ class MessageHandler: # Time-based cache for recent RF log data self.recent_rf_data: list[dict[str, Any]] = [] + # Authenticated channel rows outlive the short best-effort RF window. + # A queued CHANNEL_MSG_RECV can be delayed while the mesh stays busy. + self.channel_rf_data: list[dict[str, Any]] = [] + self._channel_rf_cache_timeout = max(300.0, self.rf_data_timeout * 4) + # Message correlation system to prevent race conditions self.pending_messages: dict[str, PendingMessageEntry] = {} # Store messages waiting for RF data @@ -344,6 +351,176 @@ class MessageHandler: # If we can't parse timestamp, process the message (safer to process than skip) return False + @staticmethod + def _channel_message_identity( + channel_idx: Any, sender_timestamp: Any, txt_type: Any, text: Any + ) -> str | None: + """Return the identity shared by an authenticated RF row and CHAN event.""" + try: + idx = int(channel_idx) + timestamp = int(sender_timestamp) + message_type = int(txt_type) + except (TypeError, ValueError): + return None + if not 0 <= idx <= 0xFFFF or not 0 <= timestamp <= 0xFFFFFFFF: + return None + if not 0 <= message_type <= 0xFF or not isinstance(text, str): + return None + + digest = sha256() + digest.update(b"meshcore-channel-message-v1\0") + digest.update(idx.to_bytes(2, "little")) + digest.update(timestamp.to_bytes(4, "little")) + digest.update(bytes([message_type])) + digest.update(text.rstrip("\x00").encode("utf-8")) + return digest.hexdigest() + + def _channel_secrets(self) -> list[tuple[int, bytes]]: + """Return configured channel indexes and keys without logging secrets.""" + meshcore = getattr(self.bot, "meshcore", None) + channels = getattr(meshcore, "channels", None) + if isinstance(channels, dict): + items = channels.items() + elif isinstance(channels, list): + items = enumerate(channels) + else: + return [] + + result: list[tuple[int, bytes]] = [] + for fallback_idx, channel in items: + if not isinstance(channel, dict): + continue + try: + channel_idx = int(channel.get("channel_idx", fallback_idx)) + except (TypeError, ValueError): + continue + + secret = channel.get("channel_secret") + if isinstance(secret, str): + try: + secret = bytes.fromhex(secret) + except ValueError: + secret = None + if not isinstance(secret, (bytes, bytearray)): + key_hex = channel.get("channel_key_hex") + if isinstance(key_hex, str): + try: + secret = bytes.fromhex(key_hex) + except ValueError: + secret = None + if isinstance(secret, bytearray): + secret = bytes(secret) + if isinstance(secret, bytes) and len(secret) == 16 and any(secret): + result.append((channel_idx, secret)) + return result + + def _decode_authenticated_channel_identity( + self, packet_info: dict[str, Any] | None + ) -> dict[str, Any] | None: + """Authenticate a decoded GRP_TXT packet and derive its CHAN identity.""" + if not packet_info or packet_info.get("payload_type") != self._grp_txt_payload_type_int(): + return None + payload_hex = packet_info.get("payload_hex") + if not isinstance(payload_hex, str): + return None + try: + group_payload = bytes.fromhex(payload_hex) + except ValueError: + return None + if len(group_payload) < 3: + return None + + channel_hash = f"{group_payload[0]:02x}" + cipher_mac = group_payload[1:3] + ciphertext = group_payload[3:] + for channel_idx, secret in self._channel_secrets(): + if channel_hash_for_key(secret) != channel_hash: + continue + decrypted = decrypt_group_text(ciphertext, cipher_mac, secret) + if not decrypted: + continue + txt_type = int(decrypted["flags"]) >> 2 + identity = self._channel_message_identity( + channel_idx, + decrypted["timestamp"], + txt_type, + decrypted["message"], + ) + if identity is None: + return None + return { + "channel_message_id": identity, + "channel_idx": channel_idx, + "channel_attempt": int(decrypted["flags"]) & 0x03, + } + return None + + def _cache_authenticated_channel_rf_data( + self, + rf_data: dict[str, Any], + packet_info: dict[str, Any] | None, + current_time: float, + ) -> None: + identity = self._decode_authenticated_channel_identity(packet_info) + if identity is None: + return + rf_data.update(identity) + self.channel_rf_data.append(rf_data) + cutoff = current_time - self._channel_rf_cache_timeout + self.channel_rf_data = [ + row for row in self.channel_rf_data if row.get("timestamp", 0) >= cutoff + ] + if len(self.channel_rf_data) > self._max_rf_cache_size: + self.channel_rf_data = self.channel_rf_data[-self._max_rf_cache_size :] + + def _find_authenticated_channel_rf_data( + self, payload: dict[str, Any] | None + ) -> tuple[dict[str, Any] | None, bool]: + """Find the first RF reception for an authenticated channel message. + + The companion logs raw RF before duplicate suppression and queues only + the first unseen channel packet. Repeater echoes share the packet hash, + so the earliest row is the reception represented by CHANNEL_MSG_RECV. + The boolean prevents weaker matching after an authenticated ambiguity. + """ + if not payload: + return None, False + identity = self._channel_message_identity( + payload.get("channel_idx"), + payload.get("sender_timestamp"), + payload.get("txt_type"), + payload.get("text"), + ) + if identity is None: + return None, False + now = time.time() + max_age = getattr( + self, + "_channel_rf_cache_timeout", + max(300.0, getattr(self, "rf_data_timeout", 15.0) * 4), + ) + matches = [ + row + for row in getattr(self, "channel_rf_data", []) + if row.get("channel_message_id") == identity + and isinstance(row.get("timestamp"), (int, float)) + and now - row["timestamp"] <= max_age + ] + if not matches: + return None, False + + packet_hashes = {row.get("packet_hash") for row in matches} + if len(packet_hashes) != 1 or not all(packet_hashes): + self.logger.warning( + "Authenticated channel message matched %d distinct or unidentified packets; " + "leaving RF attribution unresolved", + len(packet_hashes), + ) + return None, True + + selected = min(matches, key=lambda row: row.get("timestamp", float("inf"))) + return {**selected, RF_MATCH_KEY: RF_MATCH_CHANNEL_AUTHENTICATED}, True + async def handle_contact_message(self, event: Any, metadata: dict[str, Any] | None = None) -> None: """Handle incoming contact message (DM). @@ -1300,6 +1477,9 @@ class MessageHandler: if _lib_pkt_hex else (decoded_packet.get("payload_hex") if decoded_packet else None), } + self._cache_authenticated_channel_rf_data( + rf_data, decoded_packet, current_time + ) if rf_data.get("route_type_int") == 0: self.logger.debug( "TC_FLOOD scope fields: tc_code1=%s payload_type=%s payload_hex_prefix=%s", @@ -1500,6 +1680,19 @@ class MessageHandler: extended_timeout: float, ) -> dict[str, Any] | None: """Correlate a channel message with cached RF log rows (strategies 1–4).""" + authenticated, authenticated_identity_seen = self._find_authenticated_channel_rf_data(payload) + if authenticated_identity_seen: + if authenticated is None: + return None + if scope_eligible_only and not self._is_rf_data_scope_eligible(authenticated): + return None + self.logger.debug( + "Authenticated channel message matched first RF reception %s (packet %s)", + (authenticated.get("packet_prefix") or "?")[:16], + authenticated.get("packet_hash") or "?", + ) + return authenticated + recent_rf_data: dict[str, Any] | None = None if message_packet_prefix: @@ -1516,6 +1709,13 @@ class MessageHandler: message_id = f"{correlation_key}_{int(time.time() * 1000)}" self.store_message_for_correlation(message_id, payload) await asyncio.sleep(0.1) + authenticated, authenticated_identity_seen = self._find_authenticated_channel_rf_data(payload) + if authenticated_identity_seen: + if authenticated is None: + return None + if scope_eligible_only and not self._is_rf_data_scope_eligible(authenticated): + return None + return authenticated recent_rf_data = self.correlate_message_with_rf_data(message_id) if not recent_rf_data: diff --git a/tests/test_message_handler.py b/tests/test_message_handler.py index 602a71c..10078cf 100644 --- a/tests/test_message_handler.py +++ b/tests/test_message_handler.py @@ -8,6 +8,7 @@ from unittest.mock import AsyncMock, Mock, patch import pytest from modules.message_handler import ( + RF_MATCH_CHANNEL_AUTHENTICATED, RF_MATCH_EXACT, RF_MATCH_FALLBACK, RF_MATCH_KEY, @@ -930,6 +931,47 @@ def _make_packet_hex( return pkt.hex() +def _make_group_text_packet( + secret: bytes, + sender_timestamp: int, + text: str, + *, + path_bytes: bytes = b"", + bytes_per_hop: int = 1, + route_type: int = 1, + transport: bytes = b"", + attempt: int = 0, +) -> tuple[str, bytes]: + """Build an authenticated GRP_TXT packet and its application payload.""" + import hashlib + import hmac + + from cryptography.hazmat.primitives.ciphers import Cipher, algorithms, modes + + plaintext = ( + sender_timestamp.to_bytes(4, "little") + + bytes([attempt & 0x03]) + + text.encode("utf-8") + ) + plaintext += b"\x00" * ((-len(plaintext)) % 16) + encryptor = Cipher(algorithms.AES(secret), modes.ECB()).encryptor() + ciphertext = encryptor.update(plaintext) + encryptor.finalize() + cipher_mac = hmac.new(secret + b"\x00" * 16, ciphertext, hashlib.sha256).digest()[:2] + channel_hash = hashlib.sha256(secret).digest()[:1] + group_payload = channel_hash + cipher_mac + ciphertext + hop_count = len(path_bytes) // bytes_per_hop + packet_hex = _make_packet_hex( + 5, + route_type, + path_bytes=path_bytes, + payload_bytes=group_payload, + hop_count=hop_count, + bytes_per_hop=bytes_per_hop, + transport=transport, + ) + return packet_hex, group_payload + + # --------------------------------------------------------------------------- # decode_meshcore_packet # --------------------------------------------------------------------------- @@ -1641,6 +1683,307 @@ class TestHandleRfLogData: assert entry["transport_code1"] is None +class TestAuthenticatedChannelCorrelation: + """Issue #255: correlate CHAN events by authenticated message identity.""" + + SECRET_1 = bytes.fromhex("eb50a1bcb3e4e5d7bf69a57c9dada211") + SECRET_2 = bytes.fromhex("8b3387e9c5cdea6ac9e5edbaa115cd72") + SENDER_TIMESTAMP = 1_788_043_003 + TEXT = "CoderNemesis-KY 🏷️: !test" + + @staticmethod + def _setup(handler, channels): + handler.logger = Mock() + handler.bot.meshcore = Mock() + handler.bot.meshcore.channels = channels + handler.bot.transmission_tracker = None + handler.bot.web_viewer_integration = None + + @staticmethod + async def _store_rf( + handler, + secret, + *, + path_bytes=b"", + bytes_per_hop=1, + snr=11.75, + rssi=-30, + route_type=1, + transport=b"", + attempt=0, + text=TEXT, + sender_timestamp=SENDER_TIMESTAMP, + ): + packet_hex, group_payload = _make_group_text_packet( + secret, + sender_timestamp, + text, + path_bytes=path_bytes, + bytes_per_hop=bytes_per_hop, + route_type=route_type, + transport=transport, + attempt=attempt, + ) + event = Mock() + payload = { + "snr": snr, + "rssi": rssi, + "raw_hex": "0000" + packet_hex, + "payload": packet_hex, + "payload_length": len(bytes.fromhex(packet_hex)), + "route_type": route_type, + "payload_type": 5, + "pkt_payload": group_payload, + } + if transport: + payload["transport_code"] = transport.hex() + event.payload = payload + await handler.handle_rf_log_data(event) + return handler.channel_rf_data[-1] + + @classmethod + def _chan(cls, channel_idx=1, **over): + payload = { + "type": "CHAN", + "SNR": 0.0, + "channel_idx": channel_idx, + "path_hash_mode": 0, + "path_len": 0, + "txt_type": 0, + "sender_timestamp": cls.SENDER_TIMESTAMP, + "text": cls.TEXT, + } + payload.update(over) + return payload + + @pytest.mark.asyncio + async def test_zero_snr_uses_authenticated_rf_row(self, handler): + self._setup( + handler, + {1: {"channel_idx": 1, "channel_secret": self.SECRET_1}}, + ) + stored = await self._store_rf(handler, self.SECRET_1) + + result = await handler._correlate_channel_message_rf_data( + None, + "", + self._chan(SNR=0.0), + scope_eligible_only=False, + extended_timeout=30.0, + ) + + assert result[RF_MATCH_KEY] == RF_MATCH_CHANNEL_AUTHENTICATED + assert result["packet_hash"] == stored["packet_hash"] + assert result["snr"] == 11.75 + + @pytest.mark.asyncio + async def test_issue_255_zero_snr_capture_resolves_direct_path(self, handler): + """CoderNemesis26's Aug 29 capture carries the exact message identity.""" + self._setup( + handler, + {1: {"channel_idx": 1, "channel_secret": self.SECRET_1}}, + ) + inner_packet = ( + "1540ca6609ada0b960f785625ad4267e43de66426f5680f9d91e7c78c48c" + "d2081440d26c083603ac6faac17069ecedfd66aa2ca9aa" + ) + packet_info = handler.decode_meshcore_packet(inner_packet) + event = Mock() + event.payload = { + "snr": 11.75, + "rssi": 0, + "raw_hex": "2f00" + inner_packet, + "payload": inner_packet, + "payload_length": len(bytes.fromhex(inner_packet)), + "route_type": 1, + "payload_type": 5, + "pkt_payload": bytes.fromhex(packet_info["payload_hex"]), + } + await handler.handle_rf_log_data(event) + + result = await handler._correlate_channel_message_rf_data( + None, + "", + { + "type": "CHAN", + "SNR": 0.0, + "channel_idx": 1, + "path_hash_mode": 1, + "path_len": 0, + "txt_type": 0, + "sender_timestamp": 1_788_042_903, + "text": "CoderNemesis-KY 🏷️: !test", + }, + scope_eligible_only=False, + extended_timeout=30.0, + ) + + assert result[RF_MATCH_KEY] == RF_MATCH_CHANNEL_AUTHENTICATED + assert result["packet_hash"] == "82F374D969AF26A3" + assert result["routing_info"]["path_length"] == 0 + assert result["routing_info"]["path_nodes"] == [] + + @pytest.mark.asyncio + async def test_first_reception_wins_over_repeater_echo(self, handler): + self._setup( + handler, + {1: {"channel_idx": 1, "channel_secret": self.SECRET_1}}, + ) + first = await self._store_rf(handler, self.SECRET_1, snr=11.75, rssi=-30) + echo = await self._store_rf( + handler, + self.SECRET_1, + path_bytes=bytes.fromhex("f0"), + snr=-4.0, + rssi=-90, + ) + assert first["packet_hash"] == echo["packet_hash"] + + result = await handler._correlate_channel_message_rf_data( + None, + "", + self._chan(SNR=0.0), + scope_eligible_only=False, + extended_timeout=30.0, + ) + + assert result["routing_info"]["path_length"] == 0 + assert result["routing_info"]["path_nodes"] == [] + assert result["snr"] == 11.75 + assert result["rssi"] == -30 + + @pytest.mark.asyncio + async def test_authenticated_cache_outlives_general_rf_window(self, handler): + self._setup( + handler, + {1: {"channel_idx": 1, "channel_secret": self.SECRET_1}}, + ) + stored = await self._store_rf(handler, self.SECRET_1) + stored["timestamp"] = time.time() - 60 + handler.recent_rf_data = [] + + result = await handler._correlate_channel_message_rf_data( + None, + "", + self._chan(), + scope_eligible_only=False, + extended_timeout=30.0, + ) + + assert result[RF_MATCH_KEY] == RF_MATCH_CHANNEL_AUTHENTICATED + + @pytest.mark.asyncio + async def test_authenticated_cache_still_expires(self, handler): + self._setup( + handler, + {1: {"channel_idx": 1, "channel_secret": self.SECRET_1}}, + ) + stored = await self._store_rf(handler, self.SECRET_1) + stored["timestamp"] = time.time() - handler._channel_rf_cache_timeout - 1 + handler.recent_rf_data = [] + + result = await handler._correlate_channel_message_rf_data( + None, + "", + self._chan(), + scope_eligible_only=False, + extended_timeout=30.0, + ) + + assert result is None + + def test_tampered_ciphertext_is_not_authenticated(self, handler): + self._setup( + handler, + {1: {"channel_idx": 1, "channel_secret": self.SECRET_1}}, + ) + packet_hex, _group_payload = _make_group_text_packet( + self.SECRET_1, self.SENDER_TIMESTAMP, self.TEXT + ) + packet_info = handler.decode_meshcore_packet(packet_hex) + payload = bytearray.fromhex(packet_info["payload_hex"]) + payload[-1] ^= 0x01 + packet_info["payload_hex"] = payload.hex() + + assert handler._decode_authenticated_channel_identity(packet_info) is None + + @pytest.mark.asyncio + async def test_channel_index_disambiguates_same_timestamp_and_text(self, handler): + self._setup( + handler, + { + 1: {"channel_idx": 1, "channel_secret": self.SECRET_1}, + 2: {"channel_idx": 2, "channel_secret": self.SECRET_2}, + }, + ) + first = await self._store_rf(handler, self.SECRET_1) + second = await self._store_rf( + handler, + self.SECRET_2, + path_bytes=bytes.fromhex("aabb"), + bytes_per_hop=2, + snr=-7.5, + ) + + result = await handler._correlate_channel_message_rf_data( + None, + "", + self._chan(channel_idx=2, path_hash_mode=1, path_len=1), + scope_eligible_only=False, + extended_timeout=30.0, + ) + + assert result["packet_hash"] == second["packet_hash"] + assert result["packet_hash"] != first["packet_hash"] + assert result["routing_info"]["path_nodes"] == ["AABB"] + + @pytest.mark.asyncio + async def test_distinct_authenticated_packets_are_ambiguous(self, handler): + self._setup( + handler, + {1: {"channel_idx": 1, "channel_secret": self.SECRET_1}}, + ) + first = await self._store_rf(handler, self.SECRET_1, attempt=0) + second = await self._store_rf(handler, self.SECRET_1, attempt=1) + assert first["packet_hash"] != second["packet_hash"] + + result = await handler._correlate_channel_message_rf_data( + None, + "", + self._chan(), + scope_eligible_only=False, + extended_timeout=30.0, + ) + + assert result is None + + @pytest.mark.asyncio + async def test_scope_lookup_uses_same_authenticated_first_reception(self, handler): + self._setup( + handler, + {1: {"channel_idx": 1, "channel_secret": self.SECRET_1}}, + ) + transport = bytes.fromhex("12340000") + stored = await self._store_rf( + handler, + self.SECRET_1, + route_type=0, + transport=transport, + ) + + result = await handler._correlate_channel_message_rf_data( + None, + "", + self._chan(), + scope_eligible_only=True, + extended_timeout=30.0, + ) + + assert result[RF_MATCH_KEY] == RF_MATCH_CHANNEL_AUTHENTICATED + assert result["packet_hash"] == stored["packet_hash"] + assert result["transport_code1"] == 0x3412 + + # --------------------------------------------------------------------------- # _get_path_from_rf_data # --------------------------------------------------------------------------- @@ -2175,7 +2518,12 @@ class TestRfCorrelationProvenance: into the mesh graph.""" def test_correlated_matches_are_attributable(self): - for kind in (RF_MATCH_EXACT, RF_MATCH_PUBKEY, RF_MATCH_PARTIAL): + for kind in ( + RF_MATCH_EXACT, + RF_MATCH_PUBKEY, + RF_MATCH_PARTIAL, + RF_MATCH_CHANNEL_AUTHENTICATED, + ): assert rf_data_is_correlated({RF_MATCH_KEY: kind}) is True def test_fallback_is_not_attributable(self):