From 417390862b843df92650adcdb6d0d7cd8b698e21 Mon Sep 17 00:00:00 2001 From: Matt Johnson Date: Thu, 2 Jul 2026 16:44:20 +0000 Subject: [PATCH] refactor(sizing): single mesh packet budget via transport.max_chars (Phase 3) Route all mesh message sizing (renderer, digest, reply chunker, NWS one-packet fit) through the active transport's max_chars instead of scattered literal 200s. Meshtastic pinned at 200 (byte-identical output); MeshCore uses its configured ~140. Sets up uniform-to-smaller sizing for the composite transport. No behavior change on the Meshtastic path. Co-Authored-By: Claude Opus 4.8 (1M context) --- work/meshai/adapter_config/__init__.py | 2 + work/meshai/adapter_config/_accessor.py | 15 +- work/meshai/config.py | 1 + work/meshai/connector.py | 4 + work/meshai/main.py | 8 +- work/meshai/notifications/channels.py | 6 +- .../meshai/notifications/pipeline/__init__.py | 1 + work/meshai/notifications/router.py | 7 +- work/meshai/router.py | 2 +- work/meshai/transport/base.py | 6 + work/meshai/transport/meshcore_transport.py | 4 + work/tests/test_uniform_sizing.py | 244 ++++++++++++++++++ 12 files changed, 292 insertions(+), 8 deletions(-) create mode 100644 work/tests/test_uniform_sizing.py diff --git a/work/meshai/adapter_config/__init__.py b/work/meshai/adapter_config/__init__.py index 83c513f..88611ec 100644 --- a/work/meshai/adapter_config/__init__.py +++ b/work/meshai/adapter_config/__init__.py @@ -31,6 +31,7 @@ from typing import Any from meshai.adapter_config._accessor import ( adapter_config, invalidate_cache, + set_runtime_override, ) from meshai.adapter_config.defaults import ( REGISTRY, @@ -42,6 +43,7 @@ from meshai.adapter_config.defaults import ( __all__ = [ "adapter_config", "invalidate_cache", + "set_runtime_override", "seed_defaults", "prune_orphans", "REGISTRY", diff --git a/work/meshai/adapter_config/_accessor.py b/work/meshai/adapter_config/_accessor.py index 77583c3..c06a277 100644 --- a/work/meshai/adapter_config/_accessor.py +++ b/work/meshai/adapter_config/_accessor.py @@ -36,6 +36,16 @@ logger = logging.getLogger(__name__) _CACHE_LOCK = threading.Lock() _cache: dict[tuple[str, str], Any] = {} +# Runtime-override store: values injected here survive cache invalidation. +# Keyed by the same (adapter, key) tuple as _cache. Never cleared by +# invalidate_cache(); cleared only by explicit deletion or process exit. +_overrides: dict[tuple[str, str], Any] = {} + + +def set_runtime_override(adapter: str, key: str, value: Any) -> None: + """Persist a runtime-derived adapter_config value that survives cache invalidation.""" + _overrides[(adapter, key)] = value + def invalidate_cache() -> None: """Drop every cached value. Called by the REST API on PUT/reset.""" @@ -143,8 +153,11 @@ class AdapterConfig: def _resolve(adapter: str, key: str) -> Any: - """Read pipeline: cache -> DB -> registry.""" + """Read pipeline: overrides -> cache -> DB -> registry.""" cache_key = (adapter, key) + override = _overrides.get(cache_key, _SENTINEL) + if override is not _SENTINEL: + return override cached = _cache.get(cache_key, _SENTINEL) if cached is not _SENTINEL: return cached diff --git a/work/meshai/config.py b/work/meshai/config.py index 1c0b23d..69205d9 100644 --- a/work/meshai/config.py +++ b/work/meshai/config.py @@ -36,6 +36,7 @@ class ConnectionConfig: reconnect_health_interval: float = 30.0 # --- transport selection (Phase 1 seam; MeshCore support is Phase 2) --- transport: str = "meshtastic" # "meshtastic" | "meshcore" | "both" + meshtastic_max_chars: int = 200 # default Meshtastic packet budget # --- MeshCore transport settings (used when transport="meshcore") --- meshcore_host: str = "100.64.0.9" # pyMC companion frame server host meshcore_port: int = 5050 # pyMC companion frame server port diff --git a/work/meshai/connector.py b/work/meshai/connector.py index 3b6f619..5771102 100644 --- a/work/meshai/connector.py +++ b/work/meshai/connector.py @@ -72,6 +72,10 @@ class MeshtasticTransport(MeshTransport): """Get our node's ID.""" return self._my_node_id + @property + def max_chars(self) -> int: + return getattr(self.config, "meshtastic_max_chars", 200) + def connect(self) -> None: """Establish connection to Meshtastic node.""" logger.info(f"Connecting to Meshtastic node via {self.config.type}...") diff --git a/work/meshai/main.py b/work/meshai/main.py index 17bc81e..2900823 100644 --- a/work/meshai/main.py +++ b/work/meshai/main.py @@ -381,6 +381,12 @@ class MeshAI: # Transport connector (factory selects backend from config.connection.transport) self.connector = build_transport(self.config.connection) + # Fit the NWS one-packet formatter to the active mesh transport's budget + # (Meshtastic default 200 -> unchanged). Durable across adapter_config cache + # invalidation via the runtime-override store. + from meshai.adapter_config import set_runtime_override + set_runtime_override("nws", "single_packet_max_chars", self.connector.max_chars) + # Passive mesh context buffer ctx_cfg = self.config.context if ctx_cfg.enabled: @@ -666,7 +672,7 @@ class MeshAI: from .chunker import chunk_response messages, remaining = chunk_response( result.response, - max_chars=self.config.response.max_length, + max_chars=min(self.config.response.max_length, self.connector.max_chars), max_messages=self.config.response.max_messages, ) if remaining: diff --git a/work/meshai/notifications/channels.py b/work/meshai/notifications/channels.py index 6ac39ff..9b1a3c2 100644 --- a/work/meshai/notifications/channels.py +++ b/work/meshai/notifications/channels.py @@ -63,7 +63,8 @@ class MeshBroadcastChannel(NotificationChannel): def __init__(self, connector: "MeshConnector", channel_index: int = 0): self._connector = connector self._channel = channel_index - self._renderer = MeshRenderer() + _mc = getattr(connector, "max_chars", 200) + self._renderer = MeshRenderer(char_limit=_mc if isinstance(_mc, int) else 200) async def deliver(self, alert: "NotificationPayload", rule: "NotificationRuleConfig") -> bool: """Send alert to mesh channel.""" @@ -174,7 +175,8 @@ class MeshDMChannel(NotificationChannel): def __init__(self, connector: "MeshConnector", node_ids: list[str]): self._connector = connector self._node_ids = node_ids - self._renderer = MeshRenderer() + _mc = getattr(connector, "max_chars", 200) + self._renderer = MeshRenderer(char_limit=_mc if isinstance(_mc, int) else 200) async def deliver(self, alert: "NotificationPayload", rule: "NotificationRuleConfig") -> bool: """Send alert via DM to configured nodes.""" diff --git a/work/meshai/notifications/pipeline/__init__.py b/work/meshai/notifications/pipeline/__init__.py index 5b2a244..7f0e67b 100644 --- a/work/meshai/notifications/pipeline/__init__.py +++ b/work/meshai/notifications/pipeline/__init__.py @@ -74,6 +74,7 @@ def build_pipeline(config, llm_backend, connector=None) -> EventBus: accumulator = DigestAccumulator( llm_backend=llm_backend, include_toggles=include_toggles, + mesh_char_limit=connector.max_chars if connector is not None else 200, ) # Tee closure: events go to BOTH dispatcher and accumulator diff --git a/work/meshai/notifications/router.py b/work/meshai/notifications/router.py index 80f033c..507cdf3 100644 --- a/work/meshai/notifications/router.py +++ b/work/meshai/notifications/router.py @@ -185,12 +185,13 @@ class NotificationRouter: delivery_alert = alert message = alert.get("message", "") if channel.channel_type in ("mesh_broadcast", "mesh_dm"): - if len(message) > 200: + _budget = getattr(self._connector, "max_chars", 200) + if len(message) > _budget: if self._summarizer: - summary = await self._summarizer.summarize(message, max_chars=195) + summary = await self._summarizer.summarize(message, max_chars=_budget - 5) delivery_alert = {**alert, "message": summary} else: - delivery_alert = {**alert, "message": message[:195] + "..."} + delivery_alert = {**alert, "message": message[:_budget - 5] + "..."} # Convert dict to NotificationPayload for channel interface payload = NotificationPayload( diff --git a/work/meshai/router.py b/work/meshai/router.py index a912e0b..f6db43f 100644 --- a/work/meshai/router.py +++ b/work/meshai/router.py @@ -925,7 +925,7 @@ class MessageRouter: # Chunk the response with sentence awareness messages, remaining = chunk_response( response, - max_chars=self.config.response.max_length, + max_chars=min(self.config.response.max_length, self.connector.max_chars), max_messages=self.config.response.max_messages, ) diff --git a/work/meshai/transport/base.py b/work/meshai/transport/base.py index 7171c30..43cccd6 100644 --- a/work/meshai/transport/base.py +++ b/work/meshai/transport/base.py @@ -77,6 +77,12 @@ class MeshTransport(abc.ABC): def connected(self) -> bool: """True when the transport has an active connection.""" + @property + @abc.abstractmethod + def max_chars(self) -> int: + """Maximum characters per mesh packet for the active transport.""" + ... + @abc.abstractmethod def get_node_name(self, node_id: str) -> str: """Return the cached display name for *node_id*, or *node_id* itself.""" diff --git a/work/meshai/transport/meshcore_transport.py b/work/meshai/transport/meshcore_transport.py index 0ea7027..4f06c28 100644 --- a/work/meshai/transport/meshcore_transport.py +++ b/work/meshai/transport/meshcore_transport.py @@ -362,6 +362,10 @@ class MeshCoreTransport(MeshTransport): """True when the transport has an active connection.""" return self._connected and self._mc is not None + @property + def max_chars(self) -> int: + return getattr(self.config, "meshcore_max_chars", 140) + def get_node_name(self, node_id: str) -> str: """Resolve a pubkey prefix to a contact display name, or return node_id.""" if self._mc is None: diff --git a/work/tests/test_uniform_sizing.py b/work/tests/test_uniform_sizing.py new file mode 100644 index 0000000..e2f3a92 --- /dev/null +++ b/work/tests/test_uniform_sizing.py @@ -0,0 +1,244 @@ +"""Tests for Phase 3 uniform mesh-packet sizing. + +Verifies that every size-sensitive code path uses the active transport's +max_chars rather than a hardcoded 200. Meshtastic default stays 200 +(byte-identical); MeshCore uses 140. +""" + +import sys +import types +from unittest.mock import MagicMock + +import pytest + +# --------------------------------------------------------------------------- +# Fake meshcore module (must precede production imports that lazy-import it) +# --------------------------------------------------------------------------- + +def _build_fake_meshcore(): + mod = types.ModuleType("meshcore") + + class EventType: + CONTACT_MSG_RECV = "CONTACT_MSG_RECV" + CHANNEL_MSG_RECV = "CHANNEL_MSG_RECV" + DISCONNECTED = "DISCONNECTED" + CONNECTED = "CONNECTED" + + mod.EventType = EventType + + class _FakeMeshCore: + self_info = {"public_key": "aabbccdd1122", "name": "FakeNode"} + contacts = {} + + async def start_auto_message_fetching(self): + pass + + async def stop_auto_message_fetching(self): + pass + + async def disconnect(self): + pass + + def subscribe(self, event_type, callback): + pass + + def get_contact_by_key_prefix(self, prefix): + return None + + @classmethod + async def create_tcp(cls, host, port, + auto_reconnect=True, max_reconnect_attempts=5): + return cls() + + class commands: + @staticmethod + async def send_chan_msg(chan_idx, text): + result = MagicMock() + result.is_error.return_value = False + return result + + @staticmethod + async def send_msg(dst, text): + result = MagicMock() + result.is_error.return_value = False + return result + + mod.MeshCore = _FakeMeshCore + return mod + + +sys.modules.setdefault("meshcore", _build_fake_meshcore()) + +# --------------------------------------------------------------------------- +# Production imports +# --------------------------------------------------------------------------- + +from meshai.config import ConnectionConfig # noqa: E402 +from meshai.connector import MeshtasticTransport # noqa: E402 +from meshai.transport.meshcore_transport import MeshCoreTransport # noqa: E402 +from meshai.transport.factory import build_transport # noqa: E402 +from meshai.notifications.channels import ( # noqa: E402 + MeshBroadcastChannel, + MeshDMChannel, +) +from meshai.notifications.pipeline import build_pipeline # noqa: E402 + + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + +def _mt_config(**overrides): + """Return a ConnectionConfig wired for meshtastic (default).""" + cfg = ConnectionConfig(transport="meshtastic") + for k, v in overrides.items(): + setattr(cfg, k, v) + return cfg + + +def _mc_config(**overrides): + """Return a ConnectionConfig wired for meshcore.""" + cfg = ConnectionConfig( + transport="meshcore", + meshcore_host="127.0.0.1", + meshcore_port=5050, + meshcore_channel_index=0, + ) + for k, v in overrides.items(): + setattr(cfg, k, v) + return cfg + + +def _fake_connector(max_chars_val: int): + """Minimal fake connector exposing only max_chars.""" + conn = MagicMock() + conn.max_chars = max_chars_val + return conn + + +def _minimal_config(): + """Minimal full Config for build_pipeline calls.""" + from meshai.config import Config + return Config() + + +# --------------------------------------------------------------------------- +# 1. Transport max_chars property +# --------------------------------------------------------------------------- + +class TestTransportMaxChars: + def test_meshtastic_default_is_200(self): + cfg = _mt_config() + t = MeshtasticTransport(cfg) + assert t.max_chars == 200 + + def test_meshtastic_respects_config_field(self): + cfg = _mt_config(meshtastic_max_chars=160) + t = MeshtasticTransport(cfg) + assert t.max_chars == 160 + + def test_meshcore_default_is_140(self): + cfg = _mc_config() + t = MeshCoreTransport(cfg) + assert t.max_chars == 140 + + def test_meshcore_respects_config_field(self): + cfg = _mc_config(meshcore_max_chars=120) + t = MeshCoreTransport(cfg) + assert t.max_chars == 120 + + def test_build_transport_meshtastic_max_chars(self): + t = build_transport(_mt_config()) + assert t.max_chars == 200 + + def test_build_transport_meshcore_max_chars(self): + t = build_transport(_mc_config()) + assert t.max_chars == 140 + + +# --------------------------------------------------------------------------- +# 2. MeshRenderer char_limit propagation via channels +# --------------------------------------------------------------------------- + +class TestChannelRendererBudget: + def test_broadcast_channel_with_meshcore_connector_uses_140(self): + conn = _fake_connector(140) + ch = MeshBroadcastChannel(connector=conn, channel_index=0) + assert ch._renderer._limit == 140 + + def test_broadcast_channel_with_meshtastic_connector_uses_200(self): + conn = _fake_connector(200) + ch = MeshBroadcastChannel(connector=conn, channel_index=0) + assert ch._renderer._limit == 200 + + def test_dm_channel_with_meshcore_connector_uses_140(self): + conn = _fake_connector(140) + ch = MeshDMChannel(connector=conn, node_ids=["!aabbccdd"]) + assert ch._renderer._limit == 140 + + def test_dm_channel_with_meshtastic_connector_uses_200(self): + conn = _fake_connector(200) + ch = MeshDMChannel(connector=conn, node_ids=["!aabbccdd"]) + assert ch._renderer._limit == 200 + + +# --------------------------------------------------------------------------- +# 3. DigestAccumulator mesh_char_limit propagation via build_pipeline +# --------------------------------------------------------------------------- + +class TestBuildPipelineMeshCharLimit: + def test_no_connector_uses_200(self): + cfg = _minimal_config() + bus = build_pipeline(cfg, llm_backend=None, connector=None) + acc = bus._pipeline_components["accumulator"] + assert acc._mesh_char_limit == 200 + + def test_meshcore_connector_uses_140(self): + cfg = _minimal_config() + conn = _fake_connector(140) + bus = build_pipeline(cfg, llm_backend=None, connector=conn) + acc = bus._pipeline_components["accumulator"] + assert acc._mesh_char_limit == 140 + + def test_meshtastic_connector_uses_200(self): + cfg = _minimal_config() + conn = _fake_connector(200) + bus = build_pipeline(cfg, llm_backend=None, connector=conn) + acc = bus._pipeline_components["accumulator"] + assert acc._mesh_char_limit == 200 + + +# --------------------------------------------------------------------------- +# 4. Runtime-override durability: survives adapter_config cache invalidation +# --------------------------------------------------------------------------- + +class TestRuntimeOverrideDurability: + """Verify set_runtime_override writes are not cleared by invalidate_cache(). + + We test at the _resolve/_overrides seam directly because the full + adapter_config.nws.single_packet_max_chars read path hits the DB, which + is not available in this hermetic test environment. Testing via _resolve + is equivalent: it is the exact call site used by _AdapterSection.__getattr__. + """ + + def test_override_survives_cache_invalidation(self): + from meshai.adapter_config._accessor import ( + _overrides, + _resolve, + set_runtime_override, + invalidate_cache, + _SENTINEL, + ) + key = ("nws", "single_packet_max_chars") + try: + set_runtime_override("nws", "single_packet_max_chars", 140) + + # Override is visible through the resolution path. + assert _resolve("nws", "single_packet_max_chars") == 140 + + # Cache invalidation must NOT wipe the override. + invalidate_cache() + assert _resolve("nws", "single_packet_max_chars") == 140 + finally: + # Clean up so this test cannot pollute subsequent tests. + _overrides.pop(key, None)