Compare commits
4 Commits
v0.1.0
...
ci-publish-fix
| Author | SHA1 | Date | |
|---|---|---|---|
| cdc67bc263 | |||
| 310c0244df | |||
| 2d9106b864 | |||
| 7136f8c0fd |
@@ -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,35 @@ 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 may be http while Gitea redirects (308) to https:
|
||||||
response=$(curl -s -k -X POST "${{ github.server_url }}/api/v1/repos/${{ github.repository }}/releases" \
|
# curl needs -L or the POST returns an empty body. Release may also
|
||||||
-H "Authorization: token $GITEA_TOKEN" \
|
# already exist (retried run) — reuse it instead of failing.
|
||||||
-H "Content-Type: application/json" \
|
fetch_id() { python3 -c "import json,sys;print(json.load(sys.stdin).get('id') or '')" 2>/dev/null; }
|
||||||
-d "{\"tag_name\": \"$tag_name\", \"name\": \"$tag_name\", \"body\": \"Release $tag_name for hermes-meshtastic plugin\"}")
|
|
||||||
|
response=$(curl -s -k -L -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 -X POST "$api_url/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=$(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 -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 -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,3 +4,7 @@
|
|||||||
.pio/
|
.pio/
|
||||||
*.log
|
*.log
|
||||||
plugin/dist/
|
plugin/dist/
|
||||||
|
plugin/.venv-contract/
|
||||||
|
plugin/.hermes-src/
|
||||||
|
__pycache__/
|
||||||
|
*.pyc
|
||||||
|
|||||||
@@ -0,0 +1,6 @@
|
|||||||
|
dist/
|
||||||
|
__pycache__/
|
||||||
|
*.pyc
|
||||||
|
.venv-contract/
|
||||||
|
.hermes-src/
|
||||||
|
.pytest_cache/
|
||||||
@@ -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="📻",
|
||||||
|
|||||||
Binary file not shown.
Binary file not shown.
+291
-225
@@ -3,279 +3,345 @@
|
|||||||
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",
|
||||||
_MARKDOWN_RULES = (
|
"check_requirements",
|
||||||
(r"\*\*(.+?)\*\*", r"\1"), # bold
|
"chunk_text",
|
||||||
(r"__(.+?)__", r"\1"),
|
"get_adapter_class",
|
||||||
(r"\*(.+?)\*", r"\1"), # italic
|
"strip_markdown",
|
||||||
(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."""
|
# Hermes gateway API — resolved lazily and cached (see module docstring).
|
||||||
if not text:
|
# ---------------------------------------------------------------------------
|
||||||
return ""
|
|
||||||
result = text
|
_gateway_api: Optional[Dict[str, Any]] = None
|
||||||
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[str]:
|
def _load_gateway_api() -> Dict[str, Any]:
|
||||||
"""Split text into smaller chunks for LoRa transmission, breaking at whitespace/newlines."""
|
"""Import the real Hermes gateway types and bind them as module globals.
|
||||||
if not text:
|
|
||||||
return []
|
|
||||||
if len(text) <= limit:
|
|
||||||
return [text]
|
|
||||||
|
|
||||||
chunks: List[str] = []
|
Called from ``get_adapter_class()`` at adapter-creation time (the gateway
|
||||||
remaining = text
|
runtime is fully loaded by then). Method bodies reference these names as
|
||||||
while remaining:
|
module globals, so they are injected into this module's namespace.
|
||||||
if len(remaining) <= limit:
|
"""
|
||||||
chunks.append(remaining.strip())
|
global _gateway_api
|
||||||
break
|
if _gateway_api is None:
|
||||||
split_at = remaining.rfind(" ", 0, limit)
|
from gateway.config import Platform
|
||||||
if split_at <= 0 or split_at < (limit // 3):
|
from gateway.platforms.base import ( # type: ignore[import-not-found]
|
||||||
newline_at = remaining.rfind("\n", 0, limit)
|
BasePlatformAdapter,
|
||||||
if newline_at > 0:
|
MessageEvent,
|
||||||
split_at = newline_at
|
MessageType,
|
||||||
else:
|
SendResult,
|
||||||
split_at = limit
|
)
|
||||||
chunk = remaining[:split_at].strip()
|
|
||||||
if chunk:
|
_gateway_api = {
|
||||||
chunks.append(chunk)
|
"Platform": Platform,
|
||||||
remaining = remaining[split_at:].lstrip()
|
"BasePlatformAdapter": BasePlatformAdapter,
|
||||||
return chunks
|
"MessageEvent": MessageEvent,
|
||||||
|
"MessageType": MessageType,
|
||||||
|
"SendResult": SendResult,
|
||||||
|
}
|
||||||
|
globals().update(_gateway_api)
|
||||||
|
return _gateway_api
|
||||||
|
|
||||||
|
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
# Dependency probe (passive — never installs anything).
|
||||||
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
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
|
||||||
|
|
||||||
|
|
||||||
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):
|
_adapter_class = None
|
||||||
super().__init__(config)
|
|
||||||
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
|
|
||||||
|
|
||||||
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:
|
def get_adapter_class():
|
||||||
backoff = 2.0
|
"""Return the ``MeshtasticAdapter`` class bound to the real gateway base.
|
||||||
while self._running:
|
|
||||||
try:
|
|
||||||
import meshtastic.tcp_interface
|
|
||||||
from pubsub import pub
|
|
||||||
|
|
||||||
logger.info("[Meshtastic] Connecting to radio at %s:%s...", self.host, self.port)
|
The class is defined (and cached) on first call so that ``BasePlatformAdapter``
|
||||||
# TCPInterface handles initial protobuf handshake synchronously
|
is the REAL Hermes class — a fallback stub at this point would reproduce the
|
||||||
self._iface = await asyncio.to_thread(
|
``set_message_handler`` AttributeError this design exists to prevent.
|
||||||
meshtastic.tcp_interface.TCPInterface,
|
"""
|
||||||
hostname=self.host,
|
global _adapter_class
|
||||||
portNumber=self.port,
|
if _adapter_class is None:
|
||||||
noProto=False,
|
_load_gateway_api()
|
||||||
)
|
|
||||||
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)
|
|
||||||
|
|
||||||
pub.subscribe(self._on_meshtastic_receive, "meshtastic.receive.text")
|
class MeshtasticAdapter(BasePlatformAdapter):
|
||||||
|
"""Hermes gateway adapter for Meshtastic LoRa radios over TCP."""
|
||||||
|
|
||||||
# Keep-alive loop
|
def __init__(self, config: Any):
|
||||||
while self._running and self._connected:
|
super().__init__(config=config, platform=Platform("meshtastic"))
|
||||||
await asyncio.sleep(5)
|
self.extra = getattr(config, "extra", {}) or {}
|
||||||
if not self._iface or not self._iface.isConnected:
|
self.host = self.extra.get("host", "meshtastic.local")
|
||||||
break
|
self.port = int(self.extra.get("port", 4403))
|
||||||
|
self.channel_index = int(self.extra.get("channel_index", 1))
|
||||||
except asyncio.CancelledError:
|
self.allowed_nodes: List[str] = [
|
||||||
break
|
n.strip() for n in self.extra.get("allowed_nodes", []) if n.strip()
|
||||||
except Exception as exc:
|
]
|
||||||
logger.warning("[Meshtastic] Radio connection failed (%s). Retrying in %.1fs...", exc, backoff)
|
self._iface = None
|
||||||
|
self._loop: Optional[asyncio.AbstractEventLoop] = None
|
||||||
self._connected = False
|
self._connected = False
|
||||||
await asyncio.sleep(backoff)
|
self._running = False
|
||||||
backoff = min(backoff * 1.5, 60.0)
|
self._task: Optional[asyncio.Task] = None
|
||||||
|
self._connecting = False
|
||||||
|
|
||||||
def _on_meshtastic_receive(self, packet: Dict[str, Any], interface: Any) -> None:
|
async def connect(self, *, is_reconnect: bool = False) -> bool:
|
||||||
"""Handle incoming text packet from Meshtastic pubsub (runs in worker thread)."""
|
"""Establish the Meshtastic TCP session; return success.
|
||||||
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
|
|
||||||
|
|
||||||
from_id = packet.get("fromId") or f"!{packet.get('from', 0):08x}"
|
Mirrors the gateway contract ``async connect(*, is_reconnect=False)
|
||||||
my_info = getattr(interface, "myInfo", None)
|
-> bool`` (see tests/contract/). On initial failure the watchdog
|
||||||
my_node_num = getattr(my_info, "my_node_num", 0)
|
keeps retrying with exponential backoff in the background.
|
||||||
if packet.get("from") == my_node_num:
|
"""
|
||||||
# Prevent self-message loops
|
self._loop = asyncio.get_running_loop()
|
||||||
return
|
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:
|
async def _try_connect(self) -> bool:
|
||||||
logger.debug("[Meshtastic] Ignoring message from unauthorized node: %s", from_id)
|
if self._connecting:
|
||||||
return
|
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", {})
|
logger.info(
|
||||||
text = decoded.get("text", "").strip()
|
"[Meshtastic] Connecting to radio at %s:%s (is_reconnect=%s)...",
|
||||||
if not text:
|
self.host, self.port, getattr(self, "_last_reconnect", False),
|
||||||
return
|
)
|
||||||
|
# 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
|
def _ensure_watchdog(self) -> None:
|
||||||
to_id = packet.get("toId")
|
if self._task is None or self._task.done():
|
||||||
is_dm = to_id != "^all" and packet.get("to") == my_node_num
|
self._task = asyncio.create_task(self._connection_watchdog())
|
||||||
chat_id = from_id if is_dm else f"channel_{ch}"
|
|
||||||
|
|
||||||
rx_snr = packet.get("rxSnr", 0.0)
|
async def _connection_watchdog(self) -> None:
|
||||||
hop_limit = packet.get("hopLimit", 0)
|
"""Keep the radio session alive; reconnect with backoff on drop."""
|
||||||
logger.info(
|
backoff = 2.0
|
||||||
"[Meshtastic] Inbound message from %s on ch=%s (is_dm=%s, SNR=%.1fdB): %s",
|
while self._running:
|
||||||
from_id, ch, is_dm, rx_snr, text
|
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 MessageEvent
|
def _on_meshtastic_receive(self, packet: Dict[str, Any], interface: Any) -> None:
|
||||||
event = MessageEvent(
|
"""Handle incoming text packet from Meshtastic pubsub (worker thread)."""
|
||||||
platform="meshtastic",
|
try:
|
||||||
message_type=MessageType.TEXT,
|
ch = packet.get("channel", 0)
|
||||||
text=text,
|
if self.channel_index is not None and ch != self.channel_index:
|
||||||
sender_id=from_id,
|
# Ignore packets from different channels (e.g. public channel 0).
|
||||||
sender_name=from_id,
|
return
|
||||||
chat_id=chat_id,
|
|
||||||
message_id=str(packet.get("id", int(time.time()))),
|
|
||||||
raw_event=packet,
|
|
||||||
)
|
|
||||||
|
|
||||||
if self._loop and self._loop.is_running():
|
from_id = packet.get("fromId") or f"!{packet.get('from', 0):08x}"
|
||||||
asyncio.run_coroutine_threadsafe(self.handle_message(event), self._loop)
|
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
|
||||||
|
|
||||||
except Exception as exc:
|
if self.allowed_nodes and from_id not in self.allowed_nodes:
|
||||||
logger.exception("[Meshtastic] Error processing received packet: %s", exc)
|
logger.debug(
|
||||||
|
"[Meshtastic] Ignoring message from unauthorized node: %s", from_id
|
||||||
|
)
|
||||||
|
return
|
||||||
|
|
||||||
async def send_message(
|
decoded = packet.get("decoded", {})
|
||||||
self,
|
text = decoded.get("text", "").strip()
|
||||||
chat_id: str,
|
if not text:
|
||||||
message: str,
|
return
|
||||||
reply_to_id: Optional[str] = None,
|
|
||||||
**kwargs: Any,
|
|
||||||
) -> SendResult:
|
|
||||||
"""Send a message out over LoRa."""
|
|
||||||
if not self._connected or not self._iface:
|
|
||||||
return SendResult(success=False, error="Meshtastic radio not connected")
|
|
||||||
|
|
||||||
clean_text = strip_markdown(message)
|
to_id = packet.get("toId")
|
||||||
chunks = chunk_text(clean_text, MAX_LORA_MESSAGE_LENGTH)
|
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"
|
||||||
|
|
||||||
# Destination: if chat_id starts with '!', it is a direct node ID
|
rx_snr = packet.get("rxSnr", 0.0)
|
||||||
destination = chat_id if chat_id.startswith("!") else "^all"
|
logger.info(
|
||||||
channel_index = self.channel_index if destination == "^all" else 0
|
"[Meshtastic] Inbound message from %s on ch=%s (is_dm=%s, SNR=%.1fdB): %s",
|
||||||
|
from_id, ch, is_dm, rx_snr, text,
|
||||||
|
)
|
||||||
|
|
||||||
last_id = None
|
source = self.build_source(
|
||||||
try:
|
chat_id=chat_id,
|
||||||
for idx, chunk in enumerate(chunks):
|
chat_name=f"Meshtastic {chat_id}",
|
||||||
if len(chunks) > 1:
|
chat_type=chat_type,
|
||||||
formatted_chunk = f"({idx+1}/{len(chunks)}) {chunk}"
|
user_id=from_id,
|
||||||
else:
|
user_name=from_id,
|
||||||
formatted_chunk = chunk
|
message_id=str(packet.get("id", int(time.time()))),
|
||||||
|
)
|
||||||
|
|
||||||
logger.info(
|
event = MessageEvent(
|
||||||
"[Meshtastic] Sending chunk to %s on ch=%s: %s",
|
text=text,
|
||||||
destination, channel_index, formatted_chunk
|
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()))),
|
||||||
|
)
|
||||||
|
|
||||||
# Send via mesh interface thread
|
if self._loop and self._loop.is_running():
|
||||||
await asyncio.to_thread(
|
asyncio.run_coroutine_threadsafe(
|
||||||
self._iface.sendText,
|
self.handle_message(event), self._loop
|
||||||
text=formatted_chunk,
|
)
|
||||||
destinationId=destination,
|
except Exception as exc:
|
||||||
channelIndex=channel_index,
|
logger.exception("[Meshtastic] Error processing received packet: %s", exc)
|
||||||
)
|
|
||||||
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())))
|
async def send(
|
||||||
except Exception as exc:
|
self,
|
||||||
logger.exception("[Meshtastic] Failed to send message: %s", exc)
|
chat_id: str,
|
||||||
return SendResult(success=False, error=str(exc))
|
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")
|
||||||
|
|
||||||
async def disconnect(self) -> None:
|
clean_text = strip_markdown(content)
|
||||||
"""Disconnect and clean up resources."""
|
chunks = chunk_text(clean_text, MAX_LORA_MESSAGE_LENGTH)
|
||||||
self._running = False
|
|
||||||
if self._task:
|
# Destination: if chat_id starts with '!', it is a direct node ID.
|
||||||
self._task.cancel()
|
destination = chat_id if chat_id.startswith("!") else "^all"
|
||||||
if self._iface:
|
channel_index = self.channel_index if destination == "^all" else 0
|
||||||
try:
|
|
||||||
from pubsub import pub
|
last_id = None
|
||||||
pub.unsubscribe(self._on_meshtastic_receive, "meshtastic.receive.text")
|
try:
|
||||||
except Exception:
|
for idx, chunk in enumerate(chunks):
|
||||||
pass
|
if len(chunks) > 1:
|
||||||
await asyncio.to_thread(self._iface.close)
|
formatted_chunk = f"({idx+1}/{len(chunks)}) {chunk}"
|
||||||
self._connected = False
|
else:
|
||||||
logger.info("[Meshtastic] Radio interface closed")
|
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
|
||||||
|
|||||||
@@ -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
|
||||||
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "hermes-meshtastic"
|
name = "hermes-meshtastic"
|
||||||
version = "0.1.0"
|
version = "0.1.2"
|
||||||
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",
|
||||||
|
]
|
||||||
|
|||||||
Binary file not shown.
@@ -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"
|
||||||
Executable
+61
@@ -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
|
||||||
|
)
|
||||||
Reference in New Issue
Block a user