diff --git a/.gitea/workflows/release.yaml b/.gitea/workflows/release.yaml index 71e8bc4..9f57c7a 100644 --- a/.gitea/workflows/release.yaml +++ b/.gitea/workflows/release.yaml @@ -8,8 +8,49 @@ on: - 'v*' jobs: + # Gateway contract gate: the plugin must bind and wire against the REAL + # Hermes gateway API (pinned to the deployed image source) before any + # release can be built or published. Without this, adapter/API drift used + # to surface only as a gateway crash after deploy (the set_message_handler + # AttributeError). HERMES_CONTRACT_REQUIRED=1 turns a missing hermes-agent + # into a hard failure so this gate can never silently skip. + contract-test: + runs-on: ubuntu-latest + steps: + - name: Checkout code + uses: actions/checkout@v4 + + - name: Set up Python 3.13 (matches Hermes image) + uses: actions/setup-python@v5 + with: + python-version: '3.13' + + - name: Install uv + run: pip install uv + + - name: Clone pinned Hermes gateway source (deployed image tag) + run: | + git clone --depth 1 --branch v2026.8.31 \ + https://github.com/NousResearch/hermes-agent.git plugin/.hermes-src + test "$(git -C plugin/.hermes-src rev-parse HEAD)" = \ + "29112bef099274229cadff79cdff7bf7b99c4b77" + + - name: Editable install hermes-agent + pytest + run: | + uv venv --python 3.13 plugin/.venv-contract + uv pip install --python plugin/.venv-contract/bin/python -e plugin/.hermes-src + uv pip install --python plugin/.venv-contract/bin/python pytest + + - name: Run gateway contract tests + env: + HERMES_CONTRACT_REQUIRED: '1' + run: | + cd plugin + .venv-contract/bin/python -m pytest tests/contract -q + test-and-release: runs-on: ubuntu-latest + needs: contract-test steps: - name: Checkout code uses: actions/checkout@v4 @@ -40,13 +81,13 @@ jobs: echo "Publishing release for $tag_name" wheel_file=$(ls plugin/dist/*.whl | head -n 1) tar_file=$(ls plugin/dist/*.tar.gz | head -n 1) - + # Create release response=$(curl -s -k -X POST "${{ github.server_url }}/api/v1/repos/${{ github.repository }}/releases" \ -H "Authorization: token $GITEA_TOKEN" \ -H "Content-Type: application/json" \ -d "{\"tag_name\": \"$tag_name\", \"name\": \"$tag_name\", \"body\": \"Release $tag_name for hermes-meshtastic plugin\"}") - + release_id=$(echo "$response" | grep -o '"id":[0-9]*' | head -n 1 | cut -d: -f2) if [ -n "$release_id" ]; then curl -s -k -X POST "${{ github.server_url }}/api/v1/repos/${{ github.repository }}/releases/$release_id/assets?name=$(basename $wheel_file)" \ diff --git a/.gitignore b/.gitignore index b240a97..5c53fd1 100644 --- a/.gitignore +++ b/.gitignore @@ -4,3 +4,7 @@ .pio/ *.log plugin/dist/ +plugin/.venv-contract/ +plugin/.hermes-src/ +__pycache__/ +*.pyc diff --git a/plugin/.gitignore b/plugin/.gitignore new file mode 100644 index 0000000..6617f8d --- /dev/null +++ b/plugin/.gitignore @@ -0,0 +1,6 @@ +dist/ +__pycache__/ +*.pyc +.venv-contract/ +.hermes-src/ +.pytest_cache/ diff --git a/plugin/hermes_meshtastic/__init__.py b/plugin/hermes_meshtastic/__init__.py index 3492fab..395e396 100644 --- a/plugin/hermes_meshtastic/__init__.py +++ b/plugin/hermes_meshtastic/__init__.py @@ -1,16 +1,23 @@ -"""Meshtastic Platform Adapter plugin for Hermes Agent.""" +"""Meshtastic Platform Adapter plugin for Hermes Agent. -from .adapter import MeshtasticAdapter, check_requirements, MAX_LORA_MESSAGE_LENGTH +The ``register(ctx)`` entry point resolves the adapter class lazily so the real +Hermes gateway API (``gateway.platforms.base``) is imported only at registration +time — never at plugin-discovery import time (see ``adapter.get_adapter_class``). +""" + +from .formatting import MAX_LORA_MESSAGE_LENGTH def register(ctx): """Hermes plugin entrypoint.""" - from gateway.platform_registry import PlatformEntry + from .adapter import check_requirements, get_adapter_class + + adapter_cls = get_adapter_class() ctx.register_platform( name="meshtastic", label="Meshtastic LoRa", - adapter_factory=lambda config: MeshtasticAdapter(config), + adapter_factory=lambda config: adapter_cls(config), check_fn=check_requirements, max_message_length=MAX_LORA_MESSAGE_LENGTH, emoji="📻", diff --git a/plugin/hermes_meshtastic/__pycache__/__init__.cpython-314.pyc b/plugin/hermes_meshtastic/__pycache__/__init__.cpython-314.pyc deleted file mode 100644 index 87d4515..0000000 Binary files a/plugin/hermes_meshtastic/__pycache__/__init__.cpython-314.pyc and /dev/null differ diff --git a/plugin/hermes_meshtastic/__pycache__/adapter.cpython-314.pyc b/plugin/hermes_meshtastic/__pycache__/adapter.cpython-314.pyc deleted file mode 100644 index baeb1e1..0000000 Binary files a/plugin/hermes_meshtastic/__pycache__/adapter.cpython-314.pyc and /dev/null differ diff --git a/plugin/hermes_meshtastic/adapter.py b/plugin/hermes_meshtastic/adapter.py index 5f8fefa..362276d 100644 --- a/plugin/hermes_meshtastic/adapter.py +++ b/plugin/hermes_meshtastic/adapter.py @@ -3,306 +3,345 @@ Connects to a Meshtastic node via TCP (e.g. meshtastic.local:4403), listening for LoRa packets, stripping formatting, and dispatching to Hermes agent sessions. Outbound replies are chunked into concise LoRa-friendly packets. + +Contract note +------------- +The Hermes gateway API (``gateway.platforms.base.BasePlatformAdapter`` and friends) +is resolved LAZILY, inside ``get_adapter_class()``, never at module import time. +Hermes' own guidance (``hermes_cli/plugins.py``) puts gateway imports inside factory +bodies because third-party plugin modules are imported during plugin discovery, which +runs before/independently of the gateway runtime being fully importable. Importing +``gateway.platforms.base`` at module scope silently failed in the deployed runtime +(an ``ImportError`` swallowed by an earlier fallback bound the adapter to a local stub +class lacking ``set_message_handler`` — see git history). Deferred resolution keeps +this module importable in isolation (unit tests / CI without Hermes) while always +binding the adapter to the REAL gateway base at registration/creation time. + +The contract tests in ``tests/contract/`` assert this invariant against the real +``hermes-agent`` package (pinned to the deployed image tag). """ from __future__ import annotations import asyncio import logging -import re import time from typing import Any, Dict, List, Optional -try: - from gateway.platforms.base import ( - BasePlatformAdapter, - MessageEvent, - MessageType, - SendResult, - ) -except ImportError: - # Allow importing and testing utility functions without gateway in environment - class BasePlatformAdapter: # type: ignore[no-redef] - def __init__(self, config: Any, platform: Any = "meshtastic"): - self.config = config - self.platform = platform - - def build_source(self, **kwargs): - return kwargs - - async def handle_message(self, event: Any) -> None: - pass - - class SendResult: # type: ignore[no-redef] - def __init__(self, success: bool, message_id: Optional[str] = None, error: Optional[str] = None): - self.success = success - self.message_id = message_id - self.error = error - - class MessageType: # type: ignore[no-redef] - TEXT = "text" - - class MessageEvent: # type: ignore[no-redef] - def __init__(self, **kwargs): - for k, v in kwargs.items(): - setattr(self, k, v) +from .formatting import MAX_LORA_MESSAGE_LENGTH, chunk_text, strip_markdown logger = logging.getLogger(__name__) -MAX_LORA_MESSAGE_LENGTH = 180 # Characters per packet to prevent LoRa airtime congestion - -_MARKDOWN_RULES = ( - (r"\*\*(.+?)\*\*", r"\1"), # bold - (r"__(.+?)__", r"\1"), - (r"\*(.+?)\*", r"\1"), # italic - (r"(? str: - """Convert markdown formatting into plain readable text suitable for radio.""" - if not text: - return "" - result = text - for pattern, replacement in _MARKDOWN_RULES: - result = re.sub(pattern, replacement, result, flags=re.MULTILINE) - # Collapse multiple blank lines - result = re.sub(r"\n{3,}", "\n\n", result) - return result.strip() +# --------------------------------------------------------------------------- +# Hermes gateway API — resolved lazily and cached (see module docstring). +# --------------------------------------------------------------------------- + +_gateway_api: Optional[Dict[str, Any]] = None -def chunk_text(text: str, limit: int = MAX_LORA_MESSAGE_LENGTH) -> List[str]: - """Split text into smaller chunks for LoRa transmission, breaking at whitespace/newlines.""" - if not text: - return [] - if len(text) <= limit: - return [text] +def _load_gateway_api() -> Dict[str, Any]: + """Import the real Hermes gateway types and bind them as module globals. - chunks: List[str] = [] - remaining = text - while remaining: - if len(remaining) <= limit: - chunks.append(remaining.strip()) - break - split_at = remaining.rfind(" ", 0, limit) - if split_at <= 0 or split_at < (limit // 3): - newline_at = remaining.rfind("\n", 0, limit) - if newline_at > 0: - split_at = newline_at - else: - split_at = limit - chunk = remaining[:split_at].strip() - if chunk: - chunks.append(chunk) - remaining = remaining[split_at:].lstrip() - return chunks + Called from ``get_adapter_class()`` at adapter-creation time (the gateway + runtime is fully loaded by then). Method bodies reference these names as + module globals, so they are injected into this module's namespace. + """ + global _gateway_api + if _gateway_api is None: + from gateway.config import Platform + from gateway.platforms.base import ( # type: ignore[import-not-found] + BasePlatformAdapter, + MessageEvent, + MessageType, + SendResult, + ) + + _gateway_api = { + "Platform": Platform, + "BasePlatformAdapter": BasePlatformAdapter, + "MessageEvent": MessageEvent, + "MessageType": MessageType, + "SendResult": SendResult, + } + globals().update(_gateway_api) + return _gateway_api + + +# --------------------------------------------------------------------------- +# Dependency probe (passive — never installs anything). +# --------------------------------------------------------------------------- def check_requirements() -> bool: """Check if meshtastic is installed.""" try: - import meshtastic - import meshtastic.tcp_interface + import meshtastic # noqa: F401 + import meshtastic.tcp_interface # noqa: F401 return True except ImportError: return False -class MeshtasticAdapter(BasePlatformAdapter): - """Hermes gateway adapter for Meshtastic LoRa radios over TCP.""" +# --------------------------------------------------------------------------- +# Adapter — class body defined once the real gateway API is available. +# --------------------------------------------------------------------------- - def __init__(self, config: Any): - super().__init__(config, platform="meshtastic") - self.extra = getattr(config, "extra", {}) or {} - self.host = self.extra.get("host", "meshtastic.local") - self.port = int(self.extra.get("port", 4403)) - self.channel_index = int(self.extra.get("channel_index", 1)) - self.allowed_nodes: List[str] = [ - n.strip() for n in self.extra.get("allowed_nodes", []) if n.strip() - ] - self._iface = None - self._loop: Optional[asyncio.AbstractEventLoop] = None - self._connected = False - self._running = False - self._task: Optional[asyncio.Task] = None +_adapter_class = None - async def connect(self) -> None: - """Connect to Meshtastic radio via TCP and start subscriber.""" - self._loop = asyncio.get_running_loop() - self._running = True - self._task = asyncio.create_task(self._connection_supervisor()) - async def _connection_supervisor(self) -> None: - backoff = 2.0 - while self._running: - try: - import meshtastic.tcp_interface - from pubsub import pub +def get_adapter_class(): + """Return the ``MeshtasticAdapter`` class bound to the real gateway base. - logger.info("[Meshtastic] Connecting to radio at %s:%s...", self.host, self.port) - # TCPInterface handles initial protobuf handshake synchronously - self._iface = await asyncio.to_thread( - meshtastic.tcp_interface.TCPInterface, - hostname=self.host, - portNumber=self.port, - noProto=False, - ) - self._connected = True - backoff = 2.0 - my_id = getattr(self._iface.myInfo, "my_node_num", "unknown") - logger.info("[Meshtastic] Connected to radio (my_node_num=%s, listening on ch=%d)", my_id, self.channel_index) + The class is defined (and cached) on first call so that ``BasePlatformAdapter`` + is the REAL Hermes class — a fallback stub at this point would reproduce the + ``set_message_handler`` AttributeError this design exists to prevent. + """ + global _adapter_class + if _adapter_class is None: + _load_gateway_api() - pub.subscribe(self._on_meshtastic_receive, "meshtastic.receive.text") + class MeshtasticAdapter(BasePlatformAdapter): + """Hermes gateway adapter for Meshtastic LoRa radios over TCP.""" - # Keep-alive loop - while self._running and self._connected: - await asyncio.sleep(5) - if not self._iface or not self._iface.isConnected: - break - - except asyncio.CancelledError: - break - except Exception as exc: - logger.warning("[Meshtastic] Radio connection failed (%s). Retrying in %.1fs...", exc, backoff) + def __init__(self, config: Any): + super().__init__(config=config, platform=Platform("meshtastic")) + self.extra = getattr(config, "extra", {}) or {} + self.host = self.extra.get("host", "meshtastic.local") + self.port = int(self.extra.get("port", 4403)) + self.channel_index = int(self.extra.get("channel_index", 1)) + self.allowed_nodes: List[str] = [ + n.strip() for n in self.extra.get("allowed_nodes", []) if n.strip() + ] + self._iface = None + self._loop: Optional[asyncio.AbstractEventLoop] = None self._connected = False - await asyncio.sleep(backoff) - backoff = min(backoff * 1.5, 60.0) + self._running = False + self._task: Optional[asyncio.Task] = None + self._connecting = False - def _on_meshtastic_receive(self, packet: Dict[str, Any], interface: Any) -> None: - """Handle incoming text packet from Meshtastic pubsub (runs in worker thread).""" - try: - ch = packet.get("channel", 0) - if self.channel_index is not None and ch != self.channel_index: - # Ignore packets from different channels (e.g. public channel 0 when listening on 1) - return + async def connect(self, *, is_reconnect: bool = False) -> bool: + """Establish the Meshtastic TCP session; return success. - from_id = packet.get("fromId") or f"!{packet.get('from', 0):08x}" - my_info = getattr(interface, "myInfo", None) - my_node_num = getattr(my_info, "my_node_num", 0) - if packet.get("from") == my_node_num: - # Prevent self-message loops - return + Mirrors the gateway contract ``async connect(*, is_reconnect=False) + -> bool`` (see tests/contract/). On initial failure the watchdog + keeps retrying with exponential backoff in the background. + """ + self._loop = asyncio.get_running_loop() + self._running = True + try: + ok = await self._try_connect() + except Exception as exc: + logger.warning( + "[Meshtastic] Initial connect failed (%s); watchdog will retry.", exc + ) + ok = False + self._connected = ok + self._ensure_watchdog() + return ok - if self.allowed_nodes and from_id not in self.allowed_nodes: - logger.debug("[Meshtastic] Ignoring message from unauthorized node: %s", from_id) - return + async def _try_connect(self) -> bool: + if self._connecting: + return False # single-flight: gateway + watchdog must not double-connect + self._connecting = True + try: + import meshtastic.tcp_interface + from pubsub import pub - decoded = packet.get("decoded", {}) - text = decoded.get("text", "").strip() - if not text: - return + logger.info( + "[Meshtastic] Connecting to radio at %s:%s (is_reconnect=%s)...", + self.host, self.port, getattr(self, "_last_reconnect", False), + ) + # TCPInterface performs the protobuf handshake synchronously. + self._iface = await asyncio.to_thread( + meshtastic.tcp_interface.TCPInterface, + hostname=self.host, + portNumber=self.port, + noProto=False, + ) + my_id = getattr(self._iface.myInfo, "my_node_num", "unknown") + logger.info( + "[Meshtastic] Connected to radio (my_node_num=%s, listening on ch=%d)", + my_id, self.channel_index, + ) + pub.subscribe(self._on_meshtastic_receive, "meshtastic.receive.text") + return True + finally: + self._connecting = False - # Determine sender and target - to_id = packet.get("toId") - is_dm = to_id != "^all" and packet.get("to") == my_node_num - chat_id = from_id if is_dm else f"channel_{ch}" - chat_type = "dm" if is_dm else "channel" + def _ensure_watchdog(self) -> None: + if self._task is None or self._task.done(): + self._task = asyncio.create_task(self._connection_watchdog()) - rx_snr = packet.get("rxSnr", 0.0) - logger.info( - "[Meshtastic] Inbound message from %s on ch=%s (is_dm=%s, SNR=%.1fdB): %s", - from_id, ch, is_dm, rx_snr, text - ) + async def _connection_watchdog(self) -> None: + """Keep the radio session alive; reconnect with backoff on drop.""" + backoff = 2.0 + while self._running: + await asyncio.sleep(5) + try: + iface = self._iface + if iface is not None and not iface.isConnected: + self._connected = False + except Exception: + self._connected = False + if self._connected: + backoff = 2.0 + continue + # Radio went away (or never came up) — reconnect. + self._last_reconnect = True + try: + ok = await self._try_connect() + self._connected = ok + backoff = 2.0 + except Exception as exc: + logger.warning( + "[Meshtastic] Radio reconnect failed (%s). Retrying in %.1fs...", + exc, backoff, + ) + self._connected = False + await asyncio.sleep(backoff) + backoff = min(backoff * 1.5, 60.0) - # Build Hermes SessionSource and MessageEvent - source = self.build_source( - chat_id=chat_id, - chat_name=f"Meshtastic {chat_id}", - chat_type=chat_type, - user_id=from_id, - user_name=from_id, - message_id=str(packet.get("id", int(time.time()))), - ) + def _on_meshtastic_receive(self, packet: Dict[str, Any], interface: Any) -> None: + """Handle incoming text packet from Meshtastic pubsub (worker thread).""" + try: + ch = packet.get("channel", 0) + if self.channel_index is not None and ch != self.channel_index: + # Ignore packets from different channels (e.g. public channel 0). + return - event = MessageEvent( - text=text, - message_type=MessageType.TEXT, - user_id=from_id, - user_name=from_id, - source=source, - raw_message=packet, - message_id=str(packet.get("id", int(time.time()))), - ) + from_id = packet.get("fromId") or f"!{packet.get('from', 0):08x}" + my_info = getattr(interface, "myInfo", None) + my_node_num = getattr(my_info, "my_node_num", 0) + if packet.get("from") == my_node_num: + # Prevent self-message loops. + return - if self._loop and self._loop.is_running(): - asyncio.run_coroutine_threadsafe(self.handle_message(event), self._loop) + if self.allowed_nodes and from_id not in self.allowed_nodes: + logger.debug( + "[Meshtastic] Ignoring message from unauthorized node: %s", from_id + ) + return - except Exception as exc: - logger.exception("[Meshtastic] Error processing received packet: %s", exc) + decoded = packet.get("decoded", {}) + text = decoded.get("text", "").strip() + if not text: + return - async def send( - self, - chat_id: str, - content: str, - reply_to: Optional[str] = None, - metadata: Optional[Dict[str, Any]] = None, - ) -> SendResult: - """Send a message out over LoRa (implements BasePlatformAdapter.send).""" - if not self._connected or not self._iface: - return SendResult(success=False, error="Meshtastic radio not connected") + to_id = packet.get("toId") + is_dm = to_id != "^all" and packet.get("to") == my_node_num + chat_id = from_id if is_dm else f"channel_{ch}" + chat_type = "dm" if is_dm else "channel" - clean_text = strip_markdown(content) - chunks = chunk_text(clean_text, MAX_LORA_MESSAGE_LENGTH) + rx_snr = packet.get("rxSnr", 0.0) + logger.info( + "[Meshtastic] Inbound message from %s on ch=%s (is_dm=%s, SNR=%.1fdB): %s", + from_id, ch, is_dm, rx_snr, text, + ) - # Destination: if chat_id starts with '!', it is a direct node ID - destination = chat_id if chat_id.startswith("!") else "^all" - channel_index = self.channel_index if destination == "^all" else 0 + source = self.build_source( + chat_id=chat_id, + chat_name=f"Meshtastic {chat_id}", + chat_type=chat_type, + user_id=from_id, + user_name=from_id, + message_id=str(packet.get("id", int(time.time()))), + ) - last_id = None - try: - for idx, chunk in enumerate(chunks): - if len(chunks) > 1: - formatted_chunk = f"({idx+1}/{len(chunks)}) {chunk}" - else: - formatted_chunk = chunk + event = MessageEvent( + text=text, + message_type=MessageType.TEXT, + user_id=from_id, + user_name=from_id, + source=source, + raw_message=packet, + message_id=str(packet.get("id", int(time.time()))), + ) - logger.info( - "[Meshtastic] Sending chunk to %s on ch=%s: %s", - destination, channel_index, formatted_chunk - ) + if self._loop and self._loop.is_running(): + asyncio.run_coroutine_threadsafe( + self.handle_message(event), self._loop + ) + except Exception as exc: + logger.exception("[Meshtastic] Error processing received packet: %s", exc) - # Send via mesh interface thread - await asyncio.to_thread( - self._iface.sendText, - text=formatted_chunk, - destinationId=destination, - channelIndex=channel_index, - ) - last_id = str(int(time.time() * 1000)) - # Inter-packet delay to comply with LoRa duty cycle - if len(chunks) > 1: - await asyncio.sleep(1.2) + async def send( + self, + chat_id: str, + content: str, + reply_to: Optional[str] = None, + metadata: Optional[Dict[str, Any]] = None, + ) -> SendResult: + """Send a message out over LoRa (implements BasePlatformAdapter.send).""" + if not self._connected or not self._iface: + return SendResult(success=False, error="Meshtastic radio not connected") - return SendResult(success=True, message_id=last_id or str(int(time.time()))) - except Exception as exc: - logger.exception("[Meshtastic] Failed to send message: %s", exc) - return SendResult(success=False, error=str(exc)) + clean_text = strip_markdown(content) + chunks = chunk_text(clean_text, MAX_LORA_MESSAGE_LENGTH) - async def get_chat_info(self, chat_id: str) -> Dict[str, Any]: - """Get chat/channel info (implements BasePlatformAdapter.get_chat_info).""" - is_dm = chat_id.startswith("!") - return { - "name": f"Node {chat_id}" if is_dm else f"Channel {chat_id}", - "type": "dm" if is_dm else "channel", - } + # Destination: if chat_id starts with '!', it is a direct node ID. + destination = chat_id if chat_id.startswith("!") else "^all" + channel_index = self.channel_index if destination == "^all" else 0 - async def disconnect(self) -> None: - """Disconnect and clean up resources.""" - self._running = False - if self._task: - self._task.cancel() - if self._iface: - try: - from pubsub import pub - pub.unsubscribe(self._on_meshtastic_receive, "meshtastic.receive.text") - except Exception: - pass - await asyncio.to_thread(self._iface.close) - self._connected = False - logger.info("[Meshtastic] Radio interface closed") + last_id = None + try: + for idx, chunk in enumerate(chunks): + if len(chunks) > 1: + formatted_chunk = f"({idx+1}/{len(chunks)}) {chunk}" + else: + formatted_chunk = chunk + + logger.info( + "[Meshtastic] Sending chunk to %s on ch=%s: %s", + destination, channel_index, formatted_chunk, + ) + await asyncio.to_thread( + self._iface.sendText, + text=formatted_chunk, + destinationId=destination, + channelIndex=channel_index, + ) + last_id = str(int(time.time() * 1000)) + # Inter-packet delay to comply with LoRa duty cycle. + if len(chunks) > 1: + await asyncio.sleep(1.2) + return SendResult(success=True, message_id=last_id or str(int(time.time()))) + except Exception as exc: + logger.exception("[Meshtastic] Failed to send message: %s", exc) + return SendResult(success=False, error=str(exc)) + + async def get_chat_info(self, chat_id: str) -> Dict[str, Any]: + """Get chat/channel info (implements BasePlatformAdapter.get_chat_info).""" + is_dm = chat_id.startswith("!") + return { + "name": f"Node {chat_id}" if is_dm else f"Channel {chat_id}", + "type": "dm" if is_dm else "channel", + } + + async def disconnect(self) -> None: + """Disconnect and clean up resources.""" + self._running = False + if self._task: + self._task.cancel() + if self._iface: + try: + from pubsub import pub + pub.unsubscribe( + self._on_meshtastic_receive, "meshtastic.receive.text" + ) + except Exception: + pass + await asyncio.to_thread(self._iface.close) + self._connected = False + logger.info("[Meshtastic] Radio interface closed") + + _adapter_class = MeshtasticAdapter + return _adapter_class diff --git a/plugin/hermes_meshtastic/formatting.py b/plugin/hermes_meshtastic/formatting.py new file mode 100644 index 0000000..9ae5330 --- /dev/null +++ b/plugin/hermes_meshtastic/formatting.py @@ -0,0 +1,61 @@ +"""Pure formatting utilities for LoRa-friendly text. + +These helpers have NO dependency on the Hermes gateway package so they can be +imported and unit-tested in isolation (CI runs them without hermes-agent). +""" + +import re + +MAX_LORA_MESSAGE_LENGTH = 180 # Characters per packet to prevent LoRa airtime congestion + +_MARKDOWN_RULES = ( + (r"\*\*(.+?)\*\*", r"\1"), # bold + (r"__(.+?)__", r"\1"), + (r"\*(.+?)\*", r"\1"), # italic + (r"(? str: + """Convert markdown formatting into plain readable text suitable for radio.""" + if not text: + return "" + result = text + for pattern, replacement in _MARKDOWN_RULES: + result = re.sub(pattern, replacement, result, flags=re.MULTILINE) + # Collapse multiple blank lines + result = re.sub(r"\n{3,}", "\n\n", result) + return result.strip() + + +def chunk_text(text: str, limit: int = MAX_LORA_MESSAGE_LENGTH) -> list: + """Split text into smaller chunks for LoRa transmission, breaking at whitespace/newlines.""" + if not text: + return [] + if len(text) <= limit: + return [text] + + chunks = [] + remaining = text + while remaining: + if len(remaining) <= limit: + chunks.append(remaining.strip()) + break + split_at = remaining.rfind(" ", 0, limit) + if split_at <= 0 or split_at < (limit // 3): + newline_at = remaining.rfind("\n", 0, limit) + if newline_at > 0: + split_at = newline_at + else: + split_at = limit + chunk = remaining[:split_at].strip() + if chunk: + chunks.append(chunk) + remaining = remaining[split_at:].lstrip() + return chunks diff --git a/plugin/pyproject.toml b/plugin/pyproject.toml index e5f95f4..b641911 100644 --- a/plugin/pyproject.toml +++ b/plugin/pyproject.toml @@ -23,3 +23,10 @@ meshtastic-platform = "hermes_meshtastic" [tool.hatch.build.targets.wheel] packages = ["hermes_meshtastic"] + +[tool.hatch.build.targets.sdist] +exclude = [ + ".venv-contract", + ".hermes-src", + "tests", +] diff --git a/plugin/tests/__pycache__/test_adapter.cpython-314-pytest-9.0.3.pyc b/plugin/tests/__pycache__/test_adapter.cpython-314-pytest-9.0.3.pyc deleted file mode 100644 index ef9af66..0000000 Binary files a/plugin/tests/__pycache__/test_adapter.cpython-314-pytest-9.0.3.pyc and /dev/null differ diff --git a/plugin/tests/contract/test_gateway_contract.py b/plugin/tests/contract/test_gateway_contract.py new file mode 100644 index 0000000..af43967 --- /dev/null +++ b/plugin/tests/contract/test_gateway_contract.py @@ -0,0 +1,230 @@ +"""Contract tests: hermes-meshtastic vs the REAL Hermes gateway API. + +These tests exist so adapter development converges WITHOUT deploying: they +replay the gateway's adapter bootstrap against the actual ``hermes-agent`` +package (pinned to the deployed image tag, see ``scripts/contract-test.sh``) +and fail loudly on any API drift that used to surface only as a gateway +crash after release. + +Runtime ordering modelled here (verified against ``gateway/run.py`` and +``gateway/config.py``): + 1. plugin discovery runs the plugin's ``register(ctx)`` (registers the + platform in ``gateway.platform_registry``), + 2. the gateway later instantiates the adapter via the registry entry and + wires it: ``set_message_handler`` -> ... -> ``connect()``, + 3. ``Platform("meshtastic")`` (auto-created enum member) is only valid + AFTER registration — the enum's ``_missing_`` queries the registry. + +Layers covered +-------------- +L0 The adapter subclasses the REAL ``gateway.platforms.base.BasePlatformAdapter`` + (regression for the fallback-stub bug that caused the + ``AttributeError: ... 'set_message_handler'`` crash at ``gateway/run.py``). +L0b ABC instantiation gate — construction with a real ``PlatformConfig`` fails + if any abstract member is missing or the platform enum is misused. +L1 Wiring replay — every ``set_*`` call ``gateway/run.py`` makes at startup + must exist and accept a handler. +L2 Signature conformance — ``connect``/``send`` match the gateway contract. +L3 Event shape — the inbound ``MessageEvent``/``SessionSource`` construction + must be valid against the real gateway types. + +Skip policy: skipped when ``hermes-agent`` is not installed so the standalone +unit suite (CI without Hermes) stays green. Set ``HERMES_CONTRACT_REQUIRED=1`` +to turn absence into a hard failure — the dedicated contract CI job sets this +so the gate can never silently skip. +""" + +import inspect +import os + +import pytest + +try: + from gateway.config import Platform, PlatformConfig + from gateway.platforms.base import BasePlatformAdapter, MessageEvent, MessageType +except ImportError: # pragma: no cover - exercised when hermes-agent is absent + if os.environ.get("HERMES_CONTRACT_REQUIRED"): + raise AssertionError( + "hermes-agent is required for contract tests (HERMES_CONTRACT_REQUIRED=1). " + "Run `scripts/contract-test.sh` or the CI contract job." + ) + pytest.skip( + "hermes-agent not installed — contract tests skipped (scripts/contract-test.sh)", + allow_module_level=True, + ) + +import hermes_meshtastic # noqa: F401,E402 (package must import without gateway side effects) +from hermes_meshtastic import adapter as plugin_adapter # noqa: E402 + +MINIMAL_CONFIG = PlatformConfig( + enabled=True, + extra={"host": "127.0.0.1", "port": 1, "channel_index": 1}, +) + +# Exactly the set_* wiring gateway/run.py performs right before connect() +# (observed at the crash site: adapter.set_message_handler(...) etc.). +GATEWAY_WIRING_SETTERS = [ + "set_message_handler", + "set_fatal_error_handler", + "set_session_store", + "set_busy_session_handler", + "set_topic_recovery_fn", + "set_authorization_check", + "set_platform_event_handler", + "set_reaction_handler", +] + +_REGISTER_KWARGS = {} + + +class _Ctx: + """Minimal ctx mirroring how discovery calls the plugin's register().""" + + def register_platform(self, **kwargs): + _REGISTER_KWARGS.update(kwargs) + + +def _noop(*_args, **_kwargs): + return None + + +@pytest.fixture(scope="module", autouse=True) +def registered_entry(): + """Replay discovery: run register(), then register the PlatformEntry. + + Mirrors how gateway/run.py later instantiates the adapter from the + registry — and why Platform("meshtastic") is valid only after this. + """ + from gateway.platform_registry import PlatformEntry, platform_registry + + hermes_meshtastic.register(_Ctx()) + kw = dict(_REGISTER_KWARGS) + assert kw.get("name") == "meshtastic" + + if not platform_registry.is_registered("meshtastic"): + entry = PlatformEntry( + name=kw["name"], + label=kw["label"], + adapter_factory=kw["adapter_factory"], + check_fn=kw["check_fn"], + required_env=[], + max_message_length=kw.get("max_message_length", 0), + emoji=kw.get("emoji", "🔌"), + platform_hint=kw.get("platform_hint", ""), + ) + platform_registry.register(entry) + return platform_registry.get("meshtastic") + + +def _make_adapter(entry=None): + from gateway.platform_registry import platform_registry + + if entry is None: + entry = platform_registry.get("meshtastic") + assert entry is not None, "meshtastic platform not registered (fixture order broken)" + return entry.adapter_factory(MINIMAL_CONFIG) + + +# --- L0: real base binding ------------------------------------------------- + + +def test_adapter_subclasses_real_gateway_base(): + """Regression: adapter must bind the real base, never a local stub.""" + cls = plugin_adapter.get_adapter_class() + assert inspect.isclass(cls) + assert issubclass(cls, BasePlatformAdapter) + assert cls.__mro__[1].__module__ == "gateway.platforms.base" + + +def test_abstract_members_are_implemented(): + """ABC gate: BasePlatformAdapter requires connect/send/disconnect/get_chat_info.""" + cls = plugin_adapter.get_adapter_class() + for member in BasePlatformAdapter.__abstractmethods__: + impl = getattr(cls, member, None) + assert callable(impl), f"adapter missing abstract member: {member}" + + +# --- L0b: instantiation against the real contract -------------------------- + + +def test_adapter_instantiates_with_real_platform_config(registered_entry): + adapter = _make_adapter(registered_entry) + assert isinstance(adapter, BasePlatformAdapter) + # super().__init__ must receive the Platform enum (auto-created member), + # not a bare string — the gateway calls Platform methods on adapter.platform. + assert isinstance(adapter.platform, Platform) + assert adapter.platform.value == "meshtastic" + + +def test_registered_entry_carries_radio_limits(registered_entry): + assert registered_entry.max_message_length == plugin_adapter.MAX_LORA_MESSAGE_LENGTH + assert registered_entry.emoji == "📻" + + +# --- L1: gateway wiring replay --------------------------------------------- + + +def test_adapter_survives_gateway_wiring(registered_entry): + """Replay run.py's startup wiring; any missing setter crashes the gateway.""" + adapter = _make_adapter(registered_entry) + for setter in GATEWAY_WIRING_SETTERS: + fn = getattr(adapter, setter, None) + assert callable(fn), ( + f"adapter missing callable '{setter}' — gateway startup crashes here " + "(the original set_message_handler AttributeError)" + ) + fn(_noop) + # run.py also pokes these attributes directly after wiring. + adapter._busy_text_mode = False + + +# --- L2: signature conformance ---------------------------------------------- + + +def test_connect_signature_matches_gateway_contract(): + cls = plugin_adapter.get_adapter_class() + sig = inspect.signature(cls.connect) + assert inspect.iscoroutinefunction(cls.connect) + assert "is_reconnect" in sig.parameters, "connect() must accept is_reconnect kwarg" + param = sig.parameters["is_reconnect"] + assert param.kind == inspect.Parameter.KEYWORD_ONLY + assert param.default is False + + +def test_send_signature_matches_gateway_contract(): + base_sig = inspect.signature(BasePlatformAdapter.send) + cls = plugin_adapter.get_adapter_class() + sig = inspect.signature(cls.send) + base_params = [p for p in base_sig.parameters if p != "self"] + own_params = [p for p in sig.parameters if p != "self"] + assert own_params == base_params, ( + f"send() params {own_params} drift from gateway contract {base_params}" + ) + assert inspect.iscoroutinefunction(cls.send) + + +# --- L3: inbound event shape against real types ----------------------------- + + +def test_inbound_event_construction_is_valid(registered_entry): + """The exact source/event the adapter builds on receive must type-check.""" + adapter = _make_adapter(registered_entry) + source = adapter.build_source( + chat_id="!abcd1234", + chat_name="Meshtastic !abcd1234", + chat_type="dm", + user_id="!abcd1234", + user_name="!abcd1234", + message_id="42", + ) + event = MessageEvent( + text="hello over LoRa", + message_type=MessageType.TEXT, + user_id="!abcd1234", + user_name="!abcd1234", + source=source, + raw_message={"id": 42}, + message_id="42", + ) + assert event.source.platform.value == "meshtastic" + assert event.text == "hello over LoRa" diff --git a/scripts/contract-test.sh b/scripts/contract-test.sh new file mode 100755 index 0000000..d11d78d --- /dev/null +++ b/scripts/contract-test.sh @@ -0,0 +1,61 @@ +#!/usr/bin/env bash +# Gateway contract tests for hermes-meshtastic. +# +# Runs the plugin's contract suite (tests/contract/) against the REAL Hermes +# gateway API, pinned to the same upstream commit the deployed container image +# is built from (NousResearch/hermes-agent @ v2026.8.31 — the image tag used in +# k8s-flux apps/base/coder/hermes-agent). +# +# Why: the adapter used to be developed by guessing the gateway API and only +# failing at gateway startup after a release (AttributeError: ... no attribute +# 'set_message_handler'). These tests replay the gateway's adapter bootstrap +# (registry registration -> factory -> set_* wiring) so API drift fails here, +# in seconds, without deploying. +# +# Requires: uv (https://docs.astral.sh/uv/) and network on first run. +# +# Usage: scripts/contract-test.sh +set -euo pipefail + +REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +PLUGIN_DIR="$REPO_ROOT/plugin" +VENV="$PLUGIN_DIR/.venv-contract" +HERMES_SRC="$PLUGIN_DIR/.hermes-src" +HERMES_TAG="${HERMES_TAG:-v2026.8.31}" # must match k8s-flux hermes image tag +HERMES_COMMIT="${HERMES_COMMIT:-29112bef099274229cadff79cdff7bf7b99c4b77}" + +if ! command -v uv >/dev/null 2>&1; then + echo "error: uv is required (https://docs.astral.sh/uv/)" >&2 + exit 1 +fi + +if [ ! -d "$HERMES_SRC/.git" ]; then + echo "== cloning hermes-agent @ $HERMES_TAG (deployed image source) ==" + git clone --depth 1 --branch "$HERMES_TAG" \ + https://github.com/NousResearch/hermes-agent.git "$HERMES_SRC" +fi +actual="$(git -C "$HERMES_SRC" rev-parse HEAD)" +if [ "$actual" != "$HERMES_COMMIT" ]; then + echo "error: $HERMES_SRC is at $actual, expected $HERMES_COMMIT ($HERMES_TAG)." >&2 + echo " Refresh deliberately: HERMES_TAG= scripts/contract-test.sh" >&2 + exit 1 +fi + +if [ ! -x "$VENV/bin/python" ]; then + echo "== creating contract venv (python 3.13, matching the Hermes image) ==" + uv venv --python 3.13 "$VENV" +fi +uv pip install --python "$VENV/bin/python" -q -e "$HERMES_SRC" +uv pip install --python "$VENV/bin/python" -q pytest + +echo "== running contract tests (HERMES_CONTRACT_REQUIRED=1) ==" +( + cd "$PLUGIN_DIR" + HERMES_CONTRACT_REQUIRED=1 "$VENV/bin/python" -m pytest tests/contract -q "$@" +) + +echo "== running standalone unit tests (no gateway needed) ==" +( + cd "$PLUGIN_DIR" + "$VENV/bin/python" -m pytest tests/test_adapter.py -q +)