From 2af04378a1d0b04b994d98c9172098c01670ffa3 Mon Sep 17 00:00:00 2001 From: Ivan Date: Thu, 7 May 2026 17:11:47 -0500 Subject: [PATCH 01/63] chore(dependencies): update rns version to 1.2.4 in poetry.lock, pyproject.toml, requirements.txt, and build.gradle --- android/app/build.gradle | 2 +- poetry.lock | 9 +++++---- pyproject.toml | 2 +- requirements.txt | 2 +- 4 files changed, 8 insertions(+), 7 deletions(-) diff --git a/android/app/build.gradle b/android/app/build.gradle index 71c75f9..82006e3 100644 --- a/android/app/build.gradle +++ b/android/app/build.gradle @@ -163,7 +163,7 @@ chaquopy { options "--find-links", vendorWheelDir.absolutePath install "packaging>=23" install "aiohttp==3.13.5" - install "rns>=1.2.1" + install "rns>=1.2.4" install "bleak==3.0.1" install "lxmf>=0.9.4" install "numpy==1.26.2" diff --git a/poetry.lock b/poetry.lock index 6719f0f..2d2b56e 100644 --- a/poetry.lock +++ b/poetry.lock @@ -2458,14 +2458,15 @@ jupyter = ["ipywidgets (>=7.5.1,<9)"] [[package]] name = "rns" -version = "1.2.3" +version = "1.2.4" description = "Self-configuring, encrypted and resilient mesh networking stack for LoRa, packet radio, WiFi and everything in between" optional = false python-versions = ">=3.7" groups = ["main"] files = [ - {file = "rns-1.2.3-py3-none-any.whl", hash = "sha256:8562130f297a6b33be9d72c449bbe6ae83cad41e1530e0fa112f5fa545a3f364"}, - {file = "rns-1.2.3.tar.gz", hash = "sha256:5388b82cf60ffa84ed18dc56e1f5e3d500b5864d08b34d4317bba766088c0a03"}, + {file = "rns-1.2.4-1-py3-none-any.whl", hash = "sha256:cbf1d2ff517c91288b2b924e62e3e5b3f2732d35331c4bbd07329bf21a7c14bf"}, + {file = "rns-1.2.4-py3-none-any.whl", hash = "sha256:e821a0b6a18d6b3263bbcdde880d0388fb4dd0c07c7eb2f83cb0dbc30eda5965"}, + {file = "rns-1.2.4.tar.gz", hash = "sha256:49798a8cc23632ec0e4d4cd88c1408ef02bee613995a6c57d52a37fb79c329ca"}, ] [package.dependencies] @@ -3376,4 +3377,4 @@ propcache = ">=0.2.1" [metadata] lock-version = "2.1" python-versions = ">=3.11" -content-hash = "11fe20cf20f8ca70a22450f3303fe85cd4b121e714b37806b3100134b83bd2f0" +content-hash = "9f9acf8091840769458b698894d3cc4997e4cea68d7e049562c7f4520746c4ee" diff --git a/pyproject.toml b/pyproject.toml index 273e54a..f429a4e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -24,7 +24,7 @@ dependencies = [ "lxmf>=0.9.6", "psutil>=7.2.2", "bleak==3.0.1", - "rns>=1.2.3", + "rns>=1.2.4", "websockets>=16.0", "bcrypt>=5.0.0,<6.0.0", "aiohttp-session>=2.12.1,<3.0.0", diff --git a/requirements.txt b/requirements.txt index 3a4f015..f17db53 100644 --- a/requirements.txt +++ b/requirements.txt @@ -23,7 +23,7 @@ pycodec2==4.1.1 ; python_version >= "3.11" pycparser==3.0 ; python_version >= "3.11" pyserial==3.5 ; python_version >= "3.11" bleak==3.0.1 ; python_version >= "3.11" -rns==1.2.3 ; python_version >= "3.11" +rns==1.2.4 ; python_version >= "3.11" typing-extensions==4.15.0 ; python_version >= "3.11" and python_version < "3.13" websockets==16.0 ; python_version >= "3.11" yarl==1.23.0 ; python_version >= "3.11" From e8a6e7dc5a3523b850aa66789221112ec7635e93 Mon Sep 17 00:00:00 2001 From: Ivan Date: Thu, 7 May 2026 20:03:22 -0500 Subject: [PATCH 02/63] feat(logging): implement console logging restoration after Reticulum initialization and add instance name corruption repair for hot reloads --- meshchatx/meshchat.py | 140 +++++++++++++++++++++++++++++++++++++----- 1 file changed, 123 insertions(+), 17 deletions(-) diff --git a/meshchatx/meshchat.py b/meshchatx/meshchat.py index a3abc92..ccdce5e 100644 --- a/meshchatx/meshchat.py +++ b/meshchatx/meshchat.py @@ -214,6 +214,29 @@ def _resolve_rns_loglevel(cli_override: str | None) -> int | None: return _parse_rns_loglevel_value(os.environ.get("MESHCHAT_RNS_LOG_LEVEL")) +def _restore_rns_console_logging_after_reticulum_init(app) -> None: + """Undo shutdown side effects from ``RNS.Reticulum.exit_handler``. + + That handler sets ``RNS.loglevel`` to ``LOG_NONE`` and points ``sys.stdout`` / + ``sys.stderr`` at ``os.devnull``. Without this, hot reload appears to stop all + announce traffic logging even though interfaces are up. + + When no CLI or ``MESHCHAT_RNS_LOG_LEVEL`` value applies and the level is still + ``LOG_NONE`` after reading config, fall back to ``LOG_WARNING`` so notices are + visible. Explicit ``none`` in the environment remains respected. + """ + try: + if hasattr(sys, "__stdout__"): + sys.stdout = sys.__stdout__ + if hasattr(sys, "__stderr__"): + sys.stderr = sys.__stderr__ + except Exception: + pass + resolved = _resolve_rns_loglevel(getattr(app, "_rns_loglevel_cli", None)) + if resolved is None and RNS.loglevel == RNS.LOG_NONE: + RNS.loglevel = RNS.LOG_WARNING + + def _python_jit_status_line() -> str: jit_runtime = getattr(sys, "_jit", None) if jit_runtime is None: @@ -673,6 +696,9 @@ class ReticulumMeshChat: self.reticulum_config_dir = config_dir if not materialize: return + if not getattr(self, "_reticulum_instance_name_startup_repair_done", False): + self._repair_reticulum_instance_name_corruption() + self._reticulum_instance_name_startup_repair_done = True config_path = os.path.join(config_dir, "config") needs_default = True if os.path.isfile(config_path): @@ -715,6 +741,7 @@ class ReticulumMeshChat: ) else: self.reticulum = RNS.Reticulum(self.reticulum_config_dir) + _restore_rns_console_logging_after_reticulum_init(self) self.page_node_manager.load_nodes() self.page_node_manager.start_all() @@ -983,6 +1010,49 @@ class ReticulumMeshChat: return closed_any + _reload_instance_suffix_re = re.compile(r"-reload-(\d+)-(\d+)$") + _meshchat_reload_pid_max = 4_194_304 + _meshchat_reload_epoch_min = 1_577_836_800 + _meshchat_reload_epoch_max = 4_102_444_800 + + @staticmethod + def _looks_like_meshchat_hot_reload_tail(pid: int, epoch: int) -> bool: + """Limit repairs to suffixes :meth:`reload_reticulum` actually writes. + + Hot reload uses ``-reload-{os.getpid()}-{int(time.time())}``. Names like + ``my-net-reload-peer`` must not be truncated. + """ + if pid < 1 or pid > ReticulumMeshChat._meshchat_reload_pid_max: + return False + if ( + epoch < ReticulumMeshChat._meshchat_reload_epoch_min + or epoch > ReticulumMeshChat._meshchat_reload_epoch_max + ): + return False + return True + + @staticmethod + def _strip_reload_instance_suffix(name): + """Remove stacked MeshChat hot-reload tails only (validated pid + unix time).""" + if not isinstance(name, str): + return None + out = name.strip() + if not out: + return None + while True: + m = ReticulumMeshChat._reload_instance_suffix_re.search(out) + if not m or m.end() != len(out): + break + try: + pid = int(m.group(1)) + epoch = int(m.group(2)) + except ValueError: + break + if not ReticulumMeshChat._looks_like_meshchat_hot_reload_tail(pid, epoch): + break + out = out[: m.start()].strip() + return out if out else None + def _read_reticulum_instance_name(self): """Return current Reticulum instance_name from config or None.""" config_dir = self._normalize_reticulum_config_dir( @@ -998,6 +1068,16 @@ class ReticulumMeshChat: return None return cp.get("reticulum", "instance_name", fallback=None) + def _repair_reticulum_instance_name_corruption(self): + """Rewrite persisted ``instance_name`` if hot-reload suffixes were left on disk.""" + raw = self._read_reticulum_instance_name() + if not raw: + return + cleaned = ReticulumMeshChat._strip_reload_instance_suffix(raw) + if cleaned == raw or cleaned is None: + return + self._write_reticulum_instance_name(cleaned) + def _write_reticulum_instance_name(self, instance_name): """Persist a Reticulum instance_name value into the config.""" config_dir = self._normalize_reticulum_config_dir( @@ -1469,13 +1549,18 @@ class ReticulumMeshChat: if hasattr(RNS.Reticulum, "_Reticulum__instance"): RNS.Reticulum._Reticulum__instance = None - original_instance_name = None switched_instance_name = None + instance_restore_name = None if abstract_unix_addr_in_use_after_wait: - original_instance_name = self._read_reticulum_instance_name() - base_name = original_instance_name or "default" + stored_instance_name = self._read_reticulum_instance_name() + stable_base = ReticulumMeshChat._strip_reload_instance_suffix( + stored_instance_name, + ) + instance_restore_name = ( + stable_base if stable_base is not None else "default" + ) switched_instance_name = ( - f"{base_name}-reload-{os.getpid()}-{int(time.time())}" + f"{instance_restore_name}-reload-{os.getpid()}-{int(time.time())}" ) self._write_reticulum_instance_name(switched_instance_name) print( @@ -1492,9 +1577,7 @@ class ReticulumMeshChat: self.setup_identity(identity_to_restore) finally: if switched_instance_name: - self._write_reticulum_instance_name( - original_instance_name or "default", - ) + self._write_reticulum_instance_name(instance_restore_name) await self._send_rns_reload_status( "done", "RNS reload complete.", @@ -1920,7 +2003,13 @@ class ReticulumMeshChat: return None @staticmethod - def apply_bootstrap_only_to_interface(interface_details, data, default_enabled): + def apply_bootstrap_only_to_interface( + interface_details, + data, + default_enabled, + *, + updating_existing=False, + ): if "bootstrap_only" in data: yn = ReticulumMeshChat._bootstrap_only_request_yes_no( data.get("bootstrap_only") @@ -1932,6 +2021,8 @@ class ReticulumMeshChat: else: interface_details.pop("bootstrap_only", None) return + if updating_existing: + return if default_enabled: interface_details["bootstrap_only"] = "yes" else: @@ -2309,7 +2400,7 @@ class ReticulumMeshChat: # uses the provided destination hash as the active propagation node def set_active_propagation_node(self, destination_hash: str | None, context=None): ctx = context or self.current_context - if not ctx: + if not ctx or not ctx.message_router: return # set outbound propagation node @@ -2332,7 +2423,7 @@ class ReticulumMeshChat: # stops the in progress propagation node sync def stop_propagation_node_sync(self, context=None): ctx = context or self.current_context - if not ctx: + if not ctx or not ctx.message_router: return ctx.message_router.cancel_propagation_node_requests() @@ -2395,6 +2486,13 @@ class ReticulumMeshChat: "messages_hidden": 0, } + if not ctx.message_router: + return { + "messages_stored": 0, + "delivery_confirmations": 0, + "messages_hidden": 0, + } + messages_received = ctx.message_router.propagation_transfer_last_result or 0 current_total_messages = ctx.database.messages.count_lxmf_messages() current_delivered_messages = ctx.database.messages.count_lxmf_messages_by_state( @@ -2433,12 +2531,13 @@ class ReticulumMeshChat: # 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) - ctx.message_router.outbound_propagation_node = None + if ctx.message_router: + ctx.message_router.outbound_propagation_node = None # enables or disables the local lxmf propagation node def enable_local_propagation_node(self, enabled: bool = True, context=None): ctx = context or self.current_context - if not ctx: + if not ctx or not ctx.message_router: return try: if enabled: @@ -2469,6 +2568,9 @@ class ReticulumMeshChat: return None router = ctx.message_router + if not router: + return None + is_running = bool(getattr(router, "propagation_node", False)) stats = None if is_running: @@ -2479,13 +2581,13 @@ class ReticulumMeshChat: return value if isinstance(value, (int, float)) else default destination_hash_raw = getattr( - ctx.message_router.propagation_destination, + router.propagation_destination, "hexhash", None, ) if destination_hash_raw is None: destination_hash_raw = getattr( - ctx.message_router.propagation_destination, + router.propagation_destination, "hash", None, ) @@ -4912,12 +5014,13 @@ class ReticulumMeshChat: ): default_boot = ReticulumMeshChat._reticulum_yes_no_preference( self._get_reticulum_section().get("default_bootstrap_only"), - default=True, + default=False, ) ReticulumMeshChat.apply_bootstrap_only_to_interface( interface_details, data, default_boot, + updating_existing=allow_overwriting_interface, ) # set common interface options @@ -6412,7 +6515,7 @@ class ReticulumMeshChat: ), "default_bootstrap_only": ReticulumMeshChat._reticulum_yes_no_preference( reticulum_config.get("default_bootstrap_only"), - default=True, + default=False, ), "network_identity": reticulum_config.get("network_identity"), } @@ -6495,7 +6598,7 @@ class ReticulumMeshChat: ), "default_bootstrap_only": ReticulumMeshChat._reticulum_yes_no_preference( reticulum_config.get("default_bootstrap_only"), - default=True, + default=False, ), "network_identity": reticulum_config.get("network_identity"), } @@ -6578,9 +6681,12 @@ class ReticulumMeshChat: "listen_ip": s.get("listen_ip"), "connected": s.get("connected"), "online": s.get("online"), + "status": s.get("status"), "transport_id": transport_id, "network_id": s.get("network_id"), "autoconnect_source": s.get("autoconnect_source"), + "txb": s.get("txb"), + "rxb": s.get("rxb"), }, ) except Exception as e: From 7867034f5d7e21595f7481f75710a1246946b49f Mon Sep 17 00:00:00 2001 From: Ivan Date: Thu, 7 May 2026 20:03:42 -0500 Subject: [PATCH 03/63] fix(auto_propagation_manager): add router existence checks in propagation methods to prevent errors --- meshchatx/src/backend/auto_propagation_manager.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/meshchatx/src/backend/auto_propagation_manager.py b/meshchatx/src/backend/auto_propagation_manager.py index 58da785..8d052ba 100644 --- a/meshchatx/src/backend/auto_propagation_manager.py +++ b/meshchatx/src/backend/auto_propagation_manager.py @@ -69,6 +69,8 @@ class AutoPropagationManager: async def _probe_propagation_sync(self, node_hex: str) -> bool: ctx = self.context router = ctx.message_router + if not router: + return False try: dest = bytes.fromhex(node_hex) if len(dest) != RNS.Identity.TRUNCATED_HASHLENGTH // 8: @@ -105,6 +107,8 @@ class AutoPropagationManager: async def check_and_update_propagation_node(self): ctx = self.context router = ctx.message_router + if not router: + return previous_hex = ( self.config.lxmf_preferred_propagation_node_destination_hash.get() From f2235f2e0bf3aed6fadd7eed9eaf67bb8354f262 Mon Sep 17 00:00:00 2001 From: Ivan Date: Thu, 7 May 2026 20:03:52 -0500 Subject: [PATCH 04/63] feat(tests): add integration tests for announce limits and spam handling, including max storage constraints and websocket broadcast behavior --- tests/backend/test_announce_dao_trim.py | 13 + tests/backend/test_announce_limits.py | 85 +++++ .../test_announce_limits_integration.py | 338 ++++++++++++++++++ tests/backend/test_announce_spam_sqlite.py | 218 +++++++++++ tests/backend/test_async_utils_critical.py | 108 ++++++ tests/backend/test_auto_propagation.py | 21 ++ tests/backend/test_interface_discovery.py | 125 ++++++- tests/backend/test_long_running_stress.py | 247 +++++++++++++ tests/backend/test_search_integration.py | 21 ++ tests/backend/test_websocket_scale.py | 51 +++ 10 files changed, 1226 insertions(+), 1 deletion(-) create mode 100644 tests/backend/test_announce_limits_integration.py create mode 100644 tests/backend/test_announce_spam_sqlite.py create mode 100644 tests/backend/test_async_utils_critical.py create mode 100644 tests/backend/test_long_running_stress.py diff --git a/tests/backend/test_announce_dao_trim.py b/tests/backend/test_announce_dao_trim.py index d78c6b2..c5ed439 100644 --- a/tests/backend/test_announce_dao_trim.py +++ b/tests/backend/test_announce_dao_trim.py @@ -53,6 +53,19 @@ def _cleanup(db, path): pass +def test_trim_announces_for_aspect_noop_when_max_rows_below_one(): + db = path = None + try: + db, path = _new_db() + aspect = "lxmf.delivery" + _insert(db, "01" * 16, aspect, "2000-01-01T00:00:00Z") + _insert(db, "02" * 16, aspect, "2000-01-02T00:00:00Z") + db.announces.trim_announces_for_aspect(aspect, 0) + assert db.announces.get_announce_count_by_aspect(aspect) == 2 + finally: + _cleanup(db, path) + + def test_trim_announces_for_aspect_drops_oldest(): db = path = None try: diff --git a/tests/backend/test_announce_limits.py b/tests/backend/test_announce_limits.py index 59bb9b4..4ce2c3c 100644 --- a/tests/backend/test_announce_limits.py +++ b/tests/backend/test_announce_limits.py @@ -166,6 +166,91 @@ def test_get_filtered_announces_resolves_default_limit(mock_db, mock_config): assert 33 in params +def test_max_stored_clamps_to_one_million(mock_db, mock_config): + mock_config.announce_max_stored_lxmf_delivery.get.return_value = 9_999_999 + + manager = _make_manager(mock_db, mock_config) + + assert manager._get_max_stored_for_aspect("lxmf.delivery") == 1_000_000 + + +def test_trim_called_with_clamped_max_stored(mock_db, mock_config): + mock_config.announce_max_stored_lxmf_delivery.get.return_value = 5_000_000 + + manager = _make_manager(mock_db, mock_config) + reticulum = MagicMock() + identity = MagicMock() + identity.hash.hex.return_value = "ab" * 16 + identity.get_public_key.return_value = b"pub_key" + + manager.upsert_announce( + reticulum, + identity, + b"\x01" * 16, + "lxmf.delivery", + b"app_data", + b"packet_hash", + ) + + mock_db.announces.trim_announces_for_aspect.assert_called_once_with( + "lxmf.delivery", + 1_000_000, + ) + + +def test_max_stored_zero_skips_cap(mock_db, mock_config): + mock_config.announce_max_stored_lxmf_delivery.get.return_value = 0 + + manager = _make_manager(mock_db, mock_config) + + assert manager._get_max_stored_for_aspect("lxmf.delivery") is None + + +def test_fetch_limit_clamps_to_hundred_thousand(mock_db, mock_config): + mock_config.announce_fetch_limit_lxmf_delivery = MagicMock() + mock_config.announce_fetch_limit_lxmf_delivery.get.return_value = 800_000 + + manager = _make_manager(mock_db, mock_config) + + assert manager._get_fetch_limit_for_aspect("lxmf.delivery") == 100_000 + + +def test_fetch_limit_invalid_falls_back_to_default(mock_db, mock_config): + mock_config.announce_fetch_limit_lxmf_delivery = MagicMock() + mock_config.announce_fetch_limit_lxmf_delivery.get.return_value = 0 + + manager = _make_manager(mock_db, mock_config) + + assert manager._get_fetch_limit_for_aspect("lxmf.delivery") == 500 + + +def test_fetch_limit_unknown_aspect_returns_default(mock_db, mock_config): + manager = _make_manager(mock_db, mock_config) + + assert manager._get_fetch_limit_for_aspect("unknown.aspect") == 500 + + +def test_fetch_limit_none_falls_back_to_default(mock_db, mock_config): + mock_config.announce_fetch_limit_lxmf_delivery = MagicMock() + mock_config.announce_fetch_limit_lxmf_delivery.get.return_value = None + + manager = _make_manager(mock_db, mock_config) + + assert manager._get_fetch_limit_for_aspect("lxmf.delivery") == 500 + + +def test_get_filtered_announces_uses_clamped_fetch_limit(mock_db, mock_config): + mock_config.announce_fetch_limit_lxmf_delivery = MagicMock() + mock_config.announce_fetch_limit_lxmf_delivery.get.return_value = 400_000 + + manager = _make_manager(mock_db, mock_config) + manager.get_filtered_announces(aspect="lxmf.delivery", limit=None) + + args, _ = mock_db.provider.fetchall.call_args + _sql, params = args + assert 100_000 in params + + def test_announce_handles_none_packet_hash(mock_db): manager = AnnounceManager(mock_db) reticulum = MagicMock() diff --git a/tests/backend/test_announce_limits_integration.py b/tests/backend/test_announce_limits_integration.py new file mode 100644 index 0000000..36b501d --- /dev/null +++ b/tests/backend/test_announce_limits_integration.py @@ -0,0 +1,338 @@ +# SPDX-License-Identifier: 0BSD + +"""Integration tests: announce row caps via AnnounceManager + real SQLite.""" + +import os +import tempfile +from unittest.mock import MagicMock + +import pytest + +from meshchatx.src.backend.announce_manager import AnnounceManager +from meshchatx.src.backend.database import Database +from meshchatx.src.backend.database.provider import DatabaseProvider + + +class _FakeIdentity: + __slots__ = ("_h",) + + def __init__(self, identity_hex32: str): + self._h = bytes.fromhex(identity_hex32) + + @property + def hash(self): + return self._h + + def get_public_key(self): + return b"\xaa\xbb" + + +def _cleanup(db, path): + if db is not None: + try: + db.close() + except Exception: + pass + DatabaseProvider._instance = None + if path: + try: + os.unlink(path) + except OSError: + pass + for suffix in ("-wal", "-shm"): + try: + os.unlink(path + suffix) + except OSError: + pass + + +def _new_db(): + with tempfile.NamedTemporaryFile(suffix=".db", delete=False) as f: + path = f.name + db = Database(path) + db.initialize() + return db, path + + +def _store_enabled_config(**max_stored): + """Config mock with storage toggles on and configurable announce_max_stored_* .get() values.""" + config = MagicMock() + for _k in ( + "announce_store_lxmf_delivery", + "announce_store_lxst_telephony", + "announce_store_nomadnetwork_node", + "announce_store_lxmf_propagation", + "announce_store_git_repositories", + ): + m = MagicMock() + m.get.return_value = True + setattr(config, _k, m) + + for key, default in ( + ("announce_max_stored_lxmf_delivery", None), + ("announce_max_stored_nomadnetwork_node", None), + ("announce_max_stored_lxmf_propagation", None), + ): + attr = MagicMock() + attr.get.return_value = max_stored.get(key, default) + setattr(config, key, attr) + + return config + + +@pytest.fixture +def sqlite_db(): + db, path = _new_db() + yield db, path + _cleanup(db, path) + + +def test_many_sequential_upserts_trims_to_max(sqlite_db): + db, _path = sqlite_db + max_keep = 12 + cfg = _store_enabled_config(announce_max_stored_lxmf_delivery=max_keep) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + + n_insert = 55 + for i in range(n_insert): + dh = f"{i:032x}" + ident = _FakeIdentity(f"{i:032x}") + mgr.upsert_announce( + ret, + ident, + bytes.fromhex(dh), + "lxmf.delivery", + b"payload", + None, + ) + + assert db.announces.get_announce_count_by_aspect("lxmf.delivery") == max_keep + rows = db.announces.get_announces(aspect="lxmf.delivery") + kept = {r["destination_hash"] for r in rows} + expect = {f"{i:032x}" for i in range(n_insert - max_keep, n_insert)} + assert kept == expect + + +def test_aspect_max_limits_are_independent(sqlite_db): + db, _path = sqlite_db + cfg = _store_enabled_config( + announce_max_stored_lxmf_delivery=7, + announce_max_stored_nomadnetwork_node=4, + ) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + + for i in range(20): + dh = f"{i:032x}" + ident = _FakeIdentity(f"{i:032x}") + mgr.upsert_announce( + ret, + ident, + bytes.fromhex(dh), + "lxmf.delivery", + b"x", + None, + ) + + for i in range(15): + dh = f"{0x70000000 + i:032x}" + ident = _FakeIdentity(f"{0x71000000 + i:032x}") + mgr.upsert_announce( + ret, + ident, + bytes.fromhex(dh), + "nomadnetwork.node", + b"y", + None, + ) + + assert db.announces.get_announce_count_by_aspect("lxmf.delivery") == 7 + assert db.announces.get_announce_count_by_aspect("nomadnetwork.node") == 4 + + +def test_repeated_upsert_same_destination_does_not_expand_table(sqlite_db): + db, _path = sqlite_db + max_keep = 10 + cfg = _store_enabled_config(announce_max_stored_lxmf_delivery=max_keep) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + + primary_dest = "f" * 32 + ident = _FakeIdentity("e" * 32) + + for _ in range(80): + mgr.upsert_announce( + ret, + ident, + bytes.fromhex(primary_dest), + "lxmf.delivery", + b"v1", + None, + ) + + for i in range(25): + dh = f"{i:032x}" + oid = _FakeIdentity(f"1{i:031x}") + mgr.upsert_announce( + ret, + oid, + bytes.fromhex(dh), + "lxmf.delivery", + b"x", + None, + ) + + assert db.announces.get_announce_count_by_aspect("lxmf.delivery") == max_keep + dup_rows = db.provider.fetchall( + """ + SELECT destination_hash FROM announces WHERE aspect = ? + GROUP BY destination_hash HAVING COUNT(*) > 1 + """, + ("lxmf.delivery",), + ) + assert dup_rows == [] + + +def test_manager_trim_skips_contact_linked_identity(sqlite_db): + db, _path = sqlite_db + cfg = _store_enabled_config(announce_max_stored_lxmf_delivery=500) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + + protected_idx = 3 + contact_ih = f"{protected_idx:032x}" + + for i in range(8): + dh = f"{i:032x}" + ident = _FakeIdentity(f"{i:032x}") + mgr.upsert_announce( + ret, + ident, + bytes.fromhex(dh), + "lxmf.delivery", + b"p", + None, + ) + + db.contacts.add_contact("peer", contact_ih) + + cfg.announce_max_stored_lxmf_delivery.get.return_value = 3 + mgr.upsert_announce( + ret, + _FakeIdentity("ffffffffffffffffffffffffffffffff"), + bytes.fromhex("ffffffffffffffffffffffffffffffff"), + "lxmf.delivery", + b"tick", + None, + ) + + assert db.announces.get_announce_count_by_aspect("lxmf.delivery") == 3 + rows = db.announces.get_announces(aspect="lxmf.delivery") + hashes = {r["destination_hash"] for r in rows} + assert f"{protected_idx:032x}" in hashes + + +def test_trim_after_prefilled_table_overflow(sqlite_db): + """Simulates a large announce backlog (direct DAO inserts), then one managed upsert.""" + db, _path = sqlite_db + aspect = "lxmf.delivery" + for i in range(220): + dh = f"{i:032x}" + db.announces.upsert_announce( + { + "destination_hash": dh, + "aspect": aspect, + "identity_hash": f"{i:032x}", + "identity_public_key": "cHVibmtleQ==", + "app_data": None, + "rssi": None, + "snr": None, + "quality": None, + }, + ) + + assert db.announces.get_announce_count_by_aspect(aspect) == 220 + + max_keep = 15 + cfg = _store_enabled_config(announce_max_stored_lxmf_delivery=max_keep) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + + mgr.upsert_announce( + ret, + _FakeIdentity("aa" * 16), + bytes.fromhex(f"{220:032x}"), + aspect, + b"flush", + None, + ) + + assert db.announces.get_announce_count_by_aspect(aspect) == max_keep + kept = {r["destination_hash"] for r in db.announces.get_announces(aspect=aspect)} + assert kept == {f"{i:032x}" for i in range(206, 221)} + + +def test_integration_respects_favourite_under_tight_cap(sqlite_db): + db, _path = sqlite_db + aspect = "lxmf.delivery" + favourite_dest = f"{5:032x}" + for i in range(24): + dh = f"{i:032x}" + db.announces.upsert_announce( + { + "destination_hash": dh, + "aspect": aspect, + "identity_hash": f"{i:032x}", + "identity_public_key": "cHVibmtleQ==", + "app_data": None, + "rssi": None, + "snr": None, + "quality": None, + }, + ) + + db.announces.upsert_favourite(favourite_dest, "Pinned", aspect) + + cfg = _store_enabled_config(announce_max_stored_lxmf_delivery=4) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + + mgr.upsert_announce( + ret, + _FakeIdentity("bb" * 16), + bytes.fromhex(f"{100:032x}"), + aspect, + b"tight", + None, + ) + + rows = db.announces.get_announces(aspect=aspect) + hashes = {r["destination_hash"] for r in rows} + assert favourite_dest in hashes + assert f"{100:032x}" in hashes + + +def test_lxst_telephony_shares_lxmf_delivery_cap(sqlite_db): + db, _path = sqlite_db + cfg = _store_enabled_config(announce_max_stored_lxmf_delivery=6) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + + for i in range(10): + dh = f"{0x60000000 + i:032x}" + mgr.upsert_announce( + ret, + _FakeIdentity(f"{0x61000000 + i:032x}"), + bytes.fromhex(dh), + "lxst.telephony", + b"t", + None, + ) + + assert db.announces.get_announce_count_by_aspect("lxst.telephony") == 6 + kept = { + r["destination_hash"] + for r in db.announces.get_announces(aspect="lxst.telephony") + } + assert kept == {f"{0x60000000 + i:032x}" for i in range(4, 10)} diff --git a/tests/backend/test_announce_spam_sqlite.py b/tests/backend/test_announce_spam_sqlite.py new file mode 100644 index 0000000..e240b30 --- /dev/null +++ b/tests/backend/test_announce_spam_sqlite.py @@ -0,0 +1,218 @@ +# SPDX-License-Identifier: 0BSD + +"""SQLite integration: announce spam and bounded storage (anti-exhaustion). + +Multi-minute soak scenarios live in ``test_long_running_stress.py`` (opt-in via +``MESHCHAT_LONG_TEST_SECONDS``). +""" + +from __future__ import annotations + +import os +import tempfile +from unittest.mock import MagicMock + +import pytest + +from meshchatx.src.backend.announce_manager import AnnounceManager +from meshchatx.src.backend.database import Database +from meshchatx.src.backend.database.provider import DatabaseProvider + + +class _FakeIdentity: + __slots__ = ("_h",) + + def __init__(self, identity_hex32: str): + self._h = bytes.fromhex(identity_hex32) + + @property + def hash(self): + return self._h + + def get_public_key(self): + return b"\xaa\xbb" + + +def _cleanup(db, path): + if db is not None: + try: + db.close() + except Exception: + pass + DatabaseProvider._instance = None + if path: + try: + os.unlink(path) + except OSError: + pass + for suffix in ("-wal", "-shm"): + try: + os.unlink(path + suffix) + except OSError: + pass + + +def _new_db(): + with tempfile.NamedTemporaryFile(suffix=".db", delete=False) as f: + path = f.name + db = Database(path) + db.initialize() + return db, path + + +def _store_enabled_config(**max_stored): + config = MagicMock() + for _k in ( + "announce_store_lxmf_delivery", + "announce_store_lxst_telephony", + "announce_store_nomadnetwork_node", + "announce_store_lxmf_propagation", + "announce_store_git_repositories", + ): + m = MagicMock() + m.get.return_value = True + setattr(config, _k, m) + + for key, default in ( + ("announce_max_stored_lxmf_delivery", None), + ("announce_max_stored_nomadnetwork_node", None), + ("announce_max_stored_lxmf_propagation", None), + ): + attr = MagicMock() + attr.get.return_value = max_stored.get(key, default) + setattr(config, key, attr) + + return config + + +@pytest.fixture +def sqlite_db(): + db, path = _new_db() + yield db, path + _cleanup(db, path) + + +def test_spam_unique_destinations_stays_within_cap(sqlite_db): + db, _path = sqlite_db + cap = 48 + cfg = _store_enabled_config(announce_max_stored_lxmf_delivery=cap) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + + total = 650 + for i in range(total): + dh = f"{i:032x}" + mgr.upsert_announce( + ret, + _FakeIdentity(f"{i:032x}"), + bytes.fromhex(dh), + "lxmf.delivery", + b"x", + None, + ) + + assert db.announces.get_announce_count_by_aspect("lxmf.delivery") == cap + + +def test_spam_same_destination_does_not_duplicate_rows(sqlite_db): + db, _path = sqlite_db + cfg = _store_enabled_config(announce_max_stored_lxmf_delivery=12) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + dest = bytes.fromhex("ab" * 16) + ident = _FakeIdentity("cd" * 16) + + for _ in range(400): + mgr.upsert_announce(ret, ident, dest, "lxmf.delivery", b"y", None) + + rows = db.provider.fetchall( + "SELECT COUNT(*) AS n FROM announces WHERE aspect = ? AND destination_hash = ?", + ("lxmf.delivery", "ab" * 16), + ) + assert rows[0]["n"] == 1 + assert db.announces.get_announce_count_by_aspect("lxmf.delivery") == 1 + + +def test_spam_interleaved_aspects_each_bounded(sqlite_db): + db, _path = sqlite_db + cfg = _store_enabled_config( + announce_max_stored_lxmf_delivery=15, + announce_max_stored_nomadnetwork_node=9, + announce_max_stored_lxmf_propagation=11, + ) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + + for round_i in range(120): + i = round_i * 3 + mgr.upsert_announce( + ret, + _FakeIdentity(f"{i + 1:032x}"), + bytes.fromhex(f"{i + 1:032x}"), + "lxmf.delivery", + b"a", + None, + ) + mgr.upsert_announce( + ret, + _FakeIdentity(f"{i + 2:032x}"), + bytes.fromhex(f"{i + 2:032x}"), + "nomadnetwork.node", + b"b", + None, + ) + mgr.upsert_announce( + ret, + _FakeIdentity(f"{i + 3:032x}"), + bytes.fromhex(f"{i + 3:032x}"), + "lxmf.propagation", + b"c", + None, + ) + + assert db.announces.get_announce_count_by_aspect("lxmf.delivery") == 15 + assert db.announces.get_announce_count_by_aspect("nomadnetwork.node") == 9 + assert db.announces.get_announce_count_by_aspect("lxmf.propagation") == 11 + + +def test_quick_check_ok_after_heavy_announce_spam(sqlite_db): + db, _path = sqlite_db + cfg = _store_enabled_config(announce_max_stored_lxmf_delivery=30) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + + for i in range(900): + mgr.upsert_announce( + ret, + _FakeIdentity(f"{i % 100:032x}"), + bytes.fromhex(f"{i:032x}"), + "lxmf.delivery", + os.urandom(64), + None, + ) + + qc = db.provider.quick_check() + assert qc + first = qc[0] + val = next(iter(first.values())) + assert val == "ok" + + +def test_large_app_data_spam_remains_bounded(sqlite_db): + db, _path = sqlite_db + cfg = _store_enabled_config(announce_max_stored_lxmf_delivery=20) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + blob = b"z" * 12000 + + for i in range(180): + mgr.upsert_announce( + ret, + _FakeIdentity(f"{i:032x}"), + bytes.fromhex(f"{i:032x}"), + "lxmf.delivery", + blob, + None, + ) + + assert db.announces.get_announce_count_by_aspect("lxmf.delivery") == 20 diff --git a/tests/backend/test_async_utils_critical.py b/tests/backend/test_async_utils_critical.py new file mode 100644 index 0000000..ffe1711 --- /dev/null +++ b/tests/backend/test_async_utils_critical.py @@ -0,0 +1,108 @@ +# SPDX-License-Identifier: 0BSD + +"""Critical-path tests for ``AsyncUtils``: cross-thread scheduling and memory caps.""" + +from __future__ import annotations + +import asyncio +import threading +import warnings + +import pytest + +from meshchatx.src.backend.async_utils import AsyncUtils + + +@pytest.fixture(autouse=True) +def _reset_async_utils(): + AsyncUtils.main_loop = None + AsyncUtils._pending_futures.clear() + AsyncUtils._pending_coroutines.clear() + yield + AsyncUtils.main_loop = None + with warnings.catch_warnings(): + warnings.simplefilter("ignore", RuntimeWarning) + AsyncUtils._pending_futures.clear() + AsyncUtils._pending_coroutines.clear() + + +async def _noop(): + return None + + +def test_buffered_coroutines_capped_when_event_loop_not_running(): + with warnings.catch_warnings(): + warnings.simplefilter("ignore", RuntimeWarning) + for _ in range(AsyncUtils._COROUTINES_MAX + 12): + AsyncUtils.run_async(_noop()) + + assert len(AsyncUtils._pending_coroutines) == AsyncUtils._COROUTINES_MAX + + +@pytest.mark.asyncio +async def test_set_main_loop_drains_buffered_coroutines(): + seen: list[bool] = [] + + async def record(): + seen.append(True) + + AsyncUtils.main_loop = None + + queued = threading.Event() + + def schedule_from_worker(): + AsyncUtils.run_async(record()) + queued.set() + + threading.Thread(target=schedule_from_worker).start() + assert queued.wait(timeout=2.0) + assert len(AsyncUtils._pending_coroutines) == 1 + + AsyncUtils.set_main_loop(asyncio.get_running_loop()) + assert AsyncUtils._pending_coroutines == [] + + await asyncio.sleep(0.15) + assert seen == [True] + + +@pytest.mark.asyncio +async def test_run_async_with_running_loop_executes_coroutine(): + outcomes: list[int] = [] + + async def work(): + outcomes.append(7) + + AsyncUtils.set_main_loop(asyncio.get_running_loop()) + + done = threading.Event() + + def schedule_from_worker(): + AsyncUtils.run_async(work()) + done.set() + + threading.Thread(target=schedule_from_worker).start() + assert done.wait(timeout=2.0) + await asyncio.sleep(0.15) + assert outcomes == [7] + + +@pytest.mark.asyncio +async def test_pending_futures_list_sheds_completed_entries(): + AsyncUtils.set_main_loop(asyncio.get_running_loop()) + with AsyncUtils._futures_lock: + AsyncUtils._pending_futures.clear() + + finished = threading.Event() + + def blast(): + for _ in range(AsyncUtils._FUTURES_SWEEP_THRESHOLD + 8): + AsyncUtils.run_async(asyncio.sleep(0)) + finished.set() + + threading.Thread(target=blast).start() + assert finished.wait(timeout=5.0) + await asyncio.sleep(0.15) + + with AsyncUtils._futures_lock: + still_pending = [f for f in AsyncUtils._pending_futures if not f.done()] + assert len(still_pending) == 0 diff --git a/tests/backend/test_auto_propagation.py b/tests/backend/test_auto_propagation.py index cbc72c2..f96595a 100644 --- a/tests/backend/test_auto_propagation.py +++ b/tests/backend/test_auto_propagation.py @@ -214,3 +214,24 @@ async def test_auto_propagation_removes_broken_node_when_all_candidates_fail(): app.set_active_propagation_node.assert_not_called() app.remove_active_propagation_node.assert_called_once_with(context=context) + + +@pytest.mark.asyncio +async def test_check_and_update_propagation_node_noops_without_message_router(): + manager, app, context, config, _database = _make_manager() + context.message_router = None + config.lxmf_preferred_propagation_node_auto_select.get.return_value = True + + await manager.check_and_update_propagation_node() + + app.set_active_propagation_node.assert_not_called() + app.remove_active_propagation_node.assert_not_called() + + +def test_stop_propagation_node_sync_noops_when_message_router_none(): + from meshchatx.meshchat import ReticulumMeshChat + + app = ReticulumMeshChat.__new__(ReticulumMeshChat) + ctx = MagicMock() + ctx.message_router = None + ReticulumMeshChat.stop_propagation_node_sync(app, context=ctx) diff --git a/tests/backend/test_interface_discovery.py b/tests/backend/test_interface_discovery.py index 70f481c..2dd9ad6 100644 --- a/tests/backend/test_interface_discovery.py +++ b/tests/backend/test_interface_discovery.py @@ -150,6 +150,40 @@ async def test_reticulum_discovery_get_and_patch(temp_dir): assert config.write_called +@pytest.mark.asyncio +async def test_reticulum_discovery_get_default_bootstrap_false_when_unset(temp_dir): + config = ConfigDict({"reticulum": {}, "interfaces": {}}) + + with ( + patch("meshchatx.meshchat.generate_ssl_certificate"), + patch("RNS.Reticulum") as mock_rns, + patch("RNS.Transport"), + patch("LXMF.LXMRouter"), + ): + mock_reticulum = mock_rns.return_value + mock_reticulum.config = config + mock_reticulum.configpath = "/tmp/mock_config" + mock_reticulum.is_connected_to_shared_instance = False + mock_reticulum.transport_enabled.return_value = True + + app_instance = ReticulumMeshChat( + identity=build_identity(), + storage_dir=temp_dir, + reticulum_config_dir=temp_dir, + ) + + get_handler = await find_route_handler( + app_instance, + "/api/v1/reticulum/discovery", + "GET", + ) + assert get_handler + + get_response = await get_handler(MagicMock()) + get_data = json.loads(get_response.body) + assert get_data["discovery"]["default_bootstrap_only"] is False + + @pytest.mark.asyncio async def test_discovery_patch_rejects_zero_autoconnect_as_unset(temp_dir): config = ConfigDict( @@ -409,7 +443,7 @@ async def test_interface_add_includes_discovery_fields(temp_dir): assert saved["discovery_frequency"] == 915000000 assert saved["discovery_bandwidth"] == 125000 assert saved["discovery_modulation"] == "LoRa" - assert saved.get("bootstrap_only") == "yes" + assert "bootstrap_only" not in saved assert config.write_called @@ -516,6 +550,66 @@ async def test_interface_add_tcp_explicit_bootstrap_only_no(temp_dir): assert config["interfaces"]["ExplicitNo"]["bootstrap_only"] == "no" +@pytest.mark.asyncio +async def test_interface_edit_tcp_preserves_bootstrap_when_key_omitted(temp_dir): + config = ConfigDict( + { + "reticulum": {"default_bootstrap_only": True}, + "interfaces": { + "KeepBoot": { + "type": "TCPClientInterface", + "target_host": "example.com", + "target_port": "4242", + "bootstrap_only": "yes", + }, + }, + }, + ) + + with ( + patch("meshchatx.meshchat.generate_ssl_certificate"), + patch("RNS.Reticulum") as mock_rns, + patch("RNS.Transport"), + patch("LXMF.LXMRouter"), + ): + mock_reticulum = mock_rns.return_value + mock_reticulum.config = config + mock_reticulum.configpath = "/tmp/mock_config" + mock_reticulum.is_connected_to_shared_instance = False + mock_reticulum.transport_enabled.return_value = True + + app_instance = ReticulumMeshChat( + identity=build_identity(), + storage_dir=temp_dir, + reticulum_config_dir=temp_dir, + ) + + add_handler = await find_route_handler( + app_instance, + "/api/v1/reticulum/interfaces/add", + "POST", + ) + assert add_handler + + payload = { + "allow_overwriting_interface": True, + "name": "KeepBoot", + "type": "TCPClientInterface", + "target_host": "example.com", + "target_port": "4242", + } + + class AddRequest: + @staticmethod + async def json(): + return payload + + response = await add_handler(AddRequest()) + data = json.loads(response.body) + assert "message" in data + assert config["interfaces"]["KeepBoot"]["bootstrap_only"] == "yes" + + def test_apply_bootstrap_only_to_interface(): details = {} ReticulumMeshChat.apply_bootstrap_only_to_interface(details, {}, True) @@ -531,6 +625,35 @@ def test_apply_bootstrap_only_to_interface(): ReticulumMeshChat.apply_bootstrap_only_to_interface(details, {}, False) assert "bootstrap_only" not in details + details = {"bootstrap_only": "yes"} + ReticulumMeshChat.apply_bootstrap_only_to_interface( + details, {}, True, updating_existing=True + ) + assert details["bootstrap_only"] == "yes" + + +def test_strip_reload_instance_suffix(): + assert ReticulumMeshChat._strip_reload_instance_suffix(None) is None + assert ReticulumMeshChat._strip_reload_instance_suffix("") is None + assert ReticulumMeshChat._strip_reload_instance_suffix("mesh") == "mesh" + assert ReticulumMeshChat._strip_reload_instance_suffix( + "production-reload-backend" + ) == ("production-reload-backend") + assert ( + ReticulumMeshChat._strip_reload_instance_suffix("my-net-reload-peer") + == "my-net-reload-peer" + ) + assert ( + ReticulumMeshChat._strip_reload_instance_suffix("node-reload-1-500") + == "node-reload-1-500" + ) + assert ( + ReticulumMeshChat._strip_reload_instance_suffix( + "default-reload-2246687-1777566181-reload-3009314-1777566481", + ) + == "default" + ) + @pytest.mark.asyncio async def test_interface_add_discoverable_without_optional_coordinates(temp_dir): diff --git a/tests/backend/test_long_running_stress.py b/tests/backend/test_long_running_stress.py new file mode 100644 index 0000000..4841a88 --- /dev/null +++ b/tests/backend/test_long_running_stress.py @@ -0,0 +1,247 @@ +# SPDX-License-Identifier: 0BSD + +"""Multi-minute soak tests (announce DB + websocket fan-out). + +These are **opt-in**: unset ``MESHCHAT_LONG_TEST_SECONDS`` skips them immediately. + +Examples:: + + MESHCHAT_LONG_TEST_SECONDS=300 uv run pytest tests/backend/test_long_running_stress.py -m long_running -v + MESHCHAT_LONG_TEST_SECONDS=600 uv run pytest tests/backend/test_long_running_stress.py -m long_running -v + +Quick smoke (seconds):: + + MESHCHAT_LONG_TEST_SECONDS=5 uv run pytest tests/backend/test_long_running_stress.py -m long_running -v +""" + +from __future__ import annotations + +import os +import tempfile +import time +from unittest.mock import AsyncMock, MagicMock + +import pytest + +from meshchatx.meshchat import ReticulumMeshChat +from meshchatx.src.backend.announce_manager import AnnounceManager +from meshchatx.src.backend.database import Database +from meshchatx.src.backend.database.provider import DatabaseProvider + + +class _FakeIdentity: + __slots__ = ("_h",) + + def __init__(self, identity_hex32: str): + self._h = bytes.fromhex(identity_hex32) + + @property + def hash(self): + return self._h + + def get_public_key(self): + return b"\xaa\xbb" + + +def _cleanup(db, path): + if db is not None: + try: + db.close() + except Exception: + pass + DatabaseProvider._instance = None + if path: + try: + os.unlink(path) + except OSError: + pass + for suffix in ("-wal", "-shm"): + try: + os.unlink(path + suffix) + except OSError: + pass + + +def _new_db(): + with tempfile.NamedTemporaryFile(suffix=".db", delete=False) as f: + path = f.name + db = Database(path) + db.initialize() + return db, path + + +def _store_enabled_config(**max_stored): + config = MagicMock() + for _k in ( + "announce_store_lxmf_delivery", + "announce_store_lxst_telephony", + "announce_store_nomadnetwork_node", + "announce_store_lxmf_propagation", + "announce_store_git_repositories", + ): + m = MagicMock() + m.get.return_value = True + setattr(config, _k, m) + + for key, default in ( + ("announce_max_stored_lxmf_delivery", None), + ("announce_max_stored_nomadnetwork_node", None), + ("announce_max_stored_lxmf_propagation", None), + ): + attr = MagicMock() + attr.get.return_value = max_stored.get(key, default) + setattr(config, key, attr) + + return config + + +def _long_test_seconds() -> float: + raw = os.environ.get("MESHCHAT_LONG_TEST_SECONDS", "").strip() + if not raw: + return 0.0 + try: + return float(raw) + except ValueError: + return 0.0 + + +def _require_long_duration(): + sec = _long_test_seconds() + if sec <= 0: + pytest.skip( + "Set MESHCHAT_LONG_TEST_SECONDS to a positive value " + "(e.g. 300 for 5 minutes, 600 for 10 minutes).", + ) + return sec + + +def _bind_real_websocket_broadcast(app): + return ReticulumMeshChat.websocket_broadcast.__get__(app, ReticulumMeshChat) + + +class _MagicWs: + __slots__ = ("send_str",) + + def __init__(self): + self.send_str = AsyncMock(return_value=None) + + +@pytest.mark.long_running +def test_soak_sqlite_announces_stay_bounded_and_quick_check(): + duration_s = _require_long_duration() + cap = 96 + batch_size = 120 + qc_interval_batches = 15 + + db, path = _new_db() + try: + cfg = _store_enabled_config(announce_max_stored_lxmf_delivery=cap) + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + deadline = time.monotonic() + duration_s + seq = 0 + batches = 0 + while time.monotonic() < deadline: + for _ in range(batch_size): + mgr.upsert_announce( + ret, + _FakeIdentity(f"{seq % 2048:032x}"), + bytes.fromhex(f"{seq:032x}"), + "lxmf.delivery", + os.urandom(48), + None, + ) + seq += 1 + batches += 1 + assert db.announces.get_announce_count_by_aspect("lxmf.delivery") == cap + if batches % qc_interval_batches == 0: + qc = db.provider.quick_check() + assert qc + assert next(iter(qc[0].values())) == "ok" + assert seq > 0 + finally: + _cleanup(db, path) + + +@pytest.mark.long_running +def test_soak_interleaved_aspects_under_cap(): + duration_s = _require_long_duration() + cfg = _store_enabled_config( + announce_max_stored_lxmf_delivery=40, + announce_max_stored_nomadnetwork_node=25, + announce_max_stored_lxmf_propagation=30, + ) + + db, path = _new_db() + try: + mgr = AnnounceManager(db, cfg) + ret = MagicMock() + deadline = time.monotonic() + duration_s + n = 0 + while time.monotonic() < deadline: + for aspect, payload in ( + ("lxmf.delivery", b"a"), + ("nomadnetwork.node", b"b"), + ("lxmf.propagation", b"c"), + ): + mgr.upsert_announce( + ret, + _FakeIdentity(f"{n % 900:032x}"), + bytes.fromhex(f"{n:032x}"), + aspect, + payload, + None, + ) + n += 1 + assert db.announces.get_announce_count_by_aspect("lxmf.delivery") <= 40 + assert db.announces.get_announce_count_by_aspect("nomadnetwork.node") <= 25 + assert db.announces.get_announce_count_by_aspect("lxmf.propagation") <= 30 + assert n > 0 + assert db.announces.get_announce_count_by_aspect("lxmf.delivery") == 40 + assert db.announces.get_announce_count_by_aspect("nomadnetwork.node") == 25 + assert db.announces.get_announce_count_by_aspect("lxmf.propagation") == 30 + finally: + _cleanup(db, path) + + +@pytest.mark.long_running +@pytest.mark.asyncio +async def test_soak_websocket_broadcast_under_load(mock_app): + duration_s = _require_long_duration() + mock_app.websocket_clients.clear() + n_clients = min(int(os.environ.get("MESHCHAT_LONG_TEST_WS_CLIENTS", "400")), 2000) + clients = [_MagicWs() for _ in range(n_clients)] + mock_app.websocket_clients.extend(clients) + real = _bind_real_websocket_broadcast(mock_app) + + deadline = time.monotonic() + duration_s + rounds = 0 + while time.monotonic() < deadline: + payload = f'{{"type":"soak","round":{rounds}}}' + await real(payload) + for c in clients: + assert c.send_str.await_args[0][0] == payload + rounds += 1 + assert rounds > 0 + + +@pytest.mark.long_running +@pytest.mark.asyncio +async def test_soak_websocket_broadcast_with_churn(mock_app): + duration_s = _require_long_duration() + real = _bind_real_websocket_broadcast(mock_app) + deadline = time.monotonic() + duration_s + wave = 0 + while time.monotonic() < deadline: + mock_app.websocket_clients.clear() + batch = [_MagicWs() for _ in range(80)] + if wave % 3 == 0: + for c in batch[:20]: + c.send_str = AsyncMock(side_effect=ConnectionError("closed")) + mock_app.websocket_clients.extend(batch) + payload = f'{{"wave":{wave}}}' + await real(payload) + for c in mock_app.websocket_clients: + assert c.send_str.await_args[0][0] == payload + wave += 1 + assert wave > 0 diff --git a/tests/backend/test_search_integration.py b/tests/backend/test_search_integration.py index 759a569..77d32d8 100644 --- a/tests/backend/test_search_integration.py +++ b/tests/backend/test_search_integration.py @@ -85,6 +85,27 @@ def test_filter_announced_dicts_by_search_query_destination_hash_substring(): assert len(out) == 1 +def test_filter_announced_dicts_empty_search_matches_all(): + """Empty substring matches every string in Python; callers should normalize UI input.""" + items = [ + {"display_name": "AAA"}, + {"destination_hash": "0123abcd"}, + {"identity_hash": "fedcba"}, + ] + out = filter_announced_dicts_by_search_query(items, "") + assert len(out) == 3 + + +def test_filter_announced_dicts_whitespace_search_can_match(): + items = [ + {"display_name": " Hi "}, + {"destination_hash": "99"}, + ] + out = filter_announced_dicts_by_search_query(items, " ") + assert len(out) == 1 + assert out[0]["display_name"] == " Hi " + + def test_filter_announced_dicts_by_search_query_case_insensitive(): items = [ {"display_name": "CamelCaseName", "destination_hash": "z" * 32}, diff --git a/tests/backend/test_websocket_scale.py b/tests/backend/test_websocket_scale.py index 1ade132..29becb3 100644 --- a/tests/backend/test_websocket_scale.py +++ b/tests/backend/test_websocket_scale.py @@ -109,6 +109,57 @@ async def test_websocket_broadcast_drops_dead_clients(mock_app): assert good.send_str.await_count == 1 +@pytest.mark.asyncio +async def test_websocket_broadcast_all_clients_dead_empties_list(mock_app): + mock_app.websocket_clients.clear() + clients = [] + for _ in range(80): + c = MagicWs() + c.send_str = AsyncMock(side_effect=ConnectionError("closed")) + clients.append(c) + mock_app.websocket_clients.extend(clients) + + real = _bind_real_websocket_broadcast(mock_app) + await real("{}") + + assert mock_app.websocket_clients == [] + + +@pytest.mark.asyncio +async def test_websocket_broadcast_fanout_large_client_pool(mock_app): + mock_app.websocket_clients.clear() + n = 1500 + clients = [MagicWs() for _ in range(n)] + mock_app.websocket_clients.extend(clients) + + real = _bind_real_websocket_broadcast(mock_app) + payload = '{"type":"metrics","x":1}' + await real(payload) + + for c in clients: + assert c.send_str.await_count == 1 + assert c.send_str.await_args[0][0] == payload + + +@pytest.mark.asyncio +async def test_websocket_broadcast_mixed_failures_still_delivers_to_healthy(mock_app): + mock_app.websocket_clients.clear() + failing = [MagicWs() for _ in range(40)] + for c in failing: + c.send_str = AsyncMock(side_effect=BrokenPipeError()) + healthy = [MagicWs() for _ in range(60)] + mock_app.websocket_clients.extend(failing + healthy) + + real = _bind_real_websocket_broadcast(mock_app) + await real('{"type":"ping"}') + + for c in failing: + assert c not in mock_app.websocket_clients + assert mock_app.websocket_clients == healthy + for c in healthy: + assert c.send_str.await_count == 1 + + class MagicWs: def __init__(self): self.send_str = AsyncMock(return_value=None) From 02b8695726f59e1d59743cdb32e43c780b2756ef Mon Sep 17 00:00:00 2001 From: Ivan Date: Thu, 7 May 2026 20:04:57 -0500 Subject: [PATCH 05/63] refactor(ConversationViewer): streamline message loading UI --- .../messages/ConversationViewer.vue | 294 ++++++++++++------ .../ConversationViewer.scroll.test.js | 39 ++- 2 files changed, 238 insertions(+), 95 deletions(-) diff --git a/meshchatx/src/frontend/components/messages/ConversationViewer.vue b/meshchatx/src/frontend/components/messages/ConversationViewer.vue index f55ce80..f4f1394 100644 --- a/meshchatx/src/frontend/components/messages/ConversationViewer.vue +++ b/meshchatx/src/frontend/components/messages/ConversationViewer.vue @@ -2,10 +2,7 @@