fix(path): authenticate channel RF correlation (#255)

This commit is contained in:
agessaman
2026-09-14 21:31:11 -07:00
parent c41501eb48
commit 8d941adf20
3 changed files with 560 additions and 4 deletions
+11 -3
View File
@@ -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]:
+200
View File
@@ -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:
+349 -1
View File
@@ -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):