10 Commits

Author SHA1 Message Date
eric 2a7b71fc78 Merge pull request 'fix: accept PKI direct messages without channel' (#4) from codex/accept-pki-direct-messages into main
release / contract-test (push) Successful in 14s
release / test-and-release (push) Successful in 13s
2026-09-08 05:00:03 +00:00
eric facf1742ef fix: accept PKI direct messages without channel 2026-09-07 21:59:18 -07:00
eric d7aa18682b Merge pull request 'release v0.1.3: bump version, fix publish auth across http->https redirect' (#3) from v0.1.3-final into main
release / test-and-release (push) Successful in 12s
release / contract-test (push) Successful in 17s
2026-09-08 04:01:44 +00:00
eric b8e9b65995 fix(ci): forward auth token across http->https redirect on release publish
curl drops the Authorization header when -L follows Gitea's 308 redirect from
http server_url to https, so the API replied 'token is required' and publish
failed with no release id. Add --location-trusted to every publish curl.
2026-09-07 21:01:34 -07:00
eric 29b048e066 chore(plugin): release v0.1.3 2026-09-07 21:01:21 -07:00
eric 932171c34d Merge pull request 'fix(ci): follow Gitea http->https redirect when publishing releases' (#1) from ci-publish-fix into main
release / contract-test (push) Successful in 45s
release / test-and-release (push) Failing after 12s
2026-09-08 03:56:33 +00:00
eric cdc67bc263 fix(ci): follow Gitea http->https redirect when publishing releases
server_url resolves to http and Gitea answers the create-release POST with a
308 redirect; curl without -L returned an empty body and the grep-based id
extraction failed under set -e/pipefail, failing v0.1.3 publish. Add -L,
reuse an existing release on retry, and parse with python3.
2026-09-07 20:56:20 -07:00
eric 310c0244df fix(plugin): bind adapter to real Hermes gateway API; add contract-test gate
release / contract-test (push) Successful in 54s
release / test-and-release (push) Successful in 11s
- Resolve the Hermes gateway API lazily via get_adapter_class() so the adapter
  always inherits the real gateway.platforms.base.BasePlatformAdapter instead
  of the ImportError fallback stub (root cause of the AttributeError:
  'set_message_handler' crash at gateway startup).
- Pass Platform('meshtastic') enum to super().__init__ and align
  connect(*, is_reconnect=False) -> bool with the gateway contract.
- Split pure formatting utils into formatting.py; remove the stub fallback.
- Add tests/contract (L0-L3): real-base binding, ABC instantiation, gateway
  run.py wiring replay, signature and event-shape conformance, run against
  hermes-agent pinned to the deployed image source (v2026.8.31).
- Add scripts/contract-test.sh and a contract-test CI job; release publish
  now requires contract tests to pass.
- Untrack stray __pycache__ artifacts.
2026-09-07 20:50:21 -07:00
eric 2d9106b864 fix(plugin): implement BasePlatformAdapter abstract methods and proper build_source event
release / test-and-release (push) Failing after 12s
2026-09-07 19:45:37 -07:00
eric 7136f8c0fd fix: entry point should reference package module hermes_meshtastic not function
release / test-and-release (push) Successful in 13s
2026-09-07 19:32:15 -07:00
12 changed files with 786 additions and 241 deletions
+62 -5
View File
@@ -8,8 +8,49 @@ on:
- 'v*' - 'v*'
jobs: 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: test-and-release:
runs-on: ubuntu-latest runs-on: ubuntu-latest
needs: contract-test
steps: steps:
- name: Checkout code - name: Checkout code
uses: actions/checkout@v4 uses: actions/checkout@v4
@@ -40,21 +81,37 @@ jobs:
echo "Publishing release for $tag_name" echo "Publishing release for $tag_name"
wheel_file=$(ls plugin/dist/*.whl | head -n 1) wheel_file=$(ls plugin/dist/*.whl | head -n 1)
tar_file=$(ls plugin/dist/*.tar.gz | head -n 1) tar_file=$(ls plugin/dist/*.tar.gz | head -n 1)
api_url="${{ github.server_url }}/api/v1/repos/${{ github.repository }}"
# Create release # server_url is http and Gitea answers with a 308 redirect to https.
response=$(curl -s -k -X POST "${{ github.server_url }}/api/v1/repos/${{ github.repository }}/releases" \ # curl needs -L to follow, AND --location-trusted: without it curl
# drops the Authorization header on the scheme-changing redirect and
# Gitea replies "token is required". Release may already exist
# (retried run) — reuse it instead of failing.
fetch_id() { python3 -c "import json,sys;print(json.load(sys.stdin).get('id') or '')" 2>/dev/null; }
response=$(curl -s -k -L --location-trusted -H "Authorization: token $GITEA_TOKEN" \
"$api_url/releases/tags/$tag_name")
release_id=$(printf '%s' "$response" | fetch_id)
if [ -z "$release_id" ]; then
response=$(curl -s -k -L --location-trusted -X POST "$api_url/releases" \
-H "Authorization: token $GITEA_TOKEN" \ -H "Authorization: token $GITEA_TOKEN" \
-H "Content-Type: application/json" \ -H "Content-Type: application/json" \
-d "{\"tag_name\": \"$tag_name\", \"name\": \"$tag_name\", \"body\": \"Release $tag_name for hermes-meshtastic plugin\"}") -d "{\"tag_name\": \"$tag_name\", \"name\": \"$tag_name\", \"body\": \"Release $tag_name for hermes-meshtastic plugin\"}")
release_id=$(printf '%s' "$response" | fetch_id)
fi
release_id=$(echo "$response" | grep -o '"id":[0-9]*' | head -n 1 | cut -d: -f2)
if [ -n "$release_id" ]; then 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)" \ curl -s -k -L --location-trusted -X POST "$api_url/releases/$release_id/assets?name=$(basename "$wheel_file")" \
-H "Authorization: token $GITEA_TOKEN" \ -H "Authorization: token $GITEA_TOKEN" \
-H "Content-Type: application/octet-stream" \ -H "Content-Type: application/octet-stream" \
--data-binary "@$wheel_file" --data-binary "@$wheel_file"
curl -s -k -X POST "${{ github.server_url }}/api/v1/repos/${{ github.repository }}/releases/$release_id/assets?name=$(basename $tar_file)" \ curl -s -k -L --location-trusted -X POST "$api_url/releases/$release_id/assets?name=$(basename "$tar_file")" \
-H "Authorization: token $GITEA_TOKEN" \ -H "Authorization: token $GITEA_TOKEN" \
-H "Content-Type: application/octet-stream" \ -H "Content-Type: application/octet-stream" \
--data-binary "@$tar_file" --data-binary "@$tar_file"
else
echo "ERROR: could not create or fetch release $tag_name" >&2
exit 1
fi fi
+4
View File
@@ -4,3 +4,7 @@
.pio/ .pio/
*.log *.log
plugin/dist/ plugin/dist/
plugin/.venv-contract/
plugin/.hermes-src/
__pycache__/
*.pyc
+6
View File
@@ -0,0 +1,6 @@
dist/
__pycache__/
*.pyc
.venv-contract/
.hermes-src/
.pytest_cache/
+11 -4
View File
@@ -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): def register(ctx):
"""Hermes plugin entrypoint.""" """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( ctx.register_platform(
name="meshtastic", name="meshtastic",
label="Meshtastic LoRa", label="Meshtastic LoRa",
adapter_factory=lambda config: MeshtasticAdapter(config), adapter_factory=lambda config: adapter_cls(config),
check_fn=check_requirements, check_fn=check_requirements,
max_message_length=MAX_LORA_MESSAGE_LENGTH, max_message_length=MAX_LORA_MESSAGE_LENGTH,
emoji="📻", emoji="📻",
+202 -128
View File
@@ -3,111 +3,117 @@
Connects to a Meshtastic node via TCP (e.g. meshtastic.local:4403), listening for 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 LoRa packets, stripping formatting, and dispatching to Hermes agent sessions. Outbound
replies are chunked into concise LoRa-friendly packets. 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 from __future__ import annotations
import asyncio import asyncio
import logging import logging
import re
import time import time
from typing import Any, Dict, List, Optional from typing import Any, Dict, List, Optional
try: from .formatting import MAX_LORA_MESSAGE_LENGTH, chunk_text, strip_markdown
from gateway.platforms.base import BasePlatformAdapter, SendResult
from gateway.platforms.event import MessageEvent, MessageType
except ImportError:
# Allow importing and testing utility functions without gateway in environment
class BasePlatformAdapter: # type: ignore[no-redef]
def __init__(self, config: Any):
self.config = config
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)
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
MAX_LORA_MESSAGE_LENGTH = 180 # Characters per packet to prevent LoRa airtime congestion __all__ = [
"MAX_LORA_MESSAGE_LENGTH",
"check_requirements",
"chunk_text",
"get_adapter_class",
"strip_markdown",
]
_MARKDOWN_RULES = (
(r"\*\*(.+?)\*\*", r"\1"), # bold # ---------------------------------------------------------------------------
(r"__(.+?)__", r"\1"), # Hermes gateway API — resolved lazily and cached (see module docstring).
(r"\*(.+?)\*", r"\1"), # italic # ---------------------------------------------------------------------------
(r"(?<!\w)_(.+?)_(?!\w)", r"\1"),
(r"`(.+?)`", r"\1"), # inline code _gateway_api: Optional[Dict[str, Any]] = None
(r"```[\w]*\n?", ""), # code fences
(r"!\[([^\]]*)\]\(([^)]+)\)", r"\1"),
(r"\[([^\]]+)\]\(([^)]+)\)", r"\1 (\2)"), def _load_gateway_api() -> Dict[str, Any]:
(r"^#+\s*", ""), # headings """Import the real Hermes gateway types and bind them as module globals.
(r"^\s*[-*+]\s+", "- "), # normalize lists
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 = {
def strip_markdown(text: str) -> str: "Platform": Platform,
"""Convert markdown formatting into plain readable text suitable for radio.""" "BasePlatformAdapter": BasePlatformAdapter,
if not text: "MessageEvent": MessageEvent,
return "" "MessageType": MessageType,
result = text "SendResult": SendResult,
for pattern, replacement in _MARKDOWN_RULES: }
result = re.sub(pattern, replacement, result, flags=re.MULTILINE) globals().update(_gateway_api)
# Collapse multiple blank lines return _gateway_api
result = re.sub(r"\n{3,}", "\n\n", result)
return result.strip()
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.""" # Dependency probe (passive — never installs anything).
if not text: # ---------------------------------------------------------------------------
return []
if len(text) <= limit:
return [text]
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
def check_requirements() -> bool: def check_requirements() -> bool:
"""Check if meshtastic is installed.""" """Check if meshtastic is installed."""
try: try:
import meshtastic import meshtastic # noqa: F401
import meshtastic.tcp_interface import meshtastic.tcp_interface # noqa: F401
return True return True
except ImportError: except ImportError:
return False return False
# ---------------------------------------------------------------------------
# Adapter — class body defined once the real gateway API is available.
# ---------------------------------------------------------------------------
_adapter_class = None
def get_adapter_class():
"""Return the ``MeshtasticAdapter`` class bound to the real gateway base.
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()
class MeshtasticAdapter(BasePlatformAdapter): class MeshtasticAdapter(BasePlatformAdapter):
"""Hermes gateway adapter for Meshtastic LoRa radios over TCP.""" """Hermes gateway adapter for Meshtastic LoRa radios over TCP."""
def __init__(self, config: Any): def __init__(self, config: Any):
super().__init__(config) super().__init__(config=config, platform=Platform("meshtastic"))
self.extra = getattr(config, "extra", {}) or {} self.extra = getattr(config, "extra", {}) or {}
self.host = self.extra.get("host", "meshtastic.local") self.host = self.extra.get("host", "meshtastic.local")
self.port = int(self.extra.get("port", 4403)) self.port = int(self.extra.get("port", 4403))
@@ -120,66 +126,117 @@ class MeshtasticAdapter(BasePlatformAdapter):
self._connected = False self._connected = False
self._running = False self._running = False
self._task: Optional[asyncio.Task] = None self._task: Optional[asyncio.Task] = None
self._connecting = False
async def connect(self) -> None: async def connect(self, *, is_reconnect: bool = False) -> bool:
"""Connect to Meshtastic radio via TCP and start subscriber.""" """Establish the Meshtastic TCP session; return success.
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._loop = asyncio.get_running_loop()
self._running = True self._running = True
self._task = asyncio.create_task(self._connection_supervisor()) 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
async def _connection_supervisor(self) -> None: async def _try_connect(self) -> bool:
backoff = 2.0 if self._connecting:
while self._running: return False # single-flight: gateway + watchdog must not double-connect
self._connecting = True
try: try:
import meshtastic.tcp_interface import meshtastic.tcp_interface
from pubsub import pub from pubsub import pub
logger.info("[Meshtastic] Connecting to radio at %s:%s...", self.host, self.port) logger.info(
# TCPInterface handles initial protobuf handshake synchronously "[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( self._iface = await asyncio.to_thread(
meshtastic.tcp_interface.TCPInterface, meshtastic.tcp_interface.TCPInterface,
hostname=self.host, hostname=self.host,
portNumber=self.port, portNumber=self.port,
noProto=False, noProto=False,
) )
self._connected = True
backoff = 2.0
my_id = getattr(self._iface.myInfo, "my_node_num", "unknown") 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) 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") pub.subscribe(self._on_meshtastic_receive, "meshtastic.receive.text")
return True
finally:
self._connecting = False
# Keep-alive loop def _ensure_watchdog(self) -> None:
while self._running and self._connected: if self._task is None or self._task.done():
self._task = asyncio.create_task(self._connection_watchdog())
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) await asyncio.sleep(5)
if not self._iface or not self._iface.isConnected: try:
break iface = self._iface
if iface is not None and not iface.isConnected:
except asyncio.CancelledError: self._connected = False
break 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: except Exception as exc:
logger.warning("[Meshtastic] Radio connection failed (%s). Retrying in %.1fs...", exc, backoff) logger.warning(
"[Meshtastic] Radio reconnect failed (%s). Retrying in %.1fs...",
exc, backoff,
)
self._connected = False self._connected = False
await asyncio.sleep(backoff) await asyncio.sleep(backoff)
backoff = min(backoff * 1.5, 60.0) backoff = min(backoff * 1.5, 60.0)
def _on_meshtastic_receive(self, packet: Dict[str, Any], interface: Any) -> None: def _on_meshtastic_receive(self, packet: Dict[str, Any], interface: Any) -> None:
"""Handle incoming text packet from Meshtastic pubsub (runs in worker thread).""" """Handle incoming text packet from Meshtastic pubsub (worker thread)."""
try: try:
ch = packet.get("channel", 0) ch = packet.get("channel", 0)
if self.channel_index is not None and ch != self.channel_index: my_info = getattr(interface, "myInfo", None)
# Ignore packets from different channels (e.g. public channel 0 when listening on 1) my_node_num = getattr(my_info, "my_node_num", 0)
is_direct_to_me = packet.get("to") == my_node_num
is_pki_dm = is_direct_to_me and packet.get("pkiEncrypted", False)
if (
self.channel_index is not None
and ch != self.channel_index
and not is_pki_dm
):
# Ignore packets from different channels (e.g. public channel 0).
# PKI direct messages omit the channel field and decode as
# channel 0 even when sent from a secondary channel.
return return
from_id = packet.get("fromId") or f"!{packet.get('from', 0):08x}" 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: if packet.get("from") == my_node_num:
# Prevent self-message loops # Prevent self-message loops.
return return
if self.allowed_nodes and from_id not in self.allowed_nodes: if self.allowed_nodes and from_id not in self.allowed_nodes:
logger.debug("[Meshtastic] Ignoring message from unauthorized node: %s", from_id) logger.debug(
"[Meshtastic] Ignoring message from unauthorized node: %s", from_id
)
return return
decoded = packet.get("decoded", {}) decoded = packet.get("decoded", {})
@@ -187,51 +244,58 @@ class MeshtasticAdapter(BasePlatformAdapter):
if not text: if not text:
return return
# Determine sender and target
to_id = packet.get("toId") to_id = packet.get("toId")
is_dm = to_id != "^all" and packet.get("to") == my_node_num is_dm = to_id != "^all" and is_direct_to_me
chat_id = from_id if is_dm else f"channel_{ch}" chat_id = from_id if is_dm else f"channel_{ch}"
chat_type = "dm" if is_dm else "channel"
rx_snr = packet.get("rxSnr", 0.0) rx_snr = packet.get("rxSnr", 0.0)
hop_limit = packet.get("hopLimit", 0)
logger.info( logger.info(
"[Meshtastic] Inbound message from %s on ch=%s (is_dm=%s, SNR=%.1fdB): %s", "[Meshtastic] Inbound message from %s on ch=%s (is_dm=%s, SNR=%.1fdB): %s",
from_id, ch, is_dm, rx_snr, text from_id, ch, is_dm, rx_snr, text,
) )
# Build Hermes MessageEvent source = self.build_source(
event = MessageEvent(
platform="meshtastic",
message_type=MessageType.TEXT,
text=text,
sender_id=from_id,
sender_name=from_id,
chat_id=chat_id, 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()))),
)
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()))), message_id=str(packet.get("id", int(time.time()))),
raw_event=packet,
) )
if self._loop and self._loop.is_running(): if self._loop and self._loop.is_running():
asyncio.run_coroutine_threadsafe(self.handle_message(event), self._loop) asyncio.run_coroutine_threadsafe(
self.handle_message(event), self._loop
)
except Exception as exc: except Exception as exc:
logger.exception("[Meshtastic] Error processing received packet: %s", exc) logger.exception("[Meshtastic] Error processing received packet: %s", exc)
async def send_message( async def send(
self, self,
chat_id: str, chat_id: str,
message: str, content: str,
reply_to_id: Optional[str] = None, reply_to: Optional[str] = None,
**kwargs: Any, metadata: Optional[Dict[str, Any]] = None,
) -> SendResult: ) -> SendResult:
"""Send a message out over LoRa.""" """Send a message out over LoRa (implements BasePlatformAdapter.send)."""
if not self._connected or not self._iface: if not self._connected or not self._iface:
return SendResult(success=False, error="Meshtastic radio not connected") return SendResult(success=False, error="Meshtastic radio not connected")
clean_text = strip_markdown(message) clean_text = strip_markdown(content)
chunks = chunk_text(clean_text, MAX_LORA_MESSAGE_LENGTH) chunks = chunk_text(clean_text, MAX_LORA_MESSAGE_LENGTH)
# Destination: if chat_id starts with '!', it is a direct node ID # Destination: if chat_id starts with '!', it is a direct node ID.
destination = chat_id if chat_id.startswith("!") else "^all" destination = chat_id if chat_id.startswith("!") else "^all"
channel_index = self.channel_index if destination == "^all" else 0 channel_index = self.channel_index if destination == "^all" else 0
@@ -245,10 +309,8 @@ class MeshtasticAdapter(BasePlatformAdapter):
logger.info( logger.info(
"[Meshtastic] Sending chunk to %s on ch=%s: %s", "[Meshtastic] Sending chunk to %s on ch=%s: %s",
destination, channel_index, formatted_chunk destination, channel_index, formatted_chunk,
) )
# Send via mesh interface thread
await asyncio.to_thread( await asyncio.to_thread(
self._iface.sendText, self._iface.sendText,
text=formatted_chunk, text=formatted_chunk,
@@ -256,15 +318,22 @@ class MeshtasticAdapter(BasePlatformAdapter):
channelIndex=channel_index, channelIndex=channel_index,
) )
last_id = str(int(time.time() * 1000)) last_id = str(int(time.time() * 1000))
# Inter-packet delay to comply with LoRa duty cycle # Inter-packet delay to comply with LoRa duty cycle.
if len(chunks) > 1: if len(chunks) > 1:
await asyncio.sleep(1.2) await asyncio.sleep(1.2)
return SendResult(success=True, message_id=last_id or str(int(time.time()))) return SendResult(success=True, message_id=last_id or str(int(time.time())))
except Exception as exc: except Exception as exc:
logger.exception("[Meshtastic] Failed to send message: %s", exc) logger.exception("[Meshtastic] Failed to send message: %s", exc)
return SendResult(success=False, error=str(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: async def disconnect(self) -> None:
"""Disconnect and clean up resources.""" """Disconnect and clean up resources."""
self._running = False self._running = False
@@ -273,9 +342,14 @@ class MeshtasticAdapter(BasePlatformAdapter):
if self._iface: if self._iface:
try: try:
from pubsub import pub from pubsub import pub
pub.unsubscribe(self._on_meshtastic_receive, "meshtastic.receive.text") pub.unsubscribe(
self._on_meshtastic_receive, "meshtastic.receive.text"
)
except Exception: except Exception:
pass pass
await asyncio.to_thread(self._iface.close) await asyncio.to_thread(self._iface.close)
self._connected = False self._connected = False
logger.info("[Meshtastic] Radio interface closed") logger.info("[Meshtastic] Radio interface closed")
_adapter_class = MeshtasticAdapter
return _adapter_class
+61
View File
@@ -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"(?<!\w)_(.+?)_(?!\w)", r"\1"),
(r"`(.+?)`", r"\1"), # inline code
(r"```[\w]*\n?", ""), # code fences
(r"!\[([^\]]*)\]\(([^)]+)\)", r"\1"),
(r"\[([^\]]+)\]\(([^)]+)\)", r"\1 (\2)"),
(r"^#+\s*", ""), # headings
(r"^\s*[-*+]\s+", "- "), # normalize lists
)
def strip_markdown(text: str) -> 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
+9 -2
View File
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
[project] [project]
name = "hermes-meshtastic" name = "hermes-meshtastic"
version = "0.1.0" version = "0.1.4"
description = "Meshtastic LoRa mesh gateway adapter plugin for Hermes Agent" description = "Meshtastic LoRa mesh gateway adapter plugin for Hermes Agent"
readme = "README.md" readme = "README.md"
requires-python = ">=3.10" requires-python = ">=3.10"
@@ -19,7 +19,14 @@ dependencies = [
] ]
[project.entry-points."hermes_agent.plugins"] [project.entry-points."hermes_agent.plugins"]
meshtastic-platform = "hermes_meshtastic:register" meshtastic-platform = "hermes_meshtastic"
[tool.hatch.build.targets.wheel] [tool.hatch.build.targets.wheel]
packages = ["hermes_meshtastic"] packages = ["hermes_meshtastic"]
[tool.hatch.build.targets.sdist]
exclude = [
".venv-contract",
".hermes-src",
"tests",
]
@@ -0,0 +1,268 @@
"""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
from types import SimpleNamespace
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"
def test_pki_direct_message_without_channel_reaches_gateway(
registered_entry, monkeypatch
):
"""PKI DMs omit channel and must bypass the secondary-channel filter."""
adapter = _make_adapter(registered_entry)
adapter._loop = SimpleNamespace(is_running=lambda: True)
captured = {}
async def handle_message(event):
captured["event"] = event
def submit(coro, loop):
captured["loop"] = loop
with pytest.raises(StopIteration):
coro.send(None)
adapter.handle_message = handle_message
monkeypatch.setattr(plugin_adapter.asyncio, "run_coroutine_threadsafe", submit)
interface = SimpleNamespace(myInfo=SimpleNamespace(my_node_num=0x1BA1BC60))
adapter._on_meshtastic_receive(
{
"from": 0x1BA1496C,
"to": 0x1BA1BC60,
"toId": "!1ba1bc60",
"pkiEncrypted": True,
"decoded": {"text": "hello over PKI"},
"id": 42,
},
interface,
)
assert captured["loop"] is adapter._loop
assert captured["event"].text == "hello over PKI"
assert captured["event"].source.chat_id == "!1ba1496c"
+61
View File
@@ -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=<new 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
)