From cdf2a21897e07b93e0c48ec5e9ba769f216eec06 Mon Sep 17 00:00:00 2001 From: Ivan Date: Sat, 9 May 2026 19:01:03 -0500 Subject: [PATCH] feat(backend): implement flood protection and improve propagation sync reliability --- meshchatx/meshchat.py | 256 ++++++++++++++++-- .../src/backend/auto_propagation_manager.py | 55 +++- meshchatx/src/backend/config_manager.py | 25 ++ 3 files changed, 303 insertions(+), 33 deletions(-) diff --git a/meshchatx/meshchat.py b/meshchatx/meshchat.py index da4c4b2..920db7b 100644 --- a/meshchatx/meshchat.py +++ b/meshchatx/meshchat.py @@ -336,6 +336,11 @@ class ReticulumMeshChat: # track announce timestamps for rate calculation self.announce_timestamps = [] + # track incoming lxmf message timestamps for flood protection + self._lxmf_incoming_timestamps = [] + self._flood_protection_current_cost = None + self._flood_protection_last_bump_time = 0 + # track download speeds for nomadnetwork files self.download_speeds = [] @@ -2429,6 +2434,10 @@ class ReticulumMeshChat: if not ctx or not ctx.message_router: return + # Always cancel an in-flight sync before switching nodes so we don't + # orphan a transfer or leave the router in a stuck state. + self.stop_propagation_node_sync(context=ctx) + # set outbound propagation node if destination_hash is not None and destination_hash != "": try: @@ -2553,12 +2562,15 @@ class ReticulumMeshChat: ctx = context or self.current_context if not ctx: return - # fixme: it's possible for internal transfer state to get stuck if we change propagation node during a sync - # this still happens even if we cancel the propagation node requests - # for now, the user can just manually cancel syncing in the ui if they think it's stuck... self.stop_propagation_node_sync(context=ctx) if ctx.message_router: ctx.message_router.outbound_propagation_node = None + # Force the transfer state back to idle so nothing remains stuck + # after the outbound node is removed. + with contextlib.suppress(Exception): + ctx.message_router.propagation_transfer_state = ( + ctx.message_router.PR_IDLE + ) # enables or disables the local lxmf propagation node def enable_local_propagation_node(self, enabled: bool = True, context=None): @@ -8664,7 +8676,7 @@ class ReticulumMeshChat: RNS.Transport.request_path(outbound_node) # request messages from propagation node - await self.sync_propagation_nodes() + await self.sync_propagation_nodes(force=True) return web.json_response( { @@ -12951,19 +12963,40 @@ class ReticulumMeshChat: await self.send_announced_to_websocket_clients(context=ctx) # handle syncing propagation nodes - async def sync_propagation_nodes(self, context=None): + async def sync_propagation_nodes(self, context=None, force=False): ctx = context or self.current_context if not ctx: return + router = ctx.message_router + if not router: + return + + # Prevent overlapping auto-syncs from piling up requests. + # A manual/API call can force a restart by cancelling the old sync first. + if router.propagation_transfer_state != router.PR_IDLE: + if not force: + return + self.stop_propagation_node_sync(context=ctx) + # Give the router a moment to settle back to idle + settle_deadline = time.monotonic() + 5.0 + while time.monotonic() < settle_deadline: + if router.propagation_transfer_state == router.PR_IDLE: + break + await asyncio.sleep(0.2) + else: + # Force reset if it didn't settle + with contextlib.suppress(Exception): + router.propagation_transfer_state = router.PR_IDLE + self._begin_propagation_sync_metrics(context=ctx) # update last synced at timestamp ctx.config.lxmf_preferred_propagation_node_last_synced_at.set(int(time.time())) - outbound_node = ctx.message_router.get_outbound_propagation_node() + outbound_node = router.get_outbound_propagation_node() local_propagation_destination = getattr( - ctx.message_router, + router, "propagation_destination", None, ) @@ -12976,16 +13009,14 @@ class ReticulumMeshChat: # Local node selected as preferred: no transport path lookup is needed. # Mark sync as complete immediately to avoid getting stuck in PR_PATH_REQUESTED. with contextlib.suppress(Exception): - ctx.message_router.propagation_transfer_state = ( - ctx.message_router.PR_COMPLETE - ) - ctx.message_router.propagation_transfer_progress = 1.0 - ctx.message_router.propagation_transfer_last_result = 0 + router.propagation_transfer_state = router.PR_COMPLETE + router.propagation_transfer_progress = 1.0 + router.propagation_transfer_last_result = 0 await self.send_config_to_websocket_clients(context=ctx) return # request messages from propagation node - ctx.message_router.request_messages_from_propagation_node(ctx.identity) + router.request_messages_from_propagation_node(ctx.identity) # send config to websocket clients (used to tell ui last synced at) await self.send_config_to_websocket_clients(context=ctx) @@ -13116,17 +13147,24 @@ class ReticulumMeshChat: value = 0 elif value >= 255: value = 254 - self.config.lxmf_inbound_stamp_cost.set(value) - # update the inbound stamp cost on the delivery destination - self.message_router.set_inbound_stamp_cost( - self.local_lxmf_destination.hash, - value, - ) - # re-announce to update the stamp cost in announces - self.local_lxmf_destination.display_name = self.config.display_name.get() - self.message_router.announce( - destination_hash=self.local_lxmf_destination.hash, - ) + # If block strangers is active, store the desired value for later restore + # but keep the enforced max cost active + if self.config.block_all_from_strangers.get(): + self.config.lxmf_inbound_stamp_cost_before_block.set(value) + else: + self.config.lxmf_inbound_stamp_cost.set(value) + # update the inbound stamp cost on the delivery destination + self.message_router.set_inbound_stamp_cost( + self.local_lxmf_destination.hash, + value, + ) + # re-announce to update the stamp cost in announces + self.local_lxmf_destination.display_name = ( + self.config.display_name.get() + ) + self.message_router.announce( + destination_hash=self.local_lxmf_destination.hash, + ) # update propagation node stamp cost (for messages propagated through your node) if "lxmf_propagation_node_stamp_cost" in data: @@ -13378,9 +13416,73 @@ class ReticulumMeshChat: ) if "block_all_from_strangers" in data: - self.config.block_all_from_strangers.set( - self._parse_bool(data["block_all_from_strangers"]), + new_value = self._parse_bool(data["block_all_from_strangers"]) + old_value = self.config.block_all_from_strangers.get() + self.config.block_all_from_strangers.set(new_value) + if new_value and not old_value: + # Enabling block strangers: save current stamp cost and set to max + current_cost = self.config.lxmf_inbound_stamp_cost.get() + if current_cost < 254: + self.config.lxmf_inbound_stamp_cost_before_block.set(current_cost) + self.config.lxmf_inbound_stamp_cost.set(254) + if self.message_router and self.local_lxmf_destination: + self.message_router.set_inbound_stamp_cost( + self.local_lxmf_destination.hash, + 254, + ) + self.local_lxmf_destination.display_name = ( + self.config.display_name.get() + ) + self.message_router.announce( + destination_hash=self.local_lxmf_destination.hash, + ) + elif not new_value and old_value: + # Disabling block strangers: restore previous stamp cost + saved = self.config.lxmf_inbound_stamp_cost_before_block.get() + if saved > 0 and saved < 255: + restore_cost = saved + else: + restore_cost = 8 + self.config.lxmf_inbound_stamp_cost.set(restore_cost) + self.config.lxmf_inbound_stamp_cost_before_block.set(0) + if self.message_router and self.local_lxmf_destination: + self.message_router.set_inbound_stamp_cost( + self.local_lxmf_destination.hash, + restore_cost, + ) + self.local_lxmf_destination.display_name = ( + self.config.display_name.get() + ) + self.message_router.announce( + destination_hash=self.local_lxmf_destination.hash, + ) + + # update flood protection settings + if "lxmf_flood_protection_enabled" in data: + self.config.lxmf_flood_protection_enabled.set( + self._parse_bool(data["lxmf_flood_protection_enabled"]), ) + if "lxmf_flood_threshold_per_minute" in data: + try: + value = int(data["lxmf_flood_threshold_per_minute"]) + value = max(1, min(value, 1000)) + self.config.lxmf_flood_threshold_per_minute.set(value) + except (TypeError, ValueError): + pass + if "lxmf_flood_max_stamp_cost" in data: + try: + value = int(data["lxmf_flood_max_stamp_cost"]) + value = max(1, min(value, 254)) + self.config.lxmf_flood_max_stamp_cost.set(value) + except (TypeError, ValueError): + pass + if "lxmf_flood_cooldown_seconds" in data: + try: + value = int(data["lxmf_flood_cooldown_seconds"]) + value = max(30, min(value, 3600)) + self.config.lxmf_flood_cooldown_seconds.set(value) + except (TypeError, ValueError): + pass if "show_unknown_contact_banner" in data: self.config.show_unknown_contact_banner.set( @@ -14723,6 +14825,10 @@ class ReticulumMeshChat: "lxmf_user_icon_background_colour": ctx.config.lxmf_user_icon_background_colour.get(), "lxmf_inbound_stamp_cost": ctx.config.lxmf_inbound_stamp_cost.get(), "lxmf_propagation_node_stamp_cost": ctx.config.lxmf_propagation_node_stamp_cost.get(), + "lxmf_flood_protection_enabled": ctx.config.lxmf_flood_protection_enabled.get(), + "lxmf_flood_threshold_per_minute": ctx.config.lxmf_flood_threshold_per_minute.get(), + "lxmf_flood_max_stamp_cost": ctx.config.lxmf_flood_max_stamp_cost.get(), + "lxmf_flood_cooldown_seconds": ctx.config.lxmf_flood_cooldown_seconds.get(), "page_archiver_enabled": ctx.config.page_archiver_enabled.get(), "page_archiver_max_versions": ctx.config.page_archiver_max_versions.get(), "archives_max_storage_gb": ctx.config.archives_max_storage_gb.get(), @@ -15282,6 +15388,100 @@ class ReticulumMeshChat: except Exception: return False + def _apply_lxmf_flood_stamp_cost(self, cost: int, context=None) -> None: + """Apply the given inbound stamp cost for flood protection and re-announce.""" + ctx = context or self.current_context + if not ctx or not ctx.message_router or not ctx.local_lxmf_destination: + return + cost = max(0, min(254, cost)) + if cost < 1: + cost = 0 + ctx.config.lxmf_inbound_stamp_cost.set(cost) + ctx.message_router.set_inbound_stamp_cost( + ctx.local_lxmf_destination.hash, + cost, + ) + try: + ctx.local_lxmf_destination.display_name = ctx.config.display_name.get() + ctx.message_router.announce( + destination_hash=ctx.local_lxmf_destination.hash, + ) + except Exception as e: + print(f"_apply_lxmf_flood_stamp_cost: re-announce failed: {e}") + + def _check_lxmf_flood_protection(self, context=None) -> None: + """Check incoming LXMF message rate and auto-adjust stamp cost if flooding.""" + ctx = context or self.current_context + if not ctx or not ctx.config: + return + if not ctx.config.lxmf_flood_protection_enabled.get(): + return + # Do not interfere when block strangers is active (it uses max stamp) + if ctx.config.block_all_from_strangers.get(): + return + + now = time.time() + # Clean old timestamps (> 1 hour) + self._lxmf_incoming_timestamps = [ + t for t in self._lxmf_incoming_timestamps if now - t <= 3600.0 + ] + msgs_per_minute = len( + [t for t in self._lxmf_incoming_timestamps if now - t <= 60.0], + ) + + threshold = ctx.config.lxmf_flood_threshold_per_minute.get() + max_cost = ctx.config.lxmf_flood_max_stamp_cost.get() + current_cost = ctx.config.lxmf_inbound_stamp_cost.get() + if current_cost < 0: + current_cost = 0 + + # Determine base cost (the normal non-flood cost) + if self._flood_protection_current_cost is not None: + base_cost = self._flood_protection_current_cost + else: + base_cost = current_cost + + if msgs_per_minute > threshold: + # Flood detected: bump stamp cost + new_cost = min(current_cost + 2, max_cost) + if new_cost != current_cost: + print( + f"LXMF flood detected: {msgs_per_minute} msg/min " + f"(threshold {threshold}). Raising stamp cost from " + f"{current_cost} to {new_cost}.", + ) + if self._flood_protection_current_cost is None: + self._flood_protection_current_cost = base_cost + self._flood_protection_last_bump_time = now + self._apply_lxmf_flood_stamp_cost(new_cost, context=ctx) + elif current_cost > base_cost: + cooldown = ctx.config.lxmf_flood_cooldown_seconds.get() + if now - self._flood_protection_last_bump_time > cooldown: + # Step down by 1 toward base cost + new_cost = max(current_cost - 1, base_cost) + if new_cost != current_cost: + print( + f"LXMF flood subsided: {msgs_per_minute} msg/min. " + f"Lowering stamp cost from {current_cost} to {new_cost}.", + ) + self._apply_lxmf_flood_stamp_cost(new_cost, context=ctx) + if new_cost == base_cost: + self._flood_protection_current_cost = None + self._flood_protection_last_bump_time = 0 + + async def lxmf_flood_protection_cooldown_loop(self, session_id, context=None): + """Background loop to step down flood protection stamp cost during quiet periods.""" + ctx = context or self.current_context + if not ctx: + return + await asyncio.sleep(60) + while self.running and ctx.running and ctx.session_id == session_id: + try: + self._check_lxmf_flood_protection(context=ctx) + except Exception as e: + print(f"lxmf_flood_protection_cooldown_loop error: {e}") + await asyncio.sleep(30) + def _collect_lxmf_sieve_peer_haystack( self, peer_hash: str, @@ -15453,6 +15653,10 @@ class ReticulumMeshChat: print(f"Rejecting LXMF message from blocked source: {source_hash}") return + # track incoming message timestamps for flood protection + self._lxmf_incoming_timestamps.append(time.time()) + self._check_lxmf_flood_protection(context=ctx) + is_sideband_telemetry_request = False lxmf_fields = lxmf_message.get_fields() diff --git a/meshchatx/src/backend/auto_propagation_manager.py b/meshchatx/src/backend/auto_propagation_manager.py index 8d052ba..5745d4d 100644 --- a/meshchatx/src/backend/auto_propagation_manager.py +++ b/meshchatx/src/backend/auto_propagation_manager.py @@ -1,6 +1,7 @@ # SPDX-License-Identifier: 0BSD import asyncio +import contextlib import time import RNS @@ -78,7 +79,17 @@ class AutoPropagationManager: except Exception: return False + # Ensure any previous sync is fully cancelled and state is idle self.app.stop_propagation_node_sync(context=ctx) + settle_deadline = time.monotonic() + 5.0 + while time.monotonic() < settle_deadline: + if router.propagation_transfer_state == LXMRouter.PR_IDLE: + break + await asyncio.sleep(POLL_INTERVAL_SECONDS) + else: + with contextlib.suppress(Exception): + router.propagation_transfer_state = LXMRouter.PR_IDLE + try: router.set_outbound_propagation_node(dest) except Exception: @@ -87,16 +98,29 @@ class AutoPropagationManager: router.request_messages_from_propagation_node(ctx.identity) deadline = time.monotonic() + SYNC_PROBE_TIMEOUT_SECONDS - seen_progress = False + # Wait for the sync to actually start (leave idle) while time.monotonic() < deadline: state = router.propagation_transfer_state if state in _PROP_FAILURE_STATES: self.app.stop_propagation_node_sync(context=ctx) return False if state != LXMRouter.PR_IDLE: - seen_progress = True - elif seen_progress: + break + await asyncio.sleep(POLL_INTERVAL_SECONDS) + else: + # Never left idle -> failed to start + self.app.stop_propagation_node_sync(context=ctx) + return False + + # Wait for the sync to finish (return to idle) + while time.monotonic() < deadline: + state = router.propagation_transfer_state + if state in _PROP_FAILURE_STATES: + self.app.stop_propagation_node_sync(context=ctx) + return False + if state == LXMRouter.PR_IDLE: + # Success: we left idle and now we're back self.app.stop_propagation_node_sync(context=ctx) return True await asyncio.sleep(POLL_INTERVAL_SECONDS) @@ -115,22 +139,39 @@ class AutoPropagationManager: ) # If a sync is in progress, only interrupt it when the current node - # appears unreachable. This prevents getting stuck on a node we - # cannot get a path to. + # appears unreachable or the path is stale/unresponsive. This prevents + # getting stuck on a node we cannot actually reach. if router.propagation_transfer_state != LXMRouter.PR_IDLE: - current_has_path = False + current_path_ok = False if previous_hex: try: current_dest = bytes.fromhex(previous_hex) current_has_path = RNS.Transport.has_path(current_dest) + current_path_ok = ( + current_has_path + and not RNS.Transport.path_is_unresponsive(current_dest) + and not reticulum_pathfinding.transport_path_table_entry_is_expired( + current_dest, + ) + ) except Exception: pass - if current_has_path: + if current_path_ok: # Sync is likely making progress – let it finish. return # Current node is unreachable – stop the stuck sync so we can # look for a working alternative. self.app.stop_propagation_node_sync(context=ctx) + # Wait briefly for the router to settle back to idle before we + # start probing other nodes. + settle_deadline = time.monotonic() + 5.0 + while time.monotonic() < settle_deadline: + if router.propagation_transfer_state == LXMRouter.PR_IDLE: + break + await asyncio.sleep(POLL_INTERVAL_SECONDS) + else: + with contextlib.suppress(Exception): + router.propagation_transfer_state = LXMRouter.PR_IDLE announces = self.database.announces.get_announces(aspect="lxmf.propagation") diff --git a/meshchatx/src/backend/config_manager.py b/meshchatx/src/backend/config_manager.py index e358a09..3da2ec6 100644 --- a/meshchatx/src/backend/config_manager.py +++ b/meshchatx/src/backend/config_manager.py @@ -106,6 +106,31 @@ class ConfigManager: "lxmf_propagation_node_stamp_cost", 16, ) # for propagation node messages + self.lxmf_inbound_stamp_cost_before_block = self.IntConfig( + self, + "lxmf_inbound_stamp_cost_before_block", + 0, + ) # saved stamp cost before block strangers was enabled + self.lxmf_flood_protection_enabled = self.BoolConfig( + self, + "lxmf_flood_protection_enabled", + False, + ) + self.lxmf_flood_threshold_per_minute = self.IntConfig( + self, + "lxmf_flood_threshold_per_minute", + 30, + ) + self.lxmf_flood_max_stamp_cost = self.IntConfig( + self, + "lxmf_flood_max_stamp_cost", + 24, + ) + self.lxmf_flood_cooldown_seconds = self.IntConfig( + self, + "lxmf_flood_cooldown_seconds", + 300, + ) self.page_archiver_enabled = self.BoolConfig( self, "page_archiver_enabled",