From 873801c1b49c18402ce9f9c8d56567c7751a2acf Mon Sep 17 00:00:00 2001 From: comms_engineer Date: Fri, 3 Jul 2026 00:41:46 +0000 Subject: [PATCH] refactor: extract shared utilities from duplicated MeshCore interface code Create _meshcore_shared.py containing common logic that was duplicated between MeshCore_Channel_Interface and MeshCore_Dynamic_Interface: - load_meshcore(): meshcore library import with error handling - start_asyncio_loop_thread(): event loop thread creation - create_meshcore_connection(): serial/ble/tcp transport init - configure_radio() / configure_channel(): radio+channel setup - process_incoming(): standard RNS inbound delivery - decode_tunnel_frame(): sender-prefix stripping + base64 decode - reassemble_fragment(): thread-safe fragment reassembly - cleanup_stale_fragments(): timed fragment eviction - PacketIdCounter: thread-safe rolling packet ID counter MeshCore_Interface.py is left unchanged (architecturally different, older fragmentation format, different threading model). Import mechanism uses sys.path with __file__ fallback to handle Reticulum's exec()-based interface file loading. Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- Interface/MeshCore_Channel_Interface.py | 345 +++++++----------------- Interface/MeshCore_Dynamic_Interface.py | 175 +++++------- Interface/_meshcore_shared.py | 286 ++++++++++++++++++++ 3 files changed, 457 insertions(+), 349 deletions(-) create mode 100644 Interface/_meshcore_shared.py diff --git a/Interface/MeshCore_Channel_Interface.py b/Interface/MeshCore_Channel_Interface.py index e9bd894..de76ed7 100644 --- a/Interface/MeshCore_Channel_Interface.py +++ b/Interface/MeshCore_Channel_Interface.py @@ -128,11 +128,34 @@ import base64 import hashlib import json +import os import queue import socket +import sys import threading import time from collections import OrderedDict + +# Shared MeshCore utilities (handles Reticulum's exec()-based file loading) +try: + _iface_dir = os.path.dirname(os.path.abspath(__file__)) +except NameError: + _iface_dir = os.path.expanduser("~/.reticulum/interfaces") +if _iface_dir not in sys.path: + sys.path.insert(0, _iface_dir) + +from _meshcore_shared import ( + load_meshcore, + start_asyncio_loop_thread, + create_meshcore_connection, + configure_radio, + configure_channel, + process_incoming, + decode_tunnel_frame, + reassemble_fragment, + cleanup_stale_fragments, + PacketIdCounter, +) # urllib.error and urllib.request are imported lazily inside the three methods # that use them (_send_via_remoteterm, _rt_get_sync, _rt_post_sync). # Top-level submodule imports (dotted names) trigger a Python 3.13 bug in @@ -281,8 +304,7 @@ def __init__(self, owner, configuration): self._worker_thread = None self._own_src_id = self._derive_local_src_id() - self._pkt_id = 0 - self._pkt_id_lock = threading.Lock() + self._pkt_counter = PacketIdCounter(bits=8) self._outqueue = queue.Queue(maxsize=self.OUTQUEUE_MAXSIZE) @@ -302,12 +324,10 @@ def __init__(self, owner, configuration): self._load_meshcore_or_panic() # ---- start asyncio event loop thread ---- - self._loop = asyncio.new_event_loop() - self._loop_thread = threading.Thread( - target=self._run_loop, daemon=True, - name=f"MCChan-loop-{self.name}" + self._loop, self._loop_thread = start_asyncio_loop_thread( + f"MCChan-loop-{self.name}", + f"MeshCore_Channel_Interface [{self.name}]", ) - self._loop_thread.start() # ---- start outgoing worker thread ---- self._worker_thread = threading.Thread( @@ -336,16 +356,10 @@ def __init__(self, owner, configuration): def _load_meshcore_or_panic(self): try: - import meshcore as _mc_module - self._mc_module = _mc_module - self._EventType = _mc_module.EventType - except ImportError: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - "meshcore library not found. " - "Install with: pip install meshcore [--break-system-packages]", - RNS.LOG_CRITICAL + self._mc_module, self._EventType = load_meshcore( + f"MeshCore_Channel_Interface [{self.name}]" ) + except ImportError: RNS.panic() def _load_websockets_or_panic(self): @@ -373,17 +387,6 @@ def _derive_local_src_id(self) -> bytes: # Event loop management # ----------------------------------------------------------------------- - def _run_loop(self): - asyncio.set_event_loop(self._loop) - try: - self._loop.run_forever() - except Exception as exc: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"event loop crashed: {exc}", - RNS.LOG_ERROR - ) - def _run_coro(self, coro, timeout: float = 20.0): if self._loop is None or not self._loop.is_running(): return None @@ -418,40 +421,21 @@ async def _async_setup(self): # ======================================================================= async def _async_setup_direct(self): - MeshCore = self._mc_module.MeshCore - ET = self._EventType + ET = self._EventType + iname = f"MeshCore_Channel_Interface [{self.name}]" try: - if self.transport == "serial": - self._mc = await MeshCore.create_serial(self.port, self.baudrate) - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"connected via serial {self.port}", RNS.LOG_INFO - ) - elif self.transport == "ble": - self._mc = await MeshCore.create_ble(self.ble_name or None) - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - "connected via BLE", RNS.LOG_INFO - ) - elif self.transport == "tcp": - self._mc = await MeshCore.create_tcp(self.host, self.tcp_port) - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"connected via TCP {self.host}:{self.tcp_port}", RNS.LOG_INFO - ) - else: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"unknown transport '{self.transport}'", - RNS.LOG_CRITICAL - ) - return - except Exception as exc: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"connection failed: {exc}", RNS.LOG_ERROR + self._mc = await create_meshcore_connection( + self._mc_module, self.transport, + port=self.port, baudrate=self.baudrate, + host=self.host, tcp_port=self.tcp_port, + ble_name=self.ble_name, interface_name=iname, ) + except ValueError as exc: + RNS.log(f"{iname}: {exc}", RNS.LOG_CRITICAL) + return + except Exception as exc: + RNS.log(f"{iname}: Connection failed: {exc}", RNS.LOG_ERROR) return try: @@ -470,39 +454,15 @@ async def _async_setup_direct(self): f"send_appstart error: {exc}", RNS.LOG_WARNING ) - if self.radio_freq and self.radio_bw and self.radio_sf and self.radio_cr: - try: - result = await self._mc.commands.set_radio( - self.radio_freq, self.radio_bw, self.radio_sf, self.radio_cr) - if result.type == ET.OK: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"radio set freq={self.radio_freq} bw={self.radio_bw} " - f"sf={self.radio_sf} cr={self.radio_cr}", RNS.LOG_INFO - ) - except Exception as exc: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"radio config error: {exc}", RNS.LOG_WARNING - ) + await configure_radio( + self._mc, self.radio_freq, self.radio_bw, + self.radio_sf, self.radio_cr, interface_name=iname, + ) - try: - secret_bytes = bytes.fromhex(self.channel_secret_hex) - if len(secret_bytes) != 16: - raise ValueError(f"channel_secret must be 16 bytes, got {len(secret_bytes)}") - result = await self._mc.commands.set_channel( - self.channel_idx, self.channel_name, secret_bytes) - if result.type == ET.OK: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"channel {self.channel_idx} ('{self.channel_name}') configured", - RNS.LOG_INFO - ) - except Exception as exc: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"channel config error: {exc}", RNS.LOG_WARNING - ) + await configure_channel( + self._mc, self.channel_idx, self.channel_name, + self.channel_secret_hex, interface_name=iname, + ) # Subscribe WITHOUT attribute_filters — some firmware versions report # all received channel messages with channel_idx=0 regardless of the @@ -804,48 +764,31 @@ async def _process_tunnel_text(self, text: str): Decode, validate header, and reassemble one channel message. All silent drops are logged in debug mode. """ - # MeshCore prepends the sender's node name to channel message text, - # e.g. "Janus39: RNS:..." instead of "RNS:...". - # Find the RNS: marker wherever it appears and strip everything before it. - rns_idx = text.find(self.MSG_PREFIX) - if rns_idx == -1: - if self.debug_level == "debug" and text: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"pipeline drop: no RNS: marker found text={text[:40]}", - RNS.LOG_DEBUG - ) + iname = f"MeshCore_Channel_Interface [{self.name}]" + + # Decode the tunnel frame (strip sender prefix + base64 decode) + raw, sender_prefix = decode_tunnel_frame( + text, self.MSG_PREFIX, interface_name=iname, + ) + if raw is None: + if sender_prefix is None: + if self.debug_level == "debug" and text: + RNS.log( + f"{iname}: " + f"pipeline drop: no RNS: marker found text={text[:40]}", + RNS.LOG_DEBUG + ) return - if rns_idx > 0: + if sender_prefix: RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"stripped sender prefix {text[:rns_idx]!r}", + f"{iname}: " + f"stripped sender prefix {sender_prefix!r}", RNS.LOG_DEBUG ) - text = text[rns_idx:] - - b64 = text[len(self.MSG_PREFIX):].strip() - - b64 += "=" * (-len(b64) % 4) - - try: - - raw = base64.urlsafe_b64decode(b64) - - except Exception as exc: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"base64 decode error: {exc} " - f"len={len(b64)} " - f"mod4={len(b64)%4} " - f"text={text[:80]}", - RNS.LOG_WARNING - ) - return if len(raw) < self.HEADER_SIZE: RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " + f"{iname}: " f"pipeline drop: frame too short ({len(raw)} < {self.HEADER_SIZE})", RNS.LOG_WARNING ) @@ -859,18 +802,12 @@ async def _process_tunnel_text(self, text: str): payload = raw[self.HEADER_SIZE:] if frag_total == 0: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - "invalid fragment count 0", - RNS.LOG_WARNING - ) + RNS.log(f"{iname}: invalid fragment count 0", RNS.LOG_WARNING) return if frag_idx >= frag_total: RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"invalid fragment index " - f"{frag_idx}/{frag_total}", + f"{iname}: invalid fragment index {frag_idx}/{frag_total}", RNS.LOG_WARNING ) return @@ -878,8 +815,7 @@ async def _process_tunnel_text(self, text: str): if magic != self.MAGIC: if self.debug_level == "debug": RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"pipeline drop: bad magic {magic!r}", + f"{iname}: pipeline drop: bad magic {magic!r}", RNS.LOG_DEBUG ) return @@ -887,8 +823,7 @@ async def _process_tunnel_text(self, text: str): if src_id == self._own_src_id: if self.debug_level == "debug": RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"pipeline drop: own echo src={src_id.hex()}", + f"{iname}: pipeline drop: own echo src={src_id.hex()}", RNS.LOG_DEBUG ) return @@ -905,7 +840,7 @@ async def _process_tunnel_text(self, text: str): if self.debug_level == "debug": RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " + f"{iname}: " f"dedupe check " f"src={src_hex} " f"pkt={pkt_id} " @@ -917,7 +852,7 @@ async def _process_tunnel_text(self, text: str): if cache_hit: if self.debug_level == "debug": RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " + f"{iname}: " f"pipeline drop: already delivered " f"(src={src_hex} pkt={pkt_id})", RNS.LOG_DEBUG @@ -925,83 +860,27 @@ async def _process_tunnel_text(self, text: str): return RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " + f"{iname}: " f"RX frag src={src_hex} pkt={pkt_id} " f"{frag_idx+1}/{frag_total} payload={len(payload)}B", RNS.LOG_DEBUG ) - with self._asm_lock: - if key not in self._assembly: - self._assembly[key] = {} - self._assembly_meta[key] = (frag_total, time.monotonic()) - - if frag_idx in self._assembly[key]: - if self.debug_level == "debug": - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"assembly state src={src_hex} " - f"pkt={pkt_id} " - f"stored={len(self._assembly[key])}/{frag_total}", - RNS.LOG_DEBUG - ) - return - - self._assembly[key][frag_idx] = payload - - expected_total = self._assembly_meta[key][0] - if len(self._assembly[key]) < expected_total: - return - - missing = [ - i for i in range(expected_total) - if i not in self._assembly[key] - ] - - if missing: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"reassembly incomplete despite count match " - f"missing={missing}", - RNS.LOG_WARNING - ) - return - - try: - full_packet = b"".join( - self._assembly[key][i] for i in range(expected_total) - ) - - del self._assembly[key] - del self._assembly_meta[key] - - except KeyError as exc: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"reassembly gap — missing fragment {exc}", - RNS.LOG_WARNING - ) - del self._assembly[key] - del self._assembly_meta[key] - return - except Exception as exc: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"reassembly failed: {exc}", - RNS.LOG_ERROR - ) - del self._assembly[key] - del self._assembly_meta[key] - return + full_packet = reassemble_fragment( + self._assembly, self._assembly_meta, self._asm_lock, + key, frag_idx, frag_total, payload, interface_name=iname, + ) + if full_packet is None: + return - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"reassembly complete " - f"src={src_hex} " - f"pkt={pkt_id} " - f"len={len(full_packet)}", - RNS.LOG_DEBUG - ) + RNS.log( + f"{iname}: " + f"reassembly complete " + f"src={src_hex} " + f"pkt={pkt_id} " + f"len={len(full_packet)}", + RNS.LOG_DEBUG + ) with self._seen_lock: self._seen_pkts[key] = time.monotonic() @@ -1011,31 +890,27 @@ async def _process_tunnel_text(self, text: str): self._seen_pkts.popitem(last=False) RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " + f"{iname}: " f"RX reassembled {len(full_packet)}B from src={src_hex}", RNS.LOG_INFO ) try: if len(full_packet) == 0: - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - "empty reassembled packet", - RNS.LOG_WARNING - ) + RNS.log(f"{iname}: empty reassembled packet", RNS.LOG_WARNING) return RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " + f"{iname}: " f"reassembled len={len(full_packet)} " - f"from {expected_total} fragments", + f"from {frag_total} fragments", RNS.LOG_DEBUG ) self.processIncoming(full_packet) RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " + f"{iname}: " f"delivered packet " f"src={src_hex} " f"pkt={pkt_id} " @@ -1045,26 +920,18 @@ async def _process_tunnel_text(self, text: str): except Exception as exc: RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"processIncoming failed: {exc}", + f"{iname}: processIncoming failed: {exc}", RNS.LOG_ERROR ) async def _cleanup_loop(self): + iname = f"MeshCore_Channel_Interface [{self.name}]" while True: await asyncio.sleep(60) - deadline = time.monotonic() - self.fragment_timeout_s - with self._asm_lock: - stale = [k for k, (_, ts) in self._assembly_meta.items() - if ts < deadline] - for k in stale: - del self._assembly[k] - del self._assembly_meta[k] - RNS.log( - f"MeshCore_Channel_Interface [{self.name}]: " - f"evicted stale assembly {k}", - RNS.LOG_WARNING - ) + cleanup_stale_fragments( + self._assembly, self._assembly_meta, self._asm_lock, + self.fragment_timeout_s, interface_name=iname, + ) # ======================================================================= # OUTGOING PATH @@ -1077,9 +944,7 @@ def processOutgoing(self, data): if not self.online: return - with self._pkt_id_lock: - pkt_id = self._pkt_id - self._pkt_id = (self._pkt_id + 1) & 0xFF + pkt_id = self._pkt_counter.next_id() handler = _PacketHandler(data, self._own_src_id, pkt_id) @@ -1253,13 +1118,7 @@ def _rt_post_sync(self, path: str, data: dict) -> dict: # ======================================================================= def processIncoming(self, data: bytes): - # RNS 1.x Interface base class has no processIncoming method. - # The correct pattern (matching TCPClientInterface and all other RNS - # interfaces) is: update rxb, then call owner.inbound() directly. - # super() also does not work in exec()'d files under Python 3.13. - if self.online and not self.detached: - self.rxb += len(data) - self.owner.inbound(data, self) + process_incoming(self, data) def __str__(self): return f"MeshCore_Channel_Interface[{self.name}]" diff --git a/Interface/MeshCore_Dynamic_Interface.py b/Interface/MeshCore_Dynamic_Interface.py index 6c2ccbd..e547f68 100644 --- a/Interface/MeshCore_Dynamic_Interface.py +++ b/Interface/MeshCore_Dynamic_Interface.py @@ -216,12 +216,35 @@ import asyncio import base64 import hashlib +import os import struct import queue import random +import sys import threading import time +# Shared MeshCore utilities (handles Reticulum's exec()-based file loading) +try: + _iface_dir = os.path.dirname(os.path.abspath(__file__)) +except NameError: + _iface_dir = os.path.expanduser("~/.reticulum/interfaces") +if _iface_dir not in sys.path: + sys.path.insert(0, _iface_dir) + +from _meshcore_shared import ( + load_meshcore, + start_asyncio_loop_thread, + create_meshcore_connection, + configure_radio, + configure_channel, + process_incoming, + decode_tunnel_frame, + reassemble_fragment, + cleanup_stale_fragments, + PacketIdCounter, +) + # ───────────────────────────────────────────────────────────────────────────── # Fragmentation helper @@ -402,8 +425,7 @@ def __init__(self, owner, configuration): self._own_node_name = "" self._own_mc_key = "" - self._pkt_id = 0 - self._pkt_id_lock = threading.Lock() + self._pkt_counter = PacketIdCounter(bits=32) self._assembly = {} self._assembly_meta = {} @@ -431,12 +453,10 @@ def __init__(self, owner, configuration): self._setup_done = threading.Event() self._load_meshcore_or_panic() - self._loop = asyncio.new_event_loop() - self._loop_thread = threading.Thread( - target=self._run_loop, daemon=True, - name=f"MCDyn-loop-{self.name}" + self._loop, self._loop_thread = start_asyncio_loop_thread( + f"MCDyn-loop-{self.name}", + f"MeshCore_Dynamic_Interface [{self.name}]", ) - self._loop_thread.start() _setup_future = asyncio.run_coroutine_threadsafe( self._async_setup(), self._loop @@ -478,49 +498,28 @@ def _on_setup_done(fut): def _load_meshcore_or_panic(self): try: - import meshcore as _mc_mod - self._mc_module = _mc_mod - self._EventType = _mc_mod.EventType - except ImportError: - RNS.log( - f"MeshCore_Dynamic_Interface [{self.name}]: " - f"meshcore library not found — cannot continue.", - RNS.LOG_CRITICAL + self._mc_module, self._EventType = load_meshcore( + f"MeshCore_Dynamic_Interface [{self.name}]" ) + except ImportError: self.owner.panic() - def _run_loop(self): - asyncio.set_event_loop(self._loop) - try: - self._loop.run_forever() - except Exception as exc: - RNS.log( - f"MeshCore_Dynamic_Interface [{self.name}]: Loop crashed: {exc}", - RNS.LOG_ERROR - ) - async def _async_setup(self): - MeshCore = self._mc_module.MeshCore - ET = self._EventType + ET = self._EventType + iname = f"MeshCore_Dynamic_Interface [{self.name}]" try: - if self.transport == "serial": - self._mc = await MeshCore.create_serial(self.port, self.baudrate) - elif self.transport == "ble": - self._mc = await MeshCore.create_ble(self.ble_name or None) - elif self.transport == "tcp": - self._mc = await MeshCore.create_tcp(self.host, self.tcp_port) - else: - RNS.log( - f"MeshCore_Dynamic_Interface [{self.name}]: " - f"Unknown transport '{self.transport}'.", RNS.LOG_ERROR - ) - return - except Exception as exc: - RNS.log( - f"MeshCore_Dynamic_Interface [{self.name}]: " - f"Driver init error: {exc}", RNS.LOG_ERROR + self._mc = await create_meshcore_connection( + self._mc_module, self.transport, + port=self.port, baudrate=self.baudrate, + host=self.host, tcp_port=self.tcp_port, + ble_name=self.ble_name, interface_name=iname, ) + except ValueError as exc: + RNS.log(f"{iname}: {exc}", RNS.LOG_ERROR) + return + except Exception as exc: + RNS.log(f"{iname}: Driver init error: {exc}", RNS.LOG_ERROR) return try: @@ -543,24 +542,15 @@ async def _async_setup(self): f"Identity fetch failed: {exc}", RNS.LOG_WARNING ) - if self.radio_freq and self.radio_bw and self.radio_sf and self.radio_cr: - try: - await self._mc.commands.set_radio( - self.radio_freq, self.radio_bw, self.radio_sf, self.radio_cr - ) - except Exception: - pass + await configure_radio( + self._mc, self.radio_freq, self.radio_bw, + self.radio_sf, self.radio_cr, interface_name=iname, + ) - try: - secret_bytes = bytes.fromhex(self.channel_secret_hex) - await self._mc.commands.set_channel( - self.channel_idx, self.channel_name, secret_bytes - ) - except Exception as exc: - RNS.log( - f"MeshCore_Dynamic_Interface [{self.name}]: " - f"Channel init error: {exc}", RNS.LOG_WARNING - ) + await configure_channel( + self._mc, self.channel_idx, self.channel_name, + self.channel_secret_hex, interface_name=iname, + ) if self.allow_direct: self._has_direct_api = hasattr(self._mc.commands, "send_msg") @@ -694,20 +684,16 @@ async def _delayed_bind_response(self): # ------------------------------------------------------------------------- async def _cleanup_loop(self): + iname = f"MeshCore_Dynamic_Interface [{self.name}]" while True: - await asyncio.sleep(30) + await asyncio.sleep(30) now = time.monotonic() # --- Stale fragment buffers ------------------------------------ - frag_deadline = now - self.fragment_timeout_s - with self._asm_lock: - stale = [ - k for k, (_, ts) in self._assembly_meta.items() - if ts < frag_deadline - ] - for k in stale: - del self._assembly[k] - del self._assembly_meta[k] + cleanup_stale_fragments( + self._assembly, self._assembly_meta, self._asm_lock, + self.fragment_timeout_s, interface_name=iname, + ) # --- Expired sliding window deduplication records -------------- with self._seen_lock: @@ -870,14 +856,13 @@ def _resolve_sender_key(self, key_str: str) -> str: return key_str async def _process_tunnel_text(self, text: str, sender: str = "", rx_mode: str = "UNKNOWN"): + iname = f"MeshCore_Dynamic_Interface [{self.name}]" + if sender and sender == self._own_node_name: return - b64 = text[len(self.MSG_PREFIX):].strip() - b64 += "=" * (-len(b64) % 4) - try: - raw = base64.urlsafe_b64decode(b64) - except Exception: + raw, _ = decode_tunnel_frame(text, self.MSG_PREFIX, interface_name=iname) + if raw is None: return if len(raw) < self.HEADER_SIZE: @@ -901,31 +886,13 @@ async def _process_tunnel_text(self, text: str, sender: str = "", rx_mode: str = else: del self._seen_pkts[key] - # Fragment reassembly - with self._asm_lock: - if key not in self._assembly: - self._assembly[key] = {} - self._assembly_meta[key] = (frag_total, now) - - if frag_idx in self._assembly[key]: - return - - self._assembly[key][frag_idx] = payload - - if len(self._assembly[key]) < self._assembly_meta[key][0]: - return - - try: - expected = self._assembly_meta[key][0] - full_packet = b"".join( - self._assembly[key][i] for i in range(expected) - ) - del self._assembly[key] - del self._assembly_meta[key] - except Exception: - self._assembly.pop(key, None) - self._assembly_meta.pop(key, None) - return + # Fragment reassembly (shared utility) + full_packet = reassemble_fragment( + self._assembly, self._assembly_meta, self._asm_lock, + key, frag_idx, frag_total, payload, interface_name=iname, + ) + if full_packet is None: + return # Mark as completely reassembled inside sliding time window with self._seen_lock: @@ -1082,9 +1049,7 @@ def processOutgoing(self, data): # Cooldown expired -- this starts a fresh burst. self._path_req_sent_times[dest_id] = (now, now) - with self._pkt_id_lock: - pkt_id = self._pkt_id - self._pkt_id = (self._pkt_id + 1) & 0xFFFFFFFF # 32-bit bound integer tracking + pkt_id = self._pkt_counter.next_id() handler = _PacketHandler(data, pkt_id, self.payload_size) broadcast = self._is_broadcast_packet(data) @@ -1259,9 +1224,7 @@ async def _async_outgoing_worker(self): # ------------------------------------------------------------------------- def processIncoming(self, data: bytes): - if self.online and not self.detached: - self.rxb += len(data) - self.owner.inbound(data, self) + process_incoming(self, data) def __str__(self): return f"MeshCore_Dynamic_Interface[{self.name}]" diff --git a/Interface/_meshcore_shared.py b/Interface/_meshcore_shared.py new file mode 100644 index 0000000..5bd310a --- /dev/null +++ b/Interface/_meshcore_shared.py @@ -0,0 +1,286 @@ +""" +_meshcore_shared.py — Shared utilities for MeshCore RNS interfaces. + +Place alongside MeshCore_Channel_Interface.py and MeshCore_Dynamic_Interface.py +in ~/.reticulum/interfaces/ (or the repo's Interface/ directory). + +Extracted from duplicated logic that was present in both +MeshCore_Channel_Interface and MeshCore_Dynamic_Interface. +""" + +import RNS +import asyncio +import base64 +import threading +import time + + +# --------------------------------------------------------------------------- +# Library loading +# --------------------------------------------------------------------------- + +def load_meshcore(interface_name): + """ + Import the meshcore library. + + Returns ``(meshcore_module, EventType)`` on success. + Logs CRITICAL and re-raises ``ImportError`` on failure so the caller + can invoke its own panic handler. + """ + try: + import meshcore + return meshcore, meshcore.EventType + except ImportError: + RNS.log( + f"{interface_name}: " + "meshcore library not found. " + "Install with: pip install meshcore", + RNS.LOG_CRITICAL, + ) + raise + + +# --------------------------------------------------------------------------- +# Asyncio event-loop helpers +# --------------------------------------------------------------------------- + +def _run_asyncio_loop(loop, interface_name): + """Thread target that runs an asyncio event loop forever.""" + asyncio.set_event_loop(loop) + try: + loop.run_forever() + except Exception as exc: + RNS.log( + f"{interface_name}: Event loop crashed: {exc}", + RNS.LOG_ERROR, + ) + + +def start_asyncio_loop_thread(thread_name, interface_name): + """ + Create a new asyncio event loop and start it in a daemon thread. + + Returns ``(loop, thread)``. + """ + loop = asyncio.new_event_loop() + thread = threading.Thread( + target=_run_asyncio_loop, + args=(loop, interface_name), + daemon=True, + name=thread_name, + ) + thread.start() + return loop, thread + + +# --------------------------------------------------------------------------- +# MeshCore connection / radio / channel setup +# --------------------------------------------------------------------------- + +async def create_meshcore_connection(mc_module, transport, *, + port="/dev/ttyUSB0", baudrate=115200, + host="127.0.0.1", tcp_port=4403, + ble_name="", interface_name=""): + """ + Create a MeshCore connection for the given transport type. + + Raises ``ValueError`` for unknown transports; propagates driver errors. + """ + MeshCore = mc_module.MeshCore + if transport == "serial": + mc = await MeshCore.create_serial(port, baudrate) + RNS.log(f"{interface_name}: Connected via serial {port}", RNS.LOG_INFO) + elif transport == "ble": + mc = await MeshCore.create_ble(ble_name or None) + RNS.log(f"{interface_name}: Connected via BLE", RNS.LOG_INFO) + elif transport == "tcp": + mc = await MeshCore.create_tcp(host, tcp_port) + RNS.log( + f"{interface_name}: Connected via TCP {host}:{tcp_port}", + RNS.LOG_INFO, + ) + else: + raise ValueError(f"Unknown transport '{transport}'") + return mc + + +async def configure_radio(mc, freq, bw, sf, cr, interface_name=""): + """Apply radio parameter overrides if all four values are non-zero.""" + if not (freq and bw and sf and cr): + return + try: + await mc.commands.set_radio(freq, bw, sf, cr) + RNS.log( + f"{interface_name}: " + f"Radio set freq={freq} bw={bw} sf={sf} cr={cr}", + RNS.LOG_INFO, + ) + except Exception as exc: + RNS.log( + f"{interface_name}: Radio config error: {exc}", + RNS.LOG_WARNING, + ) + + +async def configure_channel(mc, channel_idx, channel_name, channel_secret_hex, + interface_name=""): + """Configure a MeshCore channel from a hex-encoded secret.""" + try: + secret_bytes = bytes.fromhex(channel_secret_hex) + await mc.commands.set_channel(channel_idx, channel_name, secret_bytes) + RNS.log( + f"{interface_name}: " + f"Channel {channel_idx} ('{channel_name}') configured", + RNS.LOG_INFO, + ) + except Exception as exc: + RNS.log( + f"{interface_name}: Channel config error: {exc}", + RNS.LOG_WARNING, + ) + + +# --------------------------------------------------------------------------- +# RNS interface helpers +# --------------------------------------------------------------------------- + +def process_incoming(interface, data): + """ + Standard RNS inbound delivery used by MeshCore interfaces. + + Updates the byte counter and hands the packet to the RNS owner. + """ + if interface.online and not interface.detached: + interface.rxb += len(data) + interface.owner.inbound(data, interface) + + +# --------------------------------------------------------------------------- +# Tunnel-frame codec +# --------------------------------------------------------------------------- + +def decode_tunnel_frame(text, msg_prefix="RNS:", interface_name=""): + """ + Locate *msg_prefix* in *text*, strip any sender name prepended by + MeshCore firmware, and base64url-decode the binary frame. + + Returns ``(raw_bytes, sender_prefix)`` on success, or + ``(None, None)`` when the prefix is absent, or + ``(None, sender_prefix)`` when decoding fails. + """ + rns_idx = text.find(msg_prefix) + if rns_idx == -1: + return None, None + + sender_prefix = text[:rns_idx].rstrip(": ") if rns_idx > 0 else "" + b64 = text[rns_idx + len(msg_prefix):].strip() + b64 += "=" * (-len(b64) % 4) + + try: + raw = base64.urlsafe_b64decode(b64) + return raw, sender_prefix + except Exception as exc: + if interface_name: + RNS.log( + f"{interface_name}: " + f"base64 decode error: {exc} " + f"len={len(b64)} text={text[:80]}", + RNS.LOG_WARNING, + ) + return None, sender_prefix + + +# --------------------------------------------------------------------------- +# Fragment reassembly +# --------------------------------------------------------------------------- + +def reassemble_fragment(assembly, assembly_meta, asm_lock, key, + frag_idx, frag_total, payload, interface_name=""): + """ + Insert one fragment into the reassembly buffer and return the + complete packet when all fragments have arrived. + + Returns the reassembled ``bytes`` on completion, or ``None`` if + the packet is still incomplete (or the fragment is a duplicate). + + Thread-safe: acquires *asm_lock* internally. + """ + now = time.monotonic() + with asm_lock: + if key not in assembly: + assembly[key] = {} + assembly_meta[key] = (frag_total, now) + + if frag_idx in assembly[key]: + return None + + assembly[key][frag_idx] = payload + expected = assembly_meta[key][0] + + if len(assembly[key]) < expected: + return None + + try: + full_packet = b"".join( + assembly[key][i] for i in range(expected) + ) + del assembly[key] + del assembly_meta[key] + return full_packet + except KeyError as exc: + RNS.log( + f"{interface_name}: " + f"Reassembly gap — missing fragment {exc}", + RNS.LOG_WARNING, + ) + assembly.pop(key, None) + assembly_meta.pop(key, None) + return None + except Exception as exc: + RNS.log( + f"{interface_name}: Reassembly failed: {exc}", + RNS.LOG_ERROR, + ) + assembly.pop(key, None) + assembly_meta.pop(key, None) + return None + + +def cleanup_stale_fragments(assembly, assembly_meta, asm_lock, timeout_s, + interface_name=""): + """ + Evict fragment assemblies that have exceeded *timeout_s* seconds. + + Returns the number of evicted entries. + """ + deadline = time.monotonic() - timeout_s + with asm_lock: + stale = [k for k, (_, ts) in assembly_meta.items() if ts < deadline] + for k in stale: + del assembly[k] + del assembly_meta[k] + RNS.log( + f"{interface_name}: Evicted stale assembly {k}", + RNS.LOG_WARNING, + ) + return len(stale) + + +# --------------------------------------------------------------------------- +# Packet-ID counter +# --------------------------------------------------------------------------- + +class PacketIdCounter: + """Thread-safe rolling packet ID counter with configurable bit width.""" + + def __init__(self, bits=8): + self._value = 0 + self._mask = (1 << bits) - 1 + self._lock = threading.Lock() + + def next_id(self): + """Return the current value and advance the counter.""" + with self._lock: + val = self._value + self._value = (self._value + 1) & self._mask + return val