From 3605b9621fb319ef3b61e89f89747634a159efa6 Mon Sep 17 00:00:00 2001 From: fizzlepoof <251490314+fizzlepoof@users.noreply.github.com> Date: Fri, 4 Sep 2026 16:29:14 +0000 Subject: [PATCH 01/27] [verified] Harden NOMAD bridge runtime --- README.md | 32 ++- config.default.json | 7 +- meshcore_nomad_bridge/__init__.py | 2 +- meshcore_nomad_bridge/config.py | 103 ++++++++- meshcore_nomad_bridge/main.py | 127 ++++++++-- meshcore_nomad_bridge/nomad_client.py | 272 +++++++++++++++++++--- openhop-plugin.json | 9 +- pyproject.toml | 2 +- tests/test_config.py | 132 ++++++++++- tests/test_message_handler.py | 213 ++++++++++++++++- tests/test_nomad_client.py | 320 ++++++++++++++++++++++++++ tests/test_plugin_package.py | 2 +- 12 files changed, 1145 insertions(+), 76 deletions(-) diff --git a/README.md b/README.md index b38abe1..7c386f9 100644 --- a/README.md +++ b/README.md @@ -50,7 +50,9 @@ $OPENHOP_PLUGIN_DATA/ └── nomad_sessions.json ``` -`nomad_sessions.json` is only used for sender-scoped persistent N.O.M.A.D. conversations when `one_shot` is disabled. +Persistent conversations are disabled in v0.1.2 because the upstream session lifecycle cannot +yet be bounded safely. `one_shot` must remain `true`; the session-map setting is retained only +for configuration compatibility. When `OPENHOP_PLUGIN_DATA` is not set, existing standalone behaviour is preserved and the session map defaults to `./data/nomad_sessions.json`. @@ -63,7 +65,8 @@ Minimum configuration: ```json { "nomad_url": "http://192.168.0.170:8080", - "nomad_model": "qwen2.5:3b-instruct" + "nomad_model": "qwen2.5:3b-instruct", + "allowed_sender_prefixes": ["001122334455"] } ``` @@ -78,7 +81,12 @@ Typical configuration: "nomad_collection": null, "nomad_timeout_seconds": 120, "one_shot": true, - "max_concurrent_requests": 2, + "max_concurrent_requests": 1, + "max_pending_requests": 1, + "max_requests_per_sender": 2, + "max_requests_global": 4, + "rate_limit_window_seconds": 60, + "allowed_sender_prefixes": ["001122334455"], "busy_wait_seconds": 5, "max_reply_chunks": 4, "max_chunk_bytes": 145, @@ -97,6 +105,15 @@ built-in defaults < environment variables ``` +An empty `allowed_sender_prefixes` list denies every MeshCore sender. Set it to the exact +12-character sender prefixes that may use NOMAD. `one_shot` must remain `true`. The default +limits permit one active request, two requests per sender per minute, and four requests +globally per minute. +Rejected overload and authorization traffic is dropped without an RF reply. NOMAD HTTP +responses are capped at 256 KiB, redirects are not followed, and the configured timeout is +an end-to-end request deadline implemented without non-cancellable worker threads. `NOMAD_URL` +must use an IP literal so DNS resolution cannot outlive that deadline. + This keeps environment variables available for development and existing standalone deployments. Important environment overrides include: @@ -110,6 +127,11 @@ Important environment overrides include: - `ONE_SHOT` - `NOMAD_SESSION_MAP_PATH` - `MAX_CONCURRENT_REQUESTS` +- `MAX_PENDING_REQUESTS` +- `MAX_REQUESTS_PER_SENDER` +- `MAX_REQUESTS_GLOBAL` +- `RATE_LIMIT_WINDOW_SECONDS` +- `ALLOWED_SENDER_PREFIXES` (comma-separated 12-character hexadecimal sender prefixes) - `NOMAD_BUSY_WAIT_SECONDS` - `MAX_REPLY_CHUNKS` - `MAX_CHUNK_BYTES` @@ -130,7 +152,7 @@ Important environment overrides include: "schema": 1, "id": "openhop.nomad", "name": "NOMAD Bridge", - "version": "0.1.1", + "version": "0.1.2", "runtime": { "type": "python", "entrypoint": "meshcore-nomad-bridge" @@ -174,7 +196,7 @@ python -m build --wheel The wheel is written to `dist/`, for example: ```text -dist/openhop_nomad_plugin-0.1.1-py3-none-any.whl +dist/openhop_nomad_plugin-0.1.2-py3-none-any.whl ``` If you want both wheel and source distribution, run: diff --git a/config.default.json b/config.default.json index 358ac57..584c75d 100644 --- a/config.default.json +++ b/config.default.json @@ -6,7 +6,12 @@ "nomad_collection": null, "nomad_timeout_seconds": 120, "one_shot": true, - "max_concurrent_requests": 2, + "max_concurrent_requests": 1, + "max_pending_requests": 1, + "max_requests_per_sender": 2, + "max_requests_global": 4, + "rate_limit_window_seconds": 60, + "allowed_sender_prefixes": [], "busy_wait_seconds": 5, "max_reply_chunks": 4, "max_chunk_bytes": 145, diff --git a/meshcore_nomad_bridge/__init__.py b/meshcore_nomad_bridge/__init__.py index 89d4557..ed11b9d 100644 --- a/meshcore_nomad_bridge/__init__.py +++ b/meshcore_nomad_bridge/__init__.py @@ -2,4 +2,4 @@ __all__ = ["__version__"] -__version__ = "0.1.1" +__version__ = "0.1.2" diff --git a/meshcore_nomad_bridge/config.py b/meshcore_nomad_bridge/config.py index 511dd8b..e61a0ab 100644 --- a/meshcore_nomad_bridge/config.py +++ b/meshcore_nomad_bridge/config.py @@ -2,11 +2,15 @@ from __future__ import annotations -from dataclasses import dataclass +import ipaddress import json +import math import os +from dataclasses import dataclass from pathlib import Path +from string import Formatter from typing import Any +from urllib.parse import urlsplit class ConfigError(ValueError): @@ -31,6 +35,11 @@ class Settings: nomad_session_map_path: str max_concurrent_requests: int + max_pending_requests: int + max_requests_per_sender: int + max_requests_global: int + rate_limit_window_seconds: float + allowed_sender_prefixes: tuple[str, ...] busy_wait_seconds: float max_reply_chunks: int @@ -45,7 +54,7 @@ class Settings: log_level: str @classmethod - def from_env(cls) -> "Settings": + def from_env(cls) -> Settings: plugin_data_dir, config = _load_plugin_config() meshcore_host = _get_str("MESHCORE_HOST", "127.0.0.1", config) @@ -65,7 +74,12 @@ def from_env(cls) -> "Settings": "NOMAD_SESSION_MAP_PATH", session_map_default, config ) - max_concurrent_requests = _get_int("MAX_CONCURRENT_REQUESTS", 2, config) + max_concurrent_requests = _get_int("MAX_CONCURRENT_REQUESTS", 1, config) + max_pending_requests = _get_int("MAX_PENDING_REQUESTS", 1, config) + max_requests_per_sender = _get_int("MAX_REQUESTS_PER_SENDER", 2, config) + max_requests_global = _get_int("MAX_REQUESTS_GLOBAL", 4, config) + rate_limit_window_seconds = _get_float("RATE_LIMIT_WINDOW_SECONDS", 60.0, config) + allowed_sender_prefixes = _get_sender_prefixes("ALLOWED_SENDER_PREFIXES", config) busy_wait_seconds = _get_float("NOMAD_BUSY_WAIT_SECONDS", 5.0, config) max_reply_chunks = _get_int("MAX_REPLY_CHUNKS", 4, config) @@ -100,6 +114,11 @@ def from_env(cls) -> "Settings": one_shot=one_shot, nomad_session_map_path=nomad_session_map_path, max_concurrent_requests=max_concurrent_requests, + max_pending_requests=max_pending_requests, + max_requests_per_sender=max_requests_per_sender, + max_requests_global=max_requests_global, + rate_limit_window_seconds=rate_limit_window_seconds, + allowed_sender_prefixes=allowed_sender_prefixes, busy_wait_seconds=busy_wait_seconds, max_reply_chunks=max_reply_chunks, max_chunk_bytes=max_chunk_bytes, @@ -208,21 +227,80 @@ def _get_bool(name: str, default: bool, config: dict[str, Any]) -> bool: raise ConfigError(f"{name} must be a boolean") +def _get_sender_prefixes(name: str, config: dict[str, Any]) -> tuple[str, ...]: + raw = _raw_value(name, config, []) + if isinstance(raw, str): + values = raw.split(",") + elif isinstance(raw, list): + values = raw + else: + raise ConfigError(f"{name} must be a list or comma-separated string") + + prefixes: list[str] = [] + for value in values: + if not isinstance(value, str): + raise ConfigError(f"{name} entries must be strings") + prefix = value.strip().lower() + if not prefix: + continue + if len(prefix) != 12 or any(char not in "0123456789abcdef" for char in prefix): + raise ConfigError(f"{name} entries must be 12 hexadecimal characters") + prefixes.append(prefix) + return tuple(dict.fromkeys(prefixes)) + + def _validate(settings: Settings) -> None: if not settings.meshcore_host: raise ConfigError("MESHCORE_HOST must not be empty") if not (1 <= settings.meshcore_port <= 65535): raise ConfigError("MESHCORE_PORT must be between 1 and 65535") - if not settings.nomad_url.startswith(("http://", "https://")): - raise ConfigError("NOMAD_URL must start with http:// or https://") + try: + parsed_url = urlsplit(settings.nomad_url) + if parsed_url.scheme not in {"http", "https"} or not parsed_url.hostname: + raise ValueError + ipaddress.ip_address(parsed_url.hostname) + port = parsed_url.port + if ( + parsed_url.username + or parsed_url.password + or parsed_url.fragment + or parsed_url.query + or parsed_url.path not in {"", "/"} + ): + raise ValueError + if port is not None and not 1 <= port <= 65535: + raise ValueError + except ValueError as exc: + raise ConfigError("NOMAD_URL must be a valid HTTP(S) URL with an IP address") from exc if not settings.nomad_session_map_path: raise ConfigError("NOMAD_SESSION_MAP_PATH must not be empty") + numeric_settings = { + "NOMAD_TIMEOUT_SECONDS": settings.nomad_timeout_seconds, + "RATE_LIMIT_WINDOW_SECONDS": settings.rate_limit_window_seconds, + "NOMAD_BUSY_WAIT_SECONDS": settings.busy_wait_seconds, + } + for name, value in numeric_settings.items(): + if not math.isfinite(value): + raise ConfigError(f"{name} must be finite") + if settings.nomad_timeout_seconds <= 0: raise ConfigError("NOMAD_TIMEOUT_SECONDS must be > 0") if settings.max_concurrent_requests <= 0: raise ConfigError("MAX_CONCURRENT_REQUESTS must be > 0") + if settings.max_pending_requests < settings.max_concurrent_requests: + raise ConfigError("MAX_PENDING_REQUESTS must be >= MAX_CONCURRENT_REQUESTS") + if settings.max_requests_per_sender <= 0: + raise ConfigError("MAX_REQUESTS_PER_SENDER must be > 0") + if settings.max_requests_global <= 0: + raise ConfigError("MAX_REQUESTS_GLOBAL must be > 0") + if settings.max_requests_per_sender > settings.max_requests_global: + raise ConfigError("MAX_REQUESTS_PER_SENDER must be <= MAX_REQUESTS_GLOBAL") + if settings.rate_limit_window_seconds <= 0: + raise ConfigError("RATE_LIMIT_WINDOW_SECONDS must be > 0") + if not settings.one_shot: + raise ConfigError("ONE_SHOT must be true; persistent sessions are not bounded upstream") if settings.busy_wait_seconds < 0: raise ConfigError("NOMAD_BUSY_WAIT_SECONDS must be >= 0") @@ -233,8 +311,19 @@ def _validate(settings: Settings) -> None: if settings.max_prompt_bytes <= 0: raise ConfigError("MAX_PROMPT_BYTES must be > 0") - if "{question}" not in settings.radio_prompt_template: - raise ConfigError("RADIO_PROMPT_TEMPLATE must contain {question}") + try: + parsed_template = list(Formatter().parse(settings.radio_prompt_template)) + except ValueError as exc: + raise ConfigError("RADIO_PROMPT_TEMPLATE must be a valid format string") from exc + fields = [field_name for _, field_name, _, _ in parsed_template if field_name is not None] + invalid_fields = [ + field_name + for _, field_name, format_spec, conversion in parsed_template + if field_name is not None + and (format_spec or conversion or field_name != "question") + ] + if "question" not in fields or invalid_fields: + raise ConfigError("RADIO_PROMPT_TEMPLATE must contain only the {question} field") if settings.duplicate_ttl_seconds < 60: raise ConfigError("DUPLICATE_TTL_SECONDS must be >= 60") diff --git a/meshcore_nomad_bridge/main.py b/meshcore_nomad_bridge/main.py index 2466f49..16c60ec 100644 --- a/meshcore_nomad_bridge/main.py +++ b/meshcore_nomad_bridge/main.py @@ -7,6 +7,7 @@ import logging import signal import time +from collections import deque from typing import Any from .config import ConfigError, Settings @@ -14,7 +15,6 @@ from .nomad_client import NomadClient, NomadUnavailable from .text import clean_for_radio, split_for_meshcore - logger = logging.getLogger(__name__) NOMAD_UNAVAILABLE_MESSAGE = "NOMAD is temporarily unavailable. Please try again." @@ -29,19 +29,57 @@ def __init__(self, ttl_seconds: int) -> None: self._data: dict[str, float] = {} def seen(self, key: str) -> bool: + if self.contains(key): + return True + self.add(key) + return False + + def contains(self, key: str) -> bool: + now = time.monotonic() + self._prune(now) + return key in self._data + + def add(self, key: str) -> None: now = time.monotonic() self._prune(now) - if key in self._data: - return True self._data[key] = now + self._ttl_seconds - return False def _prune(self, now: float) -> None: - if len(self._data) < 64: - return self._data = {k: exp for k, exp in self._data.items() if exp > now} +class RequestRateLimiter: + def __init__(self, *, per_sender: int, global_limit: int, window_seconds: float) -> None: + self._per_sender = per_sender + self._global_limit = global_limit + self._window_seconds = window_seconds + self._global: deque[float] = deque() + self._senders: dict[str, deque[float]] = {} + + def allow(self, sender_id: str) -> bool: + now = time.monotonic() + cutoff = now - self._window_seconds + self._prune(self._global, cutoff) + for key, events in list(self._senders.items()): + self._prune(events, cutoff) + if not events: + del self._senders[key] + + sender_events = self._senders.setdefault(sender_id, deque()) + + if len(self._global) >= self._global_limit or len(sender_events) >= self._per_sender: + return False + + self._global.append(now) + sender_events.append(now) + return True + + @staticmethod + def _prune(events: deque[float], cutoff: float) -> None: + while events and events[0] <= cutoff: + events.popleft() + + class BridgeService: def __init__(self, settings: Settings, meshcore: MeshCoreClient, nomad: NomadClient) -> None: self._settings = settings @@ -51,6 +89,11 @@ def __init__(self, settings: Settings, meshcore: MeshCoreClient, nomad: NomadCli self._stop_event = asyncio.Event() self._semaphore = asyncio.Semaphore(settings.max_concurrent_requests) self._duplicate_cache = DuplicateCache(settings.duplicate_ttl_seconds) + self._rate_limiter = RequestRateLimiter( + per_sender=settings.max_requests_per_sender, + global_limit=settings.max_requests_global, + window_seconds=settings.rate_limit_window_seconds, + ) self._inflight: set[asyncio.Task[Any]] = set() @@ -60,11 +103,47 @@ async def run(self) -> None: await self._shutdown() async def _dispatch_message(self, message: IncomingMessage) -> None: - task = asyncio.create_task(self._handle_message(message), name="handle-message") + if message.txt_type != 0 or not message.text.strip(): + return + + sender_id = message.sender_prefix.hex() + if sender_id not in self._settings.allowed_sender_prefixes: + logger.debug("Dropping message from a sender that is not allowed") + return + + dedupe_key = self._dedupe_key(message) + if self._duplicate_cache.contains(dedupe_key): + logger.debug("Skipping duplicate message from %s", sender_id) + return + + if len(self._inflight) >= self._settings.max_pending_requests: + logger.debug("Dropping message because the pending request limit is full") + return + if not self._rate_limiter.allow(sender_id): + logger.debug("Dropping message because the request rate limit was reached") + return + self._duplicate_cache.add(dedupe_key) + task = asyncio.create_task( + self._handle_message(message, dedupe_checked=True), + name="handle-message", + ) self._inflight.add(task) - task.add_done_callback(self._inflight.discard) + task.add_done_callback(self._finish_task) - async def _handle_message(self, message: IncomingMessage) -> None: + def _finish_task(self, task: asyncio.Task[Any]) -> None: + self._inflight.discard(task) + if task.cancelled(): + return + error = task.exception() + if error is not None: + logger.error( + "Unhandled message task error", + exc_info=(type(error), error, error.__traceback__), + ) + + async def _handle_message( + self, message: IncomingMessage, *, dedupe_checked: bool = False + ) -> None: if message.txt_type != 0: return @@ -74,10 +153,11 @@ async def _handle_message(self, message: IncomingMessage) -> None: sender_id = message.sender_prefix.hex() - dedupe_key = self._dedupe_key(message) - if self._duplicate_cache.seen(dedupe_key): - logger.debug("Skipping duplicate message from %s", sender_id) - return + if not dedupe_checked: + dedupe_key = self._dedupe_key(message) + if self._duplicate_cache.seen(dedupe_key): + logger.debug("Skipping duplicate message from %s", sender_id) + return if (not self._settings.one_shot) and prompt.lower() in {"/new", "/reset"}: try: @@ -93,12 +173,21 @@ async def _handle_message(self, message: IncomingMessage) -> None: return acquired = False - try: - await asyncio.wait_for(self._semaphore.acquire(), timeout=self._settings.busy_wait_seconds) + if self._settings.busy_wait_seconds == 0: + if self._semaphore.locked(): + await self._meshcore.send_text(message.sender_prefix, NOMAD_BUSY_MESSAGE) + return + await self._semaphore.acquire() acquired = True - except asyncio.TimeoutError: - await self._meshcore.send_text(message.sender_prefix, NOMAD_BUSY_MESSAGE) - return + else: + try: + await asyncio.wait_for( + self._semaphore.acquire(), timeout=self._settings.busy_wait_seconds + ) + acquired = True + except asyncio.TimeoutError: + await self._meshcore.send_text(message.sender_prefix, NOMAD_BUSY_MESSAGE) + return try: logger.info("Message received from %s", message.sender_prefix.hex()) @@ -133,7 +222,7 @@ def _build_nomad_prompt(self, prompt: str) -> str: return self._settings.radio_prompt_template.format(question=prompt) def _dedupe_key(self, message: IncomingMessage) -> str: - digest = hashlib.sha1(message.text.strip().encode("utf-8")).hexdigest() + digest = hashlib.sha256(message.text.strip().encode("utf-8")).hexdigest() return f"{message.sender_prefix.hex()}:{message.timestamp}:{digest}" def _register_signals(self) -> None: diff --git a/meshcore_nomad_bridge/nomad_client.py b/meshcore_nomad_bridge/nomad_client.py index 5bf6ef9..40b4a17 100644 --- a/meshcore_nomad_bridge/nomad_client.py +++ b/meshcore_nomad_bridge/nomad_client.py @@ -3,16 +3,27 @@ from __future__ import annotations import asyncio +import ipaddress import json import logging +import socket +import ssl +from collections.abc import Awaitable, Callable from pathlib import Path from time import monotonic -from typing import Any, Awaitable, Callable -from urllib import error as urlerror -from urllib import request as urlrequest +from typing import Any +from urllib.parse import urlsplit, urlunsplit logger = logging.getLogger(__name__) +MAX_HTTP_RESPONSE_BYTES = 262_144 +MAX_HTTP_HEADER_BYTES = 65_536 +HTTP_TOKEN_BYTES = frozenset( + b"!#$%&'*+-.^_`|~0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz" +) +MAX_SESSION_HISTORY_MESSAGES = 20 +MAX_SESSION_HISTORY_BYTES = 16_384 + class NomadUnavailable(RuntimeError): """Raised when NOMAD cannot provide a valid response.""" @@ -42,6 +53,7 @@ def __init__( self._one_shot = one_shot self._session_map_path = Path(session_map_path) self._state_lock = asyncio.Lock() + self._sender_locks: dict[str, asyncio.Lock] = {} if http_request is not None: self._http_request = http_request @@ -65,6 +77,11 @@ async def ask_for_sender(self, sender_id: str, prompt: str) -> str: if self._one_shot: return await self.ask(prompt) + sender_lock = await self._get_sender_lock(sender_id) + async with sender_lock: + return await self._ask_persistent_for_sender(sender_id, prompt) + + async def _ask_persistent_for_sender(self, sender_id: str, prompt: str) -> str: session_id = await self._get_or_create_session_id(sender_id) try: history = await self._get_session_messages(session_id) @@ -80,7 +97,9 @@ async def ask_for_sender(self, sender_id: str, prompt: str) -> str: async def reset_session_for_sender(self, sender_id: str) -> int: if self._one_shot: raise NomadUnavailable("reset_not_supported_in_one_shot") - return await self._create_fresh_session_for_sender(sender_id) + sender_lock = await self._get_sender_lock(sender_id) + async with sender_lock: + return await self._create_fresh_session_for_sender(sender_id) async def _chat( self, @@ -118,7 +137,7 @@ async def _chat( elapsed = monotonic() - start logger.info("NOMAD completed in %.2fs", elapsed) - if status_code >= 400: + if not 200 <= status_code < 300: logger.warning("NOMAD request failed with HTTP %s", status_code) raise NomadUnavailable(f"http_{status_code}") @@ -156,6 +175,10 @@ async def _get_or_create_session_id(self, sender_id: str) -> int: return existing return await self._create_fresh_session_for_sender(sender_id) + async def _get_sender_lock(self, sender_id: str) -> asyncio.Lock: + async with self._state_lock: + return self._sender_locks.setdefault(sender_id, asyncio.Lock()) + async def _create_fresh_session_for_sender(self, sender_id: str) -> int: payload: dict[str, object] = { "title": f"MeshCore {sender_id}", @@ -176,7 +199,7 @@ async def _create_fresh_session_for_sender(self, sender_id: str) -> int: logger.warning("NOMAD network error while creating session: %s", exc) raise NomadUnavailable("network") from exc - if status_code >= 400: + if not 200 <= status_code < 300: logger.warning("NOMAD create session failed with HTTP %s", status_code) raise NomadUnavailable(f"http_{status_code}") @@ -210,7 +233,7 @@ async def _get_session_messages(self, session_id: int) -> list[dict[str, str]]: if status_code == 404: raise _NomadSessionNotFound(session_id) - if status_code >= 400: + if not 200 <= status_code < 300: logger.warning("NOMAD get session failed with HTTP %s", status_code) raise NomadUnavailable(f"http_{status_code}") @@ -233,7 +256,7 @@ async def _get_session_messages(self, session_id: int) -> list[dict[str, str]]: if not isinstance(content, str): continue messages.append({"role": role, "content": content}) - return messages + return _limit_session_history(messages) def _load_session_map(self) -> dict[str, int]: try: @@ -291,6 +314,23 @@ def _parse_json_object(body: str) -> dict[str, Any]: return data +def _limit_session_history(messages: list[dict[str, str]]) -> list[dict[str, str]]: + selected: list[dict[str, str]] = [] + used_bytes = 0 + for message in reversed(messages): + if len(selected) >= MAX_SESSION_HISTORY_MESSAGES: + break + message_bytes = len(message["role"].encode("utf-8")) + len( + message["content"].encode("utf-8") + ) + if used_bytes + message_bytes > MAX_SESSION_HISTORY_BYTES: + break + selected.append(message) + used_bytes += message_bytes + selected.reverse() + return selected + + def _wrap_post_only(http_post: HttpPost) -> HttpRequest: async def _request( *, @@ -316,24 +356,204 @@ async def _default_http_request( timeout_seconds: float, ) -> tuple[int, str]: body = json.dumps(payload).encode("utf-8") if payload is not None else None + try: + return await asyncio.wait_for( + _async_http_request(method=method, url=url, body=body), + timeout=timeout_seconds, + ) + except asyncio.TimeoutError as exc: + raise TimeoutError("request_deadline_exceeded") from exc - def _send() -> tuple[int, str]: - req = urlrequest.Request( - url, - data=body, - headers={"Content-Type": "application/json"}, - method=method, + +async def _async_http_request( + *, + method: str, + url: str, + body: bytes | None, +) -> tuple[int, str]: + parsed = urlsplit(url) + if parsed.scheme not in {"http", "https"} or not parsed.hostname: + raise OSError("invalid_url") + try: + address = ipaddress.ip_address(parsed.hostname) + port = parsed.port or (443 if parsed.scheme == "https" else 80) + except ValueError as exc: + raise OSError("nomad_url_requires_ip_address") from exc + + family = socket.AF_INET6 if address.version == 6 else socket.AF_INET + sock = socket.socket(family, socket.SOCK_STREAM) + sock.setblocking(False) + writer: asyncio.StreamWriter | None = None + path = urlunsplit(("", "", parsed.path or "/", parsed.query, "")) + try: + await asyncio.get_running_loop().sock_connect(sock, (str(address), port)) + ssl_context = ssl.create_default_context() if parsed.scheme == "https" else None + reader, writer = await asyncio.open_connection( + sock=sock, + ssl=ssl_context, + server_hostname=parsed.hostname if ssl_context else None, ) + sock = None + request = _build_http_request(method=method, parsed=parsed, path=path, body=body) + writer.write(request) + await writer.drain() + return await _read_async_http_response(reader, method) + finally: + if writer is not None: + writer.transport.abort() + elif sock is not None: + sock.close() + + +def _build_http_request(*, method: str, parsed: Any, path: str, body: bytes | None) -> bytes: + if method not in {"GET", "POST"}: + raise OSError("unsupported_http_method") + if not path or any(ord(char) < 33 or ord(char) > 126 for char in path): + raise OSError("invalid_url") + default_port = 443 if parsed.scheme == "https" else 80 + host = f"[{parsed.hostname}]" if ":" in parsed.hostname else parsed.hostname + if parsed.port and parsed.port != default_port: + host = f"{host}:{parsed.port}" + headers = [ + f"{method} {path} HTTP/1.1", + f"Host: {host}", + "Accept: application/json", + "Connection: close", + ] + if body is not None: + headers.extend(("Content-Type: application/json", f"Content-Length: {len(body)}")) + try: + return ("\r\n".join(headers) + "\r\n\r\n").encode("ascii") + (body or b"") + except UnicodeEncodeError as exc: + raise OSError("invalid_url") from exc + + +async def _read_async_http_response( + reader: asyncio.StreamReader, method: str +) -> tuple[int, str]: + status_line = await _read_http_line(reader) + if not status_line.endswith(b"\r\n"): + raise OSError("invalid_http_response") + parts = status_line[:-2].split(b" ", 2) + try: + version, raw_status = parts[0], parts[1] + if version not in {b"HTTP/1.0", b"HTTP/1.1"}: + raise ValueError + if len(raw_status) != 3 or not raw_status.isdigit(): + raise ValueError + if len(parts) == 3 and any( + (byte < 32 and byte != 9) or byte == 127 for byte in parts[2] + ): + raise ValueError + status = int(raw_status) + if not 100 <= status <= 599: + raise ValueError + except (IndexError, ValueError) as exc: + raise OSError("invalid_http_response") from exc + + header_bytes = len(status_line) + headers: dict[str, str] = {} + while True: + line = await _read_http_line(reader) + header_bytes += len(line) + if not line or header_bytes > MAX_HTTP_HEADER_BYTES: + raise OSError("invalid_http_response") + if line == b"\r\n": + break + key, value = _parse_header_line(line) + if key in headers: + raise OSError("invalid_http_response") + headers[key] = value + + transfer_encoding = headers.get("transfer-encoding") + if transfer_encoding is not None and "content-length" in headers: + raise OSError("invalid_http_response") + if method == "HEAD" or status in {204, 304} or 100 <= status < 200: + raw = b"" + elif transfer_encoding is not None: + if transfer_encoding.lower() != "chunked": + raise OSError("invalid_http_response") + raw = await _read_chunked_body(reader) + elif "content-length" in headers: + raw_length = headers["content-length"] + if not raw_length.isascii() or not raw_length.isdigit(): + raise OSError("invalid_http_response") + length = int(raw_length) + if length > MAX_HTTP_RESPONSE_BYTES: + raise OSError("response_too_large") + try: + raw = await reader.readexactly(length) + except asyncio.IncompleteReadError as exc: + raise OSError("invalid_http_response") from exc + else: + try: + raw = await reader.readexactly(MAX_HTTP_RESPONSE_BYTES + 1) + except asyncio.IncompleteReadError as exc: + raw = exc.partial + if len(raw) > MAX_HTTP_RESPONSE_BYTES: + raise OSError("response_too_large") + return status, raw.decode("utf-8", errors="replace") + + +async def _read_chunked_body(reader: asyncio.StreamReader) -> bytes: + body = bytearray() + framing_bytes = 0 + while True: + size_line = await _read_http_line(reader) + framing_bytes += len(size_line) + if framing_bytes > MAX_HTTP_HEADER_BYTES: + raise OSError("invalid_http_response") + if not size_line.endswith(b"\r\n") or b";" in size_line: + raise OSError("invalid_http_response") + raw_size = size_line[:-2] + if not raw_size or any(byte not in b"0123456789abcdefABCDEF" for byte in raw_size): + raise OSError("invalid_http_response") try: - with urlrequest.urlopen(req, timeout=timeout_seconds) as response: - raw = response.read().decode("utf-8", errors="replace") - return int(response.status), raw - except urlerror.HTTPError as exc: - raw = exc.read().decode("utf-8", errors="replace") if exc.fp else "" - return int(exc.code), raw - except TimeoutError: - raise - except OSError: - raise - - return await asyncio.to_thread(_send) + size = int(raw_size, 16) + except ValueError as exc: + raise OSError("invalid_http_response") from exc + if size == 0: + while True: + trailer = await _read_http_line(reader) + framing_bytes += len(trailer) + if framing_bytes > MAX_HTTP_HEADER_BYTES: + raise OSError("invalid_http_response") + if trailer == b"\r\n": + return bytes(body) + if not trailer: + raise OSError("invalid_http_response") + key, _ = _parse_header_line(trailer) + if key in {"content-length", "transfer-encoding"}: + raise OSError("invalid_http_response") + if size < 0 or len(body) + size > MAX_HTTP_RESPONSE_BYTES: + raise OSError("response_too_large") + try: + body.extend(await reader.readexactly(size)) + if await reader.readexactly(2) != b"\r\n": + raise OSError("invalid_http_response") + except asyncio.IncompleteReadError as exc: + raise OSError("invalid_http_response") from exc + + +async def _read_http_line(reader: asyncio.StreamReader) -> bytes: + try: + line = await reader.readline() + except ValueError as exc: + raise OSError("invalid_http_response") from exc + if len(line) > MAX_HTTP_HEADER_BYTES: + raise OSError("invalid_http_response") + return line + + +def _parse_header_line(line: bytes) -> tuple[str, str]: + if not line.endswith(b"\r\n"): + raise OSError("invalid_http_response") + try: + raw_name, raw_value = line[:-2].split(b":", 1) + except ValueError as exc: + raise OSError("invalid_http_response") from exc + if not raw_name or any(byte not in HTTP_TOKEN_BYTES for byte in raw_name): + raise OSError("invalid_http_response") + if any((byte < 32 and byte != 9) or byte == 127 for byte in raw_value): + raise OSError("invalid_http_response") + return raw_name.decode("ascii").lower(), raw_value.decode("iso-8859-1").strip() diff --git a/openhop-plugin.json b/openhop-plugin.json index c921279..adbd38b 100644 --- a/openhop-plugin.json +++ b/openhop-plugin.json @@ -2,7 +2,7 @@ "schema": 1, "id": "openhop.nomad", "name": "NOMAD Bridge", - "version": "0.1.1", + "version": "0.1.2", "description": "Bridges an openHop Repeater Companion identity to Project N.O.M.A.D.", "runtime": { "type": "python", @@ -17,7 +17,12 @@ "nomad_collection": null, "nomad_timeout_seconds": 120, "one_shot": true, - "max_concurrent_requests": 2, + "max_concurrent_requests": 1, + "max_pending_requests": 1, + "max_requests_per_sender": 2, + "max_requests_global": 4, + "rate_limit_window_seconds": 60, + "allowed_sender_prefixes": [], "busy_wait_seconds": 5, "max_reply_chunks": 4, "max_chunk_bytes": 145, diff --git a/pyproject.toml b/pyproject.toml index 814d8e2..dee3bb4 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "openhop-nomad-plugin" -version = "0.1.1" +version = "0.1.2" description = "openHop service plugin bridging a MeshCore Companion identity to Project N.O.M.A.D." readme = "README.md" requires-python = ">=3.10" diff --git a/tests/test_config.py b/tests/test_config.py index 6b16bc1..9ddb185 100644 --- a/tests/test_config.py +++ b/tests/test_config.py @@ -5,7 +5,6 @@ from meshcore_nomad_bridge.config import ConfigError, Settings - _ENV_KEYS = ( "OPENHOP_PLUGIN_DATA", "MESHCORE_HOST", @@ -17,6 +16,11 @@ "ONE_SHOT", "NOMAD_SESSION_MAP_PATH", "MAX_CONCURRENT_REQUESTS", + "MAX_PENDING_REQUESTS", + "MAX_REQUESTS_PER_SENDER", + "MAX_REQUESTS_GLOBAL", + "RATE_LIMIT_WINDOW_SECONDS", + "ALLOWED_SENDER_PREFIXES", "NOMAD_BUSY_WAIT_SECONDS", "MAX_REPLY_CHUNKS", "MAX_CHUNK_BYTES", @@ -45,10 +49,15 @@ def test_settings_load_plugin_owned_config(monkeypatch: pytest.MonkeyPatch, tmp_ data_dir, meshcore_host="127.0.0.2", meshcore_port=5056, - nomad_url="http://nomad.local:8080", + nomad_url="http://10.5.30.7:8080", nomad_model="qwen-test", - one_shot=False, + one_shot=True, max_reply_chunks=3, + max_pending_requests=3, + max_requests_per_sender=2, + max_requests_global=7, + rate_limit_window_seconds=45, + allowed_sender_prefixes=["010203040506", "aabbccddeeff"], ) monkeypatch.setenv("OPENHOP_PLUGIN_DATA", str(data_dir)) @@ -56,10 +65,15 @@ def test_settings_load_plugin_owned_config(monkeypatch: pytest.MonkeyPatch, tmp_ assert settings.meshcore_host == "127.0.0.2" assert settings.meshcore_port == 5056 - assert settings.nomad_url == "http://nomad.local:8080" + assert settings.nomad_url == "http://10.5.30.7:8080" assert settings.nomad_model == "qwen-test" - assert settings.one_shot is False + assert settings.one_shot is True assert settings.max_reply_chunks == 3 + assert settings.max_pending_requests == 3 + assert settings.max_requests_per_sender == 2 + assert settings.max_requests_global == 7 + assert settings.rate_limit_window_seconds == 45 + assert settings.allowed_sender_prefixes == ("010203040506", "aabbccddeeff") def test_environment_overrides_plugin_config(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> None: @@ -67,18 +81,18 @@ def test_environment_overrides_plugin_config(monkeypatch: pytest.MonkeyPatch, tm data_dir = tmp_path / "plugin-data" _write_config( data_dir, - nomad_url="http://from-file:8080", + nomad_url="http://10.5.30.8:8080", nomad_model="file-model", meshcore_port=5001, ) monkeypatch.setenv("OPENHOP_PLUGIN_DATA", str(data_dir)) - monkeypatch.setenv("NOMAD_URL", "http://from-env:8080") + monkeypatch.setenv("NOMAD_URL", "http://10.5.30.9:8080") monkeypatch.setenv("NOMAD_MODEL", "env-model") monkeypatch.setenv("MESHCORE_PORT", "6001") settings = Settings.from_env() - assert settings.nomad_url == "http://from-env:8080" + assert settings.nomad_url == "http://10.5.30.9:8080" assert settings.nomad_model == "env-model" assert settings.meshcore_port == 6001 @@ -86,7 +100,7 @@ def test_environment_overrides_plugin_config(monkeypatch: pytest.MonkeyPatch, tm def test_session_map_defaults_to_plugin_data(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> None: _clean_env(monkeypatch) data_dir = tmp_path / "plugin-data" - _write_config(data_dir, nomad_url="http://nomad.local", nomad_model="test-model") + _write_config(data_dir, nomad_url="http://10.5.30.7", nomad_model="test-model") monkeypatch.setenv("OPENHOP_PLUGIN_DATA", str(data_dir)) settings = Settings.from_env() @@ -96,7 +110,7 @@ def test_session_map_defaults_to_plugin_data(monkeypatch: pytest.MonkeyPatch, tm def test_standalone_session_map_default_is_preserved(monkeypatch: pytest.MonkeyPatch) -> None: _clean_env(monkeypatch) - monkeypatch.setenv("NOMAD_URL", "http://nomad.local") + monkeypatch.setenv("NOMAD_URL", "http://10.5.30.7") monkeypatch.setenv("NOMAD_MODEL", "test-model") settings = Settings.from_env() @@ -104,6 +118,104 @@ def test_standalone_session_map_default_is_preserved(monkeypatch: pytest.MonkeyP assert settings.nomad_session_map_path == "./data/nomad_sessions.json" +def test_persistent_mode_is_rejected_even_with_sender_allowlist( + monkeypatch: pytest.MonkeyPatch, tmp_path: Path +) -> None: + _clean_env(monkeypatch) + data_dir = tmp_path / "plugin-data" + _write_config( + data_dir, + nomad_url="http://10.5.30.7", + nomad_model="test-model", + one_shot=False, + allowed_sender_prefixes=["010203040506"], + ) + monkeypatch.setenv("OPENHOP_PLUGIN_DATA", str(data_dir)) + + with pytest.raises(ConfigError, match="ONE_SHOT"): + Settings.from_env() + + +@pytest.mark.parametrize( + ("name", "value"), + [ + ("NOMAD_TIMEOUT_SECONDS", "nan"), + ("NOMAD_TIMEOUT_SECONDS", "inf"), + ("RATE_LIMIT_WINDOW_SECONDS", "nan"), + ("NOMAD_BUSY_WAIT_SECONDS", "-inf"), + ], +) +def test_nonfinite_numeric_settings_are_rejected( + monkeypatch: pytest.MonkeyPatch, name: str, value: str +) -> None: + _clean_env(monkeypatch) + monkeypatch.setenv("NOMAD_URL", "http://10.5.30.7:8080") + monkeypatch.setenv("NOMAD_MODEL", "test-model") + monkeypatch.setenv(name, value) + + with pytest.raises(ConfigError, match="finite"): + Settings.from_env() + + +@pytest.mark.parametrize( + "url", + [ + "http://nomad.local:8080", + "http://10.5.30.7:bad", + "http://10.5.30.7:70000", + "http://10.5.30.7/base", + "http://10.5.30.7/?query=1", + ], +) +def test_nomad_url_requires_valid_ip_literal( + monkeypatch: pytest.MonkeyPatch, url: str +) -> None: + _clean_env(monkeypatch) + monkeypatch.setenv("NOMAD_URL", url) + monkeypatch.setenv("NOMAD_MODEL", "test-model") + + with pytest.raises(ConfigError, match="NOMAD_URL"): + Settings.from_env() + + +@pytest.mark.parametrize("template", ["{question} {missing}", "{{question}}"]) +def test_radio_prompt_template_rejects_invalid_fields( + monkeypatch: pytest.MonkeyPatch, template: str +) -> None: + _clean_env(monkeypatch) + monkeypatch.setenv("NOMAD_URL", "http://10.5.30.7:8080") + monkeypatch.setenv("NOMAD_MODEL", "test-model") + monkeypatch.setenv("RADIO_PROMPT_TEMPLATE", template) + + with pytest.raises(ConfigError, match="RADIO_PROMPT_TEMPLATE"): + Settings.from_env() + + +def test_environment_parses_comma_separated_sender_allowlist( + monkeypatch: pytest.MonkeyPatch, +) -> None: + _clean_env(monkeypatch) + monkeypatch.setenv("NOMAD_URL", "http://10.5.30.7") + monkeypatch.setenv("NOMAD_MODEL", "test-model") + monkeypatch.setenv("ALLOWED_SENDER_PREFIXES", "010203040506, AABBCCDDEEFF") + + settings = Settings.from_env() + + assert settings.allowed_sender_prefixes == ("010203040506", "aabbccddeeff") + + +def test_invalid_sender_allowlist_entry_is_rejected( + monkeypatch: pytest.MonkeyPatch, +) -> None: + _clean_env(monkeypatch) + monkeypatch.setenv("NOMAD_URL", "http://10.5.30.7") + monkeypatch.setenv("NOMAD_MODEL", "test-model") + monkeypatch.setenv("ALLOWED_SENDER_PREFIXES", "not-hex") + + with pytest.raises(ConfigError, match="12 hexadecimal"): + Settings.from_env() + + def test_invalid_plugin_config_is_reported(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> None: _clean_env(monkeypatch) data_dir = tmp_path / "plugin-data" diff --git a/tests/test_message_handler.py b/tests/test_message_handler.py index 11023bd..7e6d108 100644 --- a/tests/test_message_handler.py +++ b/tests/test_message_handler.py @@ -5,10 +5,12 @@ from meshcore_nomad_bridge.config import Settings from meshcore_nomad_bridge.main import ( - BridgeService, NOMAD_BUSY_MESSAGE, NOMAD_RESET_MESSAGE, NOMAD_TOO_LONG_MESSAGE, + BridgeService, + DuplicateCache, + RequestRateLimiter, ) from meshcore_nomad_bridge.meshcore_client import IncomingMessage @@ -64,6 +66,11 @@ def _settings() -> Settings: one_shot=True, nomad_session_map_path="./data/test_nomad_sessions.json", max_concurrent_requests=2, + max_pending_requests=4, + max_requests_per_sender=10, + max_requests_global=50, + rate_limit_window_seconds=60.0, + allowed_sender_prefixes=("010203040506", "111213141516"), busy_wait_seconds=0.2, max_reply_chunks=4, max_chunk_bytes=145, @@ -75,9 +82,9 @@ def _settings() -> Settings: ) -def _msg(text: str, ts: int = 1) -> IncomingMessage: +def _msg(text: str, ts: int = 1, sender: bytes = b"\x01\x02\x03\x04\x05\x06") -> IncomingMessage: return IncomingMessage( - sender_prefix=b"\x01\x02\x03\x04\x05\x06", + sender_prefix=sender, text=text, timestamp=ts, txt_type=0, @@ -86,6 +93,25 @@ def _msg(text: str, ts: int = 1) -> IncomingMessage: ) +class ExplodingNomad(SlowNomad): + async def ask_for_sender(self, sender_id: str, prompt: str) -> str: + raise RuntimeError("unexpected failure") + + +@pytest.mark.asyncio +async def test_background_task_exception_is_retrieved_and_logged(caplog) -> None: + service = BridgeService( + settings=_settings(), meshcore=FakeMeshCore(), nomad=ExplodingNomad(delay=0) + ) + + await service._dispatch_message(_msg("explode", ts=700)) + tasks = list(service._inflight) + await asyncio.gather(*tasks, return_exceptions=True) + await asyncio.sleep(0) + + assert "Unhandled message task error" in caplog.text + + @pytest.mark.asyncio async def test_duplicate_message_processed_once() -> None: settings = _settings() @@ -99,6 +125,14 @@ async def test_duplicate_message_processed_once() -> None: assert nomad.calls == 1 +def test_duplicate_key_uses_sha256_digest() -> None: + service = BridgeService(settings=_settings(), meshcore=FakeMeshCore(), nomad=SlowNomad()) + + digest = service._dedupe_key(_msg("hello", ts=100)).rsplit(":", 1)[-1] + + assert len(digest) == 64 + + @pytest.mark.asyncio async def test_prompt_byte_limit_rejected() -> None: settings = replace(_settings(), max_prompt_bytes=10) @@ -144,6 +178,179 @@ async def test_overload_returns_busy_message() -> None: assert any(text == NOMAD_BUSY_MESSAGE for _, text in mesh.sent) +@pytest.mark.asyncio +async def test_zero_busy_wait_uses_immediately_available_capacity() -> None: + settings = replace(_settings(), max_concurrent_requests=1, busy_wait_seconds=0) + mesh = FakeMeshCore() + nomad = SlowNomad(delay=0) + service = BridgeService(settings=settings, meshcore=mesh, nomad=nomad) + + await service._handle_message(_msg("first", ts=13)) + + assert nomad.calls == 1 + assert all(text != NOMAD_BUSY_MESSAGE for _, text in mesh.sent) + + +def test_duplicate_cache_expires_entries_below_prune_threshold(monkeypatch) -> None: + now = 100.0 + monkeypatch.setattr("meshcore_nomad_bridge.main.time.monotonic", lambda: now) + cache = DuplicateCache(ttl_seconds=60) + + assert cache.seen("key") is False + now = 161.0 + + assert cache.seen("key") is False + + +@pytest.mark.asyncio +async def test_dispatch_silently_drops_when_pending_limit_is_full() -> None: + settings = replace( + _settings(), + max_concurrent_requests=1, + max_pending_requests=1, + busy_wait_seconds=1.0, + ) + mesh = FakeMeshCore() + nomad = SlowNomad(delay=0.05) + service = BridgeService(settings=settings, meshcore=mesh, nomad=nomad) + + await service._dispatch_message(_msg("first", ts=201)) + await service._dispatch_message(_msg("second", ts=202)) + await asyncio.gather(*list(service._inflight)) + + assert nomad.calls == 1 + assert [text for _, text in mesh.sent] == ["A concise answer from NOMAD."] + + +@pytest.mark.asyncio +async def test_pending_rejection_does_not_mark_message_as_duplicate() -> None: + settings = replace( + _settings(), + max_concurrent_requests=1, + max_pending_requests=1, + max_requests_per_sender=3, + busy_wait_seconds=1.0, + ) + mesh = FakeMeshCore() + nomad = SlowNomad(delay=0.05) + service = BridgeService(settings=settings, meshcore=mesh, nomad=nomad) + rejected = _msg("retry me", ts=211) + + await service._dispatch_message(_msg("first", ts=210)) + await service._dispatch_message(rejected) + await asyncio.gather(*list(service._inflight)) + await service._dispatch_message(rejected) + await asyncio.gather(*list(service._inflight)) + + assert nomad.calls == 2 + + +@pytest.mark.asyncio +async def test_dispatch_silently_rate_limits_each_sender() -> None: + settings = replace( + _settings(), + max_pending_requests=4, + max_requests_per_sender=1, + max_requests_global=10, + ) + mesh = FakeMeshCore() + nomad = SlowNomad(delay=0) + service = BridgeService(settings=settings, meshcore=mesh, nomad=nomad) + + await service._dispatch_message(_msg("first", ts=301)) + await service._dispatch_message(_msg("second", ts=302)) + await asyncio.gather(*list(service._inflight)) + + assert nomad.calls == 1 + assert [text for _, text in mesh.sent] == ["A concise answer from NOMAD."] + + +@pytest.mark.asyncio +async def test_dispatch_silently_enforces_global_rate_limit() -> None: + settings = replace( + _settings(), + max_pending_requests=4, + max_requests_per_sender=3, + max_requests_global=1, + ) + mesh = FakeMeshCore() + nomad = SlowNomad(delay=0) + service = BridgeService(settings=settings, meshcore=mesh, nomad=nomad) + + await service._dispatch_message( + _msg("first", ts=401, sender=b"\x01\x02\x03\x04\x05\x06") + ) + await service._dispatch_message( + _msg("second", ts=402, sender=b"\x11\x12\x13\x14\x15\x16") + ) + await asyncio.gather(*list(service._inflight)) + + assert nomad.calls == 1 + assert [text for _, text in mesh.sent] == ["A concise answer from NOMAD."] + + +@pytest.mark.asyncio +async def test_empty_sender_allowlist_rejects_all_requests() -> None: + settings = replace(_settings(), allowed_sender_prefixes=()) + mesh = FakeMeshCore() + nomad = SlowNomad(delay=0) + service = BridgeService(settings=settings, meshcore=mesh, nomad=nomad) + + await service._dispatch_message(_msg("blocked", ts=499)) + await asyncio.sleep(0) + + assert nomad.calls == 0 + assert mesh.sent == [] + + +@pytest.mark.asyncio +async def test_dispatch_silently_rejects_sender_outside_allowlist() -> None: + settings = replace( + _settings(), + allowed_sender_prefixes=("aabbccddeeff",), + ) + mesh = FakeMeshCore() + nomad = SlowNomad(delay=0) + service = BridgeService(settings=settings, meshcore=mesh, nomad=nomad) + + await service._dispatch_message(_msg("not allowed", ts=501)) + await asyncio.gather(*list(service._inflight)) + + assert nomad.calls == 0 + assert mesh.sent == [] + + +def test_rate_limiter_allows_requests_after_window_expires(monkeypatch) -> None: + now = [100.0] + monkeypatch.setattr("meshcore_nomad_bridge.main.time.monotonic", lambda: now[0]) + limiter = RequestRateLimiter(per_sender=1, global_limit=1, window_seconds=60) + + assert limiter.allow("sender") is True + assert limiter.allow("sender") is False + now[0] = 161.0 + assert limiter.allow("sender") is True + + +@pytest.mark.asyncio +async def test_duplicate_retransmission_does_not_consume_rate_capacity() -> None: + settings = replace( + _settings(), + max_pending_requests=4, + max_requests_per_sender=2, + max_requests_global=10, + ) + mesh = FakeMeshCore() + nomad = SlowNomad(delay=0) + service = BridgeService(settings=settings, meshcore=mesh, nomad=nomad) + + await service._dispatch_message(_msg("first", ts=601)) + await service._dispatch_message(_msg("first", ts=601)) + await service._dispatch_message(_msg("second", ts=602)) + await asyncio.gather(*list(service._inflight)) + + assert nomad.calls == 2 + + @pytest.mark.asyncio async def test_reset_command_starts_new_session_in_persistent_mode() -> None: settings = replace(_settings(), one_shot=False) diff --git a/tests/test_nomad_client.py b/tests/test_nomad_client.py index 1f6783a..8ef6b1a 100644 --- a/tests/test_nomad_client.py +++ b/tests/test_nomad_client.py @@ -1,8 +1,249 @@ +import asyncio +import json +import threading +import time +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer + import pytest +from meshcore_nomad_bridge import nomad_client from meshcore_nomad_bridge.nomad_client import NomadClient, NomadUnavailable +async def _request_raw_response(raw_response: bytes) -> tuple[int, str]: + async def handle(reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None: + await reader.readuntil(b"\r\n\r\n") + writer.write(raw_response) + await writer.drain() + writer.close() + + server = await asyncio.start_server(handle, "127.0.0.1", 0) + port = server.sockets[0].getsockname()[1] + async with server: + return await nomad_client._default_http_request( + method="GET", + url=f"http://127.0.0.1:{port}/test", + payload=None, + timeout_seconds=2, + ) + + +class RedirectHandler(BaseHTTPRequestHandler): + followed = False + + def do_GET(self) -> None: + if self.path == "/redirect": + self.send_response(302) + self.send_header("Location", "/followed") + self.send_header("Content-Length", "0") + self.end_headers() + return + type(self).followed = True + self.send_response(200) + self.send_header("Content-Length", "0") + self.end_headers() + + def log_message(self, format: str, *args: object) -> None: + return None + + +def test_http_request_rejects_control_characters_in_target() -> None: + parsed = nomad_client.urlsplit("http://127.0.0.1") + + with pytest.raises(OSError, match="invalid_url"): + nomad_client._build_http_request( + method="GET", parsed=parsed, path="/ok\r\nX-Injected: yes", body=None + ) + + +@pytest.mark.asyncio +async def test_default_http_request_does_not_follow_redirects() -> None: + RedirectHandler.followed = False + server = ThreadingHTTPServer(("127.0.0.1", 0), RedirectHandler) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + status, body = await nomad_client._default_http_request( + method="GET", + url=f"http://127.0.0.1:{server.server_port}/redirect", + payload=None, + timeout_seconds=2, + ) + finally: + server.shutdown() + server.server_close() + thread.join(timeout=1) + + assert status == 302 + assert body == "" + assert RedirectHandler.followed is False + + +class SlowTrickleHandler(BaseHTTPRequestHandler): + disconnected = threading.Event() + + def do_GET(self) -> None: + self.send_response(200) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", "20") + self.end_headers() + for _ in range(20): + try: + self.wfile.write(b"x") + self.wfile.flush() + except (BrokenPipeError, ConnectionResetError): + type(self).disconnected.set() + break + time.sleep(0.05) + + def log_message(self, format: str, *args: object) -> None: + return None + + +@pytest.mark.asyncio +async def test_default_http_request_enforces_total_deadline(monkeypatch) -> None: + async def forbid_worker_thread(*args: object, **kwargs: object): + raise AssertionError("HTTP transport must not use a non-cancellable worker thread") + + monkeypatch.setattr(asyncio, "to_thread", forbid_worker_thread) + server = ThreadingHTTPServer(("127.0.0.1", 0), SlowTrickleHandler) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + started = time.monotonic() + try: + with pytest.raises((OSError, TimeoutError)): + await nomad_client._default_http_request( + method="GET", + url=f"http://127.0.0.1:{server.server_port}/slow", + payload=None, + timeout_seconds=0.15, + ) + finally: + server.shutdown() + server.server_close() + thread.join(timeout=1) + + assert time.monotonic() - started < 0.7 + + +@pytest.mark.asyncio +async def test_connection_close_body_waits_for_all_fragments() -> None: + async def handle(reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None: + await reader.readuntil(b"\r\n\r\n") + writer.write(b"HTTP/1.1 200 OK\r\nConnection: close\r\n\r\nabc") + await writer.drain() + await asyncio.sleep(0.05) + writer.write(b"def") + await writer.drain() + writer.close() + + server = await asyncio.start_server(handle, "127.0.0.1", 0) + port = server.sockets[0].getsockname()[1] + async with server: + status, body = await nomad_client._default_http_request( + method="GET", + url=f"http://127.0.0.1:{port}/fragmented", + payload=None, + timeout_seconds=2, + ) + + assert status == 200 + assert body == "abcdef" + + +@pytest.mark.asyncio +async def test_oversized_http_header_is_a_controlled_network_error() -> None: + raw = b"HTTP/1.1 200 OK\r\nX-Large: " + (b"x" * 70_000) + b"\r\n\r\n" + + with pytest.raises(OSError, match="invalid_http_response"): + await _request_raw_response(raw) + + +@pytest.mark.asyncio +async def test_chunked_response_is_decoded() -> None: + status, body = await _request_raw_response( + b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n" + b"3\r\nabc\r\n3\r\ndef\r\n0\r\n\r\n" + ) + + assert status == 200 + assert body == "abcdef" + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "raw_response", + [ + b"HTTP/1.1 20 Weird\r\nContent-Length: 2\r\n\r\n{}", + b"HTTP/1.1 200 bad\x00reason\r\nContent-Length: 2\r\n\r\n{}", + b"HTTP/1.1 200 OK\r\nContent-Length: +2\r\n\r\n{}", + ( + b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\nContent-Length: 2\r\n\r\n" + b"2\r\n{}\r\n0\r\n\r\n" + ), + b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nContent-Length: 2\r\n\r\n{}", + ( + b"HTTP/1.1 200 OK\r\nTransfer-Encoding: notchunked\r\n\r\n" + b"2\r\n{}\r\n0\r\n\r\n" + ), + ( + b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n" + b"2\r\n{}\r\n0\r\ngarbage\r\n\r\n" + ), + ( + b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n" + b"2;bad=\x00\r\n{}\r\n0\r\n\r\n" + ), + ], +) +async def test_malformed_http_framing_is_rejected(raw_response: bytes) -> None: + with pytest.raises(OSError, match="invalid_http_response"): + await _request_raw_response(raw_response) + + +@pytest.mark.asyncio +async def test_chunked_response_rejects_missing_terminal_line() -> None: + reader = asyncio.StreamReader() + reader.feed_data(b"1\r\nx\r\n0\r\n") + reader.feed_eof() + + with pytest.raises(OSError, match="invalid_http_response"): + await nomad_client._read_chunked_body(reader) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("size", [0, 262_144]) +async def test_response_body_accepts_exact_size_boundary(size: int) -> None: + status, body = await _request_raw_response( + f"HTTP/1.1 200 OK\r\nContent-Length: {size}\r\n\r\n".encode() + b"x" * size + ) + + assert status == 200 + assert len(body) == size + + +@pytest.mark.asyncio +async def test_default_http_request_converts_malformed_http_to_oserror() -> None: + with pytest.raises(OSError, match="invalid_http_response"): + await _request_raw_response(b"broken\r\n\r\n") + + +@pytest.mark.asyncio +async def test_default_http_request_rejects_oversized_success_body() -> None: + with pytest.raises(OSError, match="response_too_large"): + await _request_raw_response( + b"HTTP/1.1 200 OK\r\nContent-Length: 262145\r\n\r\n" + ) + + +@pytest.mark.asyncio +async def test_default_http_request_rejects_oversized_error_body() -> None: + with pytest.raises(OSError, match="response_too_large"): + await _request_raw_response( + b"HTTP/1.1 500 Error\r\nContent-Length: 262145\r\n\r\n" + ) + + @pytest.mark.asyncio async def test_nomad_client_success() -> None: async def fake_post(*, url: str, payload: dict[str, object], timeout_seconds: float): @@ -28,6 +269,7 @@ async def fake_post(*, url: str, payload: dict[str, object], timeout_seconds: fl "status_code,body", [ (500, '{"error":"boom"}'), + (302, '{"message":{"content":"redirect body"}}'), (200, "not json"), (200, '{"message":{}}'), (200, '{"message":{"content":""}}'), @@ -152,3 +394,81 @@ async def fake_request( assert reset_id == 102 data = session_map.read_text(encoding="utf-8") assert data == '{"sender-b": 102}' + + +@pytest.mark.asyncio +async def test_persistent_mode_serializes_session_creation_per_sender(tmp_path) -> None: + creates = 0 + + async def fake_request( + *, + method: str, + url: str, + payload: dict[str, object] | None, + timeout_seconds: float, + ): + nonlocal creates + _ = (payload, timeout_seconds) + if method == "POST" and url == "http://nomad.local/api/chat/sessions": + creates += 1 + await asyncio.sleep(0.02) + return 201, '{"id":"42"}' + if method == "GET" and url == "http://nomad.local/api/chat/sessions/42": + return 200, '{"messages":[]}' + if method == "POST" and url == "http://nomad.local/api/ollama/chat": + return 200, '{"message":{"content":"ok"}}' + raise AssertionError(f"Unexpected request: {method} {url}") + + client = NomadClient( + base_url="http://nomad.local", + model="test-model", + timeout_seconds=5, + collection=None, + one_shot=False, + session_map_path=str(tmp_path / "sessions.json"), + http_request=fake_request, + ) + + answers = await asyncio.gather( + client.ask_for_sender("sender-c", "first"), + client.ask_for_sender("sender-c", "second"), + ) + + assert answers == ["ok", "ok"] + assert creates == 1 + + +@pytest.mark.asyncio +async def test_persistent_history_is_limited_to_recent_messages(tmp_path) -> None: + messages = [{"role": "user", "content": str(index)} for index in range(25)] + + async def fake_request(**kwargs): + _ = kwargs + return 200, json.dumps({"messages": messages}) + + client = NomadClient( + base_url="http://nomad.local", + model="test-model", + timeout_seconds=5, + collection=None, + one_shot=False, + session_map_path=str(tmp_path / "sessions.json"), + http_request=fake_request, + ) + + history = await client._get_session_messages(42) + + assert len(history) == 20 + assert history[0]["content"] == "5" + assert history[-1]["content"] == "24" + + +def test_persistent_history_is_limited_by_utf8_bytes() -> None: + messages = [ + {"role": "user", "content": "x" * 9_000}, + {"role": "assistant", "content": "y" * 9_000}, + ] + + history = nomad_client._limit_session_history(messages) + + assert history == [messages[-1]] diff --git a/tests/test_plugin_package.py b/tests/test_plugin_package.py index 71db0f7..9be569e 100644 --- a/tests/test_plugin_package.py +++ b/tests/test_plugin_package.py @@ -1,7 +1,7 @@ import json from pathlib import Path -import tomllib +import tomllib ROOT = Path(__file__).resolve().parents[1] From de76d0b870ae038d0e022d317c5aef2b543b4cfc Mon Sep 17 00:00:00 2001 From: fizzlepoof <251490314+fizzlepoof@users.noreply.github.com> Date: Fri, 4 Sep 2026 18:36:08 +0000 Subject: [PATCH 02/27] [verified] Pace NOMAD reply chunks --- README.md | 10 ++++++---- config.default.json | 3 ++- meshcore_nomad_bridge/__init__.py | 2 +- meshcore_nomad_bridge/config.py | 8 +++++++- meshcore_nomad_bridge/main.py | 4 +++- openhop-plugin.json | 5 +++-- pyproject.toml | 2 +- tests/test_config.py | 4 ++++ tests/test_message_handler.py | 24 ++++++++++++++++++++++++ 9 files changed, 51 insertions(+), 11 deletions(-) diff --git a/README.md b/README.md index 7c386f9..b7d865b 100644 --- a/README.md +++ b/README.md @@ -50,7 +50,7 @@ $OPENHOP_PLUGIN_DATA/ └── nomad_sessions.json ``` -Persistent conversations are disabled in v0.1.2 because the upstream session lifecycle cannot +Persistent conversations are disabled in v0.1.3 because the upstream session lifecycle cannot yet be bounded safely. `one_shot` must remain `true`; the session-map setting is retained only for configuration compatibility. @@ -89,7 +89,8 @@ Typical configuration: "allowed_sender_prefixes": ["001122334455"], "busy_wait_seconds": 5, "max_reply_chunks": 4, - "max_chunk_bytes": 145, + "max_chunk_bytes": 80, + "reply_chunk_delay_seconds": 2.0, "max_prompt_bytes": 1000, "radio_prompt_enabled": true, "duplicate_ttl_seconds": 600, @@ -135,6 +136,7 @@ Important environment overrides include: - `NOMAD_BUSY_WAIT_SECONDS` - `MAX_REPLY_CHUNKS` - `MAX_CHUNK_BYTES` +- `REPLY_CHUNK_DELAY_SECONDS` (0–60 seconds between multi-packet reply chunks) - `MAX_PROMPT_BYTES` - `RADIO_PROMPT_ENABLED` - `RADIO_PROMPT_TEMPLATE` @@ -152,7 +154,7 @@ Important environment overrides include: "schema": 1, "id": "openhop.nomad", "name": "NOMAD Bridge", - "version": "0.1.2", + "version": "0.1.3", "runtime": { "type": "python", "entrypoint": "meshcore-nomad-bridge" @@ -196,7 +198,7 @@ python -m build --wheel The wheel is written to `dist/`, for example: ```text -dist/openhop_nomad_plugin-0.1.2-py3-none-any.whl +dist/openhop_nomad_plugin-0.1.3-py3-none-any.whl ``` If you want both wheel and source distribution, run: diff --git a/config.default.json b/config.default.json index 584c75d..f47e6f8 100644 --- a/config.default.json +++ b/config.default.json @@ -14,7 +14,8 @@ "allowed_sender_prefixes": [], "busy_wait_seconds": 5, "max_reply_chunks": 4, - "max_chunk_bytes": 145, + "max_chunk_bytes": 80, + "reply_chunk_delay_seconds": 2.0, "max_prompt_bytes": 1000, "radio_prompt_enabled": true, "duplicate_ttl_seconds": 600, diff --git a/meshcore_nomad_bridge/__init__.py b/meshcore_nomad_bridge/__init__.py index ed11b9d..dad68a3 100644 --- a/meshcore_nomad_bridge/__init__.py +++ b/meshcore_nomad_bridge/__init__.py @@ -2,4 +2,4 @@ __all__ = ["__version__"] -__version__ = "0.1.2" +__version__ = "0.1.3" diff --git a/meshcore_nomad_bridge/config.py b/meshcore_nomad_bridge/config.py index e61a0ab..637a247 100644 --- a/meshcore_nomad_bridge/config.py +++ b/meshcore_nomad_bridge/config.py @@ -44,6 +44,7 @@ class Settings: max_reply_chunks: int max_chunk_bytes: int + reply_chunk_delay_seconds: float max_prompt_bytes: int radio_prompt_enabled: bool @@ -83,7 +84,8 @@ def from_env(cls) -> Settings: busy_wait_seconds = _get_float("NOMAD_BUSY_WAIT_SECONDS", 5.0, config) max_reply_chunks = _get_int("MAX_REPLY_CHUNKS", 4, config) - max_chunk_bytes = _get_int("MAX_CHUNK_BYTES", 145, config) + max_chunk_bytes = _get_int("MAX_CHUNK_BYTES", 80, config) + reply_chunk_delay_seconds = _get_float("REPLY_CHUNK_DELAY_SECONDS", 2.0, config) max_prompt_bytes = _get_int("MAX_PROMPT_BYTES", 1000, config) radio_prompt_enabled = _get_bool("RADIO_PROMPT_ENABLED", True, config) @@ -122,6 +124,7 @@ def from_env(cls) -> Settings: busy_wait_seconds=busy_wait_seconds, max_reply_chunks=max_reply_chunks, max_chunk_bytes=max_chunk_bytes, + reply_chunk_delay_seconds=reply_chunk_delay_seconds, max_prompt_bytes=max_prompt_bytes, radio_prompt_enabled=radio_prompt_enabled, radio_prompt_template=radio_prompt_template, @@ -280,6 +283,7 @@ def _validate(settings: Settings) -> None: "NOMAD_TIMEOUT_SECONDS": settings.nomad_timeout_seconds, "RATE_LIMIT_WINDOW_SECONDS": settings.rate_limit_window_seconds, "NOMAD_BUSY_WAIT_SECONDS": settings.busy_wait_seconds, + "REPLY_CHUNK_DELAY_SECONDS": settings.reply_chunk_delay_seconds, } for name, value in numeric_settings.items(): if not math.isfinite(value): @@ -308,6 +312,8 @@ def _validate(settings: Settings) -> None: raise ConfigError("MAX_REPLY_CHUNKS must be > 0") if settings.max_chunk_bytes < 40: raise ConfigError("MAX_CHUNK_BYTES must be >= 40") + if not 0 <= settings.reply_chunk_delay_seconds <= 60: + raise ConfigError("REPLY_CHUNK_DELAY_SECONDS must be between 0 and 60") if settings.max_prompt_bytes <= 0: raise ConfigError("MAX_PROMPT_BYTES must be > 0") diff --git a/meshcore_nomad_bridge/main.py b/meshcore_nomad_bridge/main.py index 16c60ec..c75ecb2 100644 --- a/meshcore_nomad_bridge/main.py +++ b/meshcore_nomad_bridge/main.py @@ -213,8 +213,10 @@ async def _handle_message( chunks = [NOMAD_UNAVAILABLE_MESSAGE] logger.info("Sending %d MeshCore reply packets", len(chunks)) - for chunk in chunks: + for index, chunk in enumerate(chunks): await self._meshcore.send_text(message.sender_prefix, chunk) + if index + 1 < len(chunks) and self._settings.reply_chunk_delay_seconds > 0: + await asyncio.sleep(self._settings.reply_chunk_delay_seconds) def _build_nomad_prompt(self, prompt: str) -> str: if not self._settings.radio_prompt_enabled: diff --git a/openhop-plugin.json b/openhop-plugin.json index adbd38b..6e67bb7 100644 --- a/openhop-plugin.json +++ b/openhop-plugin.json @@ -2,7 +2,7 @@ "schema": 1, "id": "openhop.nomad", "name": "NOMAD Bridge", - "version": "0.1.2", + "version": "0.1.3", "description": "Bridges an openHop Repeater Companion identity to Project N.O.M.A.D.", "runtime": { "type": "python", @@ -25,7 +25,8 @@ "allowed_sender_prefixes": [], "busy_wait_seconds": 5, "max_reply_chunks": 4, - "max_chunk_bytes": 145, + "max_chunk_bytes": 80, + "reply_chunk_delay_seconds": 2.0, "max_prompt_bytes": 1000, "radio_prompt_enabled": true, "duplicate_ttl_seconds": 600, diff --git a/pyproject.toml b/pyproject.toml index dee3bb4..2e96bb5 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "openhop-nomad-plugin" -version = "0.1.2" +version = "0.1.3" description = "openHop service plugin bridging a MeshCore Companion identity to Project N.O.M.A.D." readme = "README.md" requires-python = ">=3.10" diff --git a/tests/test_config.py b/tests/test_config.py index 9ddb185..bbcb455 100644 --- a/tests/test_config.py +++ b/tests/test_config.py @@ -24,6 +24,7 @@ "NOMAD_BUSY_WAIT_SECONDS", "MAX_REPLY_CHUNKS", "MAX_CHUNK_BYTES", + "REPLY_CHUNK_DELAY_SECONDS", "MAX_PROMPT_BYTES", "RADIO_PROMPT_ENABLED", "RADIO_PROMPT_TEMPLATE", @@ -53,6 +54,7 @@ def test_settings_load_plugin_owned_config(monkeypatch: pytest.MonkeyPatch, tmp_ nomad_model="qwen-test", one_shot=True, max_reply_chunks=3, + reply_chunk_delay_seconds=2.5, max_pending_requests=3, max_requests_per_sender=2, max_requests_global=7, @@ -69,6 +71,7 @@ def test_settings_load_plugin_owned_config(monkeypatch: pytest.MonkeyPatch, tmp_ assert settings.nomad_model == "qwen-test" assert settings.one_shot is True assert settings.max_reply_chunks == 3 + assert settings.reply_chunk_delay_seconds == 2.5 assert settings.max_pending_requests == 3 assert settings.max_requests_per_sender == 2 assert settings.max_requests_global == 7 @@ -143,6 +146,7 @@ def test_persistent_mode_is_rejected_even_with_sender_allowlist( ("NOMAD_TIMEOUT_SECONDS", "inf"), ("RATE_LIMIT_WINDOW_SECONDS", "nan"), ("NOMAD_BUSY_WAIT_SECONDS", "-inf"), + ("REPLY_CHUNK_DELAY_SECONDS", "inf"), ], ) def test_nonfinite_numeric_settings_are_rejected( diff --git a/tests/test_message_handler.py b/tests/test_message_handler.py index 7e6d108..ea1bd89 100644 --- a/tests/test_message_handler.py +++ b/tests/test_message_handler.py @@ -74,6 +74,7 @@ def _settings() -> Settings: busy_wait_seconds=0.2, max_reply_chunks=4, max_chunk_bytes=145, + reply_chunk_delay_seconds=2.0, max_prompt_bytes=1000, radio_prompt_enabled=False, radio_prompt_template="{question}", @@ -82,6 +83,29 @@ def _settings() -> Settings: ) +@pytest.mark.asyncio +async def test_reply_chunks_are_spaced_except_after_last(monkeypatch) -> None: + class ChunkedNomad(SlowNomad): + async def ask_for_sender(self, sender_id: str, prompt: str) -> str: + _ = (sender_id, prompt) + return "x" * 95 + + delays: list[float] = [] + + async def record_sleep(delay: float) -> None: + delays.append(delay) + + monkeypatch.setattr("meshcore_nomad_bridge.main.asyncio.sleep", record_sleep) + settings = replace(_settings(), max_chunk_bytes=40, reply_chunk_delay_seconds=2.0) + mesh = FakeMeshCore() + service = BridgeService(settings=settings, meshcore=mesh, nomad=ChunkedNomad()) + + await service._handle_message(_msg("chunk this", ts=701)) + + assert len(mesh.sent) == 3 + assert delays == [2.0, 2.0] + + def _msg(text: str, ts: int = 1, sender: bytes = b"\x01\x02\x03\x04\x05\x06") -> IncomingMessage: return IncomingMessage( sender_prefix=sender, From 1bb9356386a6e9a6c31240993b9e1c871aa9888b Mon Sep 17 00:00:00 2001 From: yellowcooln <12516003+yellowcooln@users.noreply.github.com> Date: Fri, 11 Sep 2026 10:45:49 -0400 Subject: [PATCH 03/27] fix: complete bounded NOMAD transport and global reply pacing --- .github/workflows/ci.yml | 29 +++ README.md | 40 +++- meshcore_nomad_bridge/config.py | 6 +- meshcore_nomad_bridge/main.py | 39 +++- meshcore_nomad_bridge/meshcore_client.py | 10 +- meshcore_nomad_bridge/nomad_client.py | 238 +++++------------------ meshcore_nomad_bridge/text.py | 1 - pyproject.toml | 4 +- scripts/release_automation.py | 6 +- tests/test_config.py | 29 ++- tests/test_meshcore_client.py | 5 +- tests/test_message_handler.py | 127 +++++++++++- tests/test_nomad_client.py | 228 ++++++++++++++++++---- tests/test_plugin_package.py | 5 +- 14 files changed, 510 insertions(+), 257 deletions(-) create mode 100644 .github/workflows/ci.yml diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 0000000..e634404 --- /dev/null +++ b/.github/workflows/ci.yml @@ -0,0 +1,29 @@ +name: PR checks +on: + pull_request: + branches: [main, dev] +permissions: + contents: read +concurrency: + group: nomad-pr-${{ github.event.pull_request.number }} + cancel-in-progress: true +jobs: + test: + runs-on: ubuntu-latest + timeout-minutes: 10 + strategy: + fail-fast: false + matrix: + python: ['3.10', '3.11', '3.12', '3.13', '3.14'] + steps: + - uses: actions/checkout@11d5960a326750d5838078e36cf38b85af677262 # v4 + with: + persist-credentials: false + - uses: actions/setup-python@a26af69be951a213d495a4c3e4e4022e16d87065 # v5 + with: + python-version: ${{ matrix.python }} + - run: python -m pip install -e '.[dev]' build + - run: python -m pytest + # Explicit correctness rules: do not inherit developer-machine lint defaults. + - run: python -m ruff check --select E4,E7,E9,F . + - run: python -m build --wheel diff --git a/README.md b/README.md index 37bb815..781e380 100644 --- a/README.md +++ b/README.md @@ -113,8 +113,29 @@ limits permit one active request, two requests per sender per minute, and four r globally per minute. Rejected overload and authorization traffic is dropped without an RF reply. NOMAD HTTP responses are capped at 256 KiB, redirects are not followed, and the configured timeout is -an end-to-end request deadline implemented without non-cancellable worker threads. `NOMAD_URL` -must use an IP literal so DNS resolution cannot outlive that deadline. +an end-to-end deadline covering DNS, connection/TLS, headers and body reads. The transport +uses `aiohttp` with a dedicated `aiodns`/c-ares resolver, closing connections and cancelling +DNS queries on timeout or cancellation rather than leaving blocking Python `getaddrinfo` +executor jobs running. c-ares may use its own native event thread; this is not a guarantee +that the process creates no threads. It uses DNS/hosts resolution, not every OS NSS/mDNS +plugin: `.local` names require a DNS/hosts entry or an IP address where mDNS is unavailable. + +`NOMAD_URL` accepts HTTP(S) hostnames (including Docker names such as +`http://nomad_admin:8080`), IPv4 and bracketed IPv6. Docker names only resolve when the +plugin container shares the appropriate network. Credentials, paths (except `/`), query +strings and fragments are not accepted in the configured origin. TLS certificates are +verified; environment proxies and cookies are not used. Repeated standard response headers +are supported. Parser limits are 8,190 bytes per header field/status line and 128 headers +and trailers; aggregate response headers are checked against 64 KiB after parsing. +Compressed responses are rejected (the request asks for identity encoding), so compression +cannot bypass the 256 KiB body cap. + +All outgoing replies share `reply_chunk_delay_seconds` pacing, including short error +responses and concurrent answers. The interval starts when the previous send finishes; +there is no trailing sleep after the last packet. Companion retries remain internal to +`send_text` with their existing backoff. A rejected or unconfirmed chunk stops the rest of +that answer. `RESP_CODE_SENT` means Companion send acceptance, **not an over-air delivery +ACK**; retrying after an acceptance timeout can produce duplicates. This keeps environment variables available for development and existing standalone deployments. @@ -137,7 +158,7 @@ Important environment overrides include: - `NOMAD_BUSY_WAIT_SECONDS` - `MAX_REPLY_CHUNKS` - `MAX_CHUNK_BYTES` -- `REPLY_CHUNK_DELAY_SECONDS` (0–60 seconds between multi-packet reply chunks) +- `REPLY_CHUNK_DELAY_SECONDS` (0–60 seconds between outgoing reply sends, globally) - `MAX_PROMPT_BYTES` - `RADIO_PROMPT_ENABLED` - `RADIO_PROMPT_TEMPLATE` @@ -146,6 +167,18 @@ Important environment overrides include: `NOMAD_URL` and `NOMAD_MODEL` must be provided by `config.json` or environment variables. +## Upgrading from v0.1.2 + +Install with dependencies (including the new `aiohttp` and `aiodns` requirements). Python +3.10 and newer remain supported. The `openhop-core==1.1.1` pin is unchanged. +Before restarting, set `allowed_sender_prefixes` to the permitted 12-hex-character sender +prefixes and ensure `one_shot` is `true`. An empty allowlist now denies everyone and logs a +startup warning; there is no public/open mode. Old persistent-session maps are not deleted, +but persistent mode is rejected at startup. Existing hostname/IP URLs remain usable subject +to the origin rules above. Review the conservative request limits and global reply pacing; +these intentionally restrict traffic compared with earlier versions. No automatic config +migration or deployment is performed. + ## Plugin manifest `openhop-plugin.json` declares this as a Python service plugin with a dashboard UI for configuration editing: @@ -180,6 +213,7 @@ pip install -e .[dev] export NOMAD_URL=http://127.0.0.1:8080 export NOMAD_MODEL=qwen2.5:3b-instruct +export ALLOWED_SENDER_PREFIXES=001122334455 # replace with the permitted sender meshcore-nomad-bridge ``` diff --git a/meshcore_nomad_bridge/config.py b/meshcore_nomad_bridge/config.py index 637a247..f53a80b 100644 --- a/meshcore_nomad_bridge/config.py +++ b/meshcore_nomad_bridge/config.py @@ -2,7 +2,6 @@ from __future__ import annotations -import ipaddress import json import math import os @@ -262,7 +261,8 @@ def _validate(settings: Settings) -> None: parsed_url = urlsplit(settings.nomad_url) if parsed_url.scheme not in {"http", "https"} or not parsed_url.hostname: raise ValueError - ipaddress.ip_address(parsed_url.hostname) + if any(ord(char) <= 32 or ord(char) == 127 for char in settings.nomad_url): + raise ValueError port = parsed_url.port if ( parsed_url.username @@ -275,7 +275,7 @@ def _validate(settings: Settings) -> None: if port is not None and not 1 <= port <= 65535: raise ValueError except ValueError as exc: - raise ConfigError("NOMAD_URL must be a valid HTTP(S) URL with an IP address") from exc + raise ConfigError("NOMAD_URL must be a valid HTTP(S) origin with a hostname or IP address") from exc if not settings.nomad_session_map_path: raise ConfigError("NOMAD_SESSION_MAP_PATH must not be empty") diff --git a/meshcore_nomad_bridge/main.py b/meshcore_nomad_bridge/main.py index c75ecb2..2c9d8fc 100644 --- a/meshcore_nomad_bridge/main.py +++ b/meshcore_nomad_bridge/main.py @@ -86,6 +86,8 @@ def __init__(self, settings: Settings, meshcore: MeshCoreClient, nomad: NomadCli self._meshcore = meshcore self._nomad = nomad + self._send_lock = asyncio.Lock() + self._next_send_at = 0.0 self._stop_event = asyncio.Event() self._semaphore = asyncio.Semaphore(settings.max_concurrent_requests) self._duplicate_cache = DuplicateCache(settings.duplicate_ttl_seconds) @@ -98,6 +100,8 @@ def __init__(self, settings: Settings, meshcore: MeshCoreClient, nomad: NomadCli self._inflight: set[asyncio.Task[Any]] = set() async def run(self) -> None: + if not self._settings.allowed_sender_prefixes: + logger.warning("allowed_sender_prefixes is empty; all senders are denied") self._register_signals() await self._meshcore.run(self._dispatch_message, self._stop_event) await self._shutdown() @@ -163,19 +167,19 @@ async def _handle_message( try: await self._nomad.reset_session_for_sender(sender_id) except NomadUnavailable: - await self._meshcore.send_text(message.sender_prefix, NOMAD_UNAVAILABLE_MESSAGE) + await self._send_text(message.sender_prefix, NOMAD_UNAVAILABLE_MESSAGE) else: - await self._meshcore.send_text(message.sender_prefix, NOMAD_RESET_MESSAGE) + await self._send_text(message.sender_prefix, NOMAD_RESET_MESSAGE) return if len(prompt.encode("utf-8")) > self._settings.max_prompt_bytes: - await self._meshcore.send_text(message.sender_prefix, NOMAD_TOO_LONG_MESSAGE) + await self._send_text(message.sender_prefix, NOMAD_TOO_LONG_MESSAGE) return acquired = False if self._settings.busy_wait_seconds == 0: if self._semaphore.locked(): - await self._meshcore.send_text(message.sender_prefix, NOMAD_BUSY_MESSAGE) + await self._send_text(message.sender_prefix, NOMAD_BUSY_MESSAGE) return await self._semaphore.acquire() acquired = True @@ -186,7 +190,7 @@ async def _handle_message( ) acquired = True except asyncio.TimeoutError: - await self._meshcore.send_text(message.sender_prefix, NOMAD_BUSY_MESSAGE) + await self._send_text(message.sender_prefix, NOMAD_BUSY_MESSAGE) return try: @@ -194,7 +198,7 @@ async def _handle_message( logger.info("Sending request to NOMAD") answer = await self._nomad.ask_for_sender(sender_id, self._build_nomad_prompt(prompt)) except NomadUnavailable: - await self._meshcore.send_text(message.sender_prefix, NOMAD_UNAVAILABLE_MESSAGE) + await self._send_text(message.sender_prefix, NOMAD_UNAVAILABLE_MESSAGE) return finally: if acquired: @@ -214,9 +218,26 @@ async def _handle_message( logger.info("Sending %d MeshCore reply packets", len(chunks)) for index, chunk in enumerate(chunks): - await self._meshcore.send_text(message.sender_prefix, chunk) - if index + 1 < len(chunks) and self._settings.reply_chunk_delay_seconds > 0: - await asyncio.sleep(self._settings.reply_chunk_delay_seconds) + if not await self._send_text(message.sender_prefix, chunk): + logger.warning( + "Stopping reply to %s: chunk %d/%d not accepted or acceptance unknown", + sender_id, + index + 1, + len(chunks), + ) + break + + async def _send_text(self, recipient_prefix: bytes, text: str) -> bool: + # Serialize all replies, including short errors, across concurrent requests. + async with self._send_lock: + delay = self._next_send_at - time.monotonic() + if delay > 0: + await asyncio.sleep(delay) + try: + return await self._meshcore.send_text(recipient_prefix, text) + finally: + # Also pace a rejected/uncertain/cancelled send; it may have reached RF. + self._next_send_at = time.monotonic() + self._settings.reply_chunk_delay_seconds def _build_nomad_prompt(self, prompt: str) -> str: if not self._settings.radio_prompt_enabled: diff --git a/meshcore_nomad_bridge/meshcore_client.py b/meshcore_nomad_bridge/meshcore_client.py index 2501d98..a2faf32 100644 --- a/meshcore_nomad_bridge/meshcore_client.py +++ b/meshcore_nomad_bridge/meshcore_client.py @@ -6,8 +6,8 @@ import contextlib import logging import struct +from collections.abc import Awaitable, Callable from dataclasses import dataclass -from typing import Awaitable, Callable from openhop_core.companion.constants import ( CMD_APP_START, @@ -125,7 +125,7 @@ async def run( raise except Exception as exc: self._connected.clear() - logger.warning("MeshCore connection lost: %s", exc) + logger.warning("MeshCore connection lost: %s", exc, exc_info=True) if stop_event.is_set() or self._stop_requested.is_set(): break @@ -162,7 +162,7 @@ async def send_text(self, recipient_prefix: bytes, text: str) -> bool: code = frame[0] if code == RESP_CODE_SENT: logger.info( - "DM ACK received for recipient=%s on attempt %s/%s", + "Companion accepted DM for recipient=%s on attempt %s/%s (not RF delivery confirmation)", recipient_prefix[:6].hex(), attempt, max_retries, @@ -180,7 +180,7 @@ async def send_text(self, recipient_prefix: bytes, text: str) -> bool: except asyncio.TimeoutError: if attempt == max_retries: logger.warning( - "No DM ACK received for recipient=%s after %s attempts", + "No send acceptance received for recipient=%s after %s attempts; outcome unknown", recipient_prefix[:6].hex(), max_retries, ) @@ -188,7 +188,7 @@ async def send_text(self, recipient_prefix: bytes, text: str) -> bool: delay = 1.0 * (2 ** (attempt - 1)) logger.warning( - "No DM ACK received for recipient=%s, retrying in %.1fs (%s/%s)", + "No send acceptance received for recipient=%s, retrying in %.1fs (%s/%s)", recipient_prefix[:6].hex(), delay, attempt, diff --git a/meshcore_nomad_bridge/nomad_client.py b/meshcore_nomad_bridge/nomad_client.py index 40b4a17..2232581 100644 --- a/meshcore_nomad_bridge/nomad_client.py +++ b/meshcore_nomad_bridge/nomad_client.py @@ -3,24 +3,19 @@ from __future__ import annotations import asyncio -import ipaddress import json import logging -import socket -import ssl from collections.abc import Awaitable, Callable from pathlib import Path from time import monotonic from typing import Any -from urllib.parse import urlsplit, urlunsplit + +import aiohttp logger = logging.getLogger(__name__) MAX_HTTP_RESPONSE_BYTES = 262_144 MAX_HTTP_HEADER_BYTES = 65_536 -HTTP_TOKEN_BYTES = frozenset( - b"!#$%&'*+-.^_`|~0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz" -) MAX_SESSION_HISTORY_MESSAGES = 20 MAX_SESSION_HISTORY_BYTES = 16_384 @@ -371,189 +366,60 @@ async def _async_http_request( url: str, body: bytes | None, ) -> tuple[int, str]: - parsed = urlsplit(url) - if parsed.scheme not in {"http", "https"} or not parsed.hostname: - raise OSError("invalid_url") - try: - address = ipaddress.ip_address(parsed.hostname) - port = parsed.port or (443 if parsed.scheme == "https" else 80) - except ValueError as exc: - raise OSError("nomad_url_requires_ip_address") from exc - - family = socket.AF_INET6 if address.version == 6 else socket.AF_INET - sock = socket.socket(family, socket.SOCK_STREAM) - sock.setblocking(False) - writer: asyncio.StreamWriter | None = None - path = urlunsplit(("", "", parsed.path or "/", parsed.query, "")) - try: - await asyncio.get_running_loop().sock_connect(sock, (str(address), port)) - ssl_context = ssl.create_default_context() if parsed.scheme == "https" else None - reader, writer = await asyncio.open_connection( - sock=sock, - ssl=ssl_context, - server_hostname=parsed.hostname if ssl_context else None, - ) - sock = None - request = _build_http_request(method=method, parsed=parsed, path=path, body=body) - writer.write(request) - await writer.drain() - return await _read_async_http_response(reader, method) - finally: - if writer is not None: - writer.transport.abort() - elif sock is not None: - sock.close() - - -def _build_http_request(*, method: str, parsed: Any, path: str, body: bytes | None) -> bytes: if method not in {"GET", "POST"}: raise OSError("unsupported_http_method") - if not path or any(ord(char) < 33 or ord(char) > 126 for char in path): + if any(ord(char) <= 32 or ord(char) == 127 for char in url): raise OSError("invalid_url") - default_port = 443 if parsed.scheme == "https" else 80 - host = f"[{parsed.hostname}]" if ":" in parsed.hostname else parsed.hostname - if parsed.port and parsed.port != default_port: - host = f"{host}:{parsed.port}" - headers = [ - f"{method} {path} HTTP/1.1", - f"Host: {host}", - "Accept: application/json", - "Connection: close", - ] - if body is not None: - headers.extend(("Content-Type: application/json", f"Content-Length: {len(body)}")) + # Dedicated c-ares resolver: no asyncio getaddrinfo executor jobs survive + # cancellation. Passing options avoids aiohttp's shared resolver lifetime. + resolver = aiohttp.AsyncResolver(tries=1) try: - return ("\r\n".join(headers) + "\r\n\r\n").encode("ascii") + (body or b"") - except UnicodeEncodeError as exc: - raise OSError("invalid_url") from exc - - -async def _read_async_http_response( - reader: asyncio.StreamReader, method: str -) -> tuple[int, str]: - status_line = await _read_http_line(reader) - if not status_line.endswith(b"\r\n"): - raise OSError("invalid_http_response") - parts = status_line[:-2].split(b" ", 2) - try: - version, raw_status = parts[0], parts[1] - if version not in {b"HTTP/1.0", b"HTTP/1.1"}: - raise ValueError - if len(raw_status) != 3 or not raw_status.isdigit(): - raise ValueError - if len(parts) == 3 and any( - (byte < 32 and byte != 9) or byte == 127 for byte in parts[2] + async with ( + aiohttp.TCPConnector(resolver=resolver, use_dns_cache=False) as connector, + aiohttp.ClientSession( + connector=connector, + timeout=aiohttp.ClientTimeout(total=None), + trust_env=False, + cookie_jar=aiohttp.DummyCookieJar(), + auto_decompress=False, + headers={"Accept": "application/json", "Accept-Encoding": "identity"}, + max_line_size=8190, + max_field_size=8190, + max_headers=128, + read_bufsize=16_384, + ) as session, + session.request( + method, + url, + data=body, + allow_redirects=False, + headers={"Content-Type": "application/json"} if body is not None else None, + ) as response, ): - raise ValueError - status = int(raw_status) - if not 100 <= status <= 599: - raise ValueError - except (IndexError, ValueError) as exc: - raise OSError("invalid_http_response") from exc - - header_bytes = len(status_line) - headers: dict[str, str] = {} - while True: - line = await _read_http_line(reader) - header_bytes += len(line) - if not line or header_bytes > MAX_HTTP_HEADER_BYTES: - raise OSError("invalid_http_response") - if line == b"\r\n": - break - key, value = _parse_header_line(line) - if key in headers: - raise OSError("invalid_http_response") - headers[key] = value - - transfer_encoding = headers.get("transfer-encoding") - if transfer_encoding is not None and "content-length" in headers: - raise OSError("invalid_http_response") - if method == "HEAD" or status in {204, 304} or 100 <= status < 200: - raw = b"" - elif transfer_encoding is not None: - if transfer_encoding.lower() != "chunked": - raise OSError("invalid_http_response") - raw = await _read_chunked_body(reader) - elif "content-length" in headers: - raw_length = headers["content-length"] - if not raw_length.isascii() or not raw_length.isdigit(): - raise OSError("invalid_http_response") - length = int(raw_length) - if length > MAX_HTTP_RESPONSE_BYTES: - raise OSError("response_too_large") - try: - raw = await reader.readexactly(length) - except asyncio.IncompleteReadError as exc: - raise OSError("invalid_http_response") from exc - else: - try: - raw = await reader.readexactly(MAX_HTTP_RESPONSE_BYTES + 1) - except asyncio.IncompleteReadError as exc: - raw = exc.partial - if len(raw) > MAX_HTTP_RESPONSE_BYTES: - raise OSError("response_too_large") - return status, raw.decode("utf-8", errors="replace") - - -async def _read_chunked_body(reader: asyncio.StreamReader) -> bytes: - body = bytearray() - framing_bytes = 0 - while True: - size_line = await _read_http_line(reader) - framing_bytes += len(size_line) - if framing_bytes > MAX_HTTP_HEADER_BYTES: - raise OSError("invalid_http_response") - if not size_line.endswith(b"\r\n") or b";" in size_line: - raise OSError("invalid_http_response") - raw_size = size_line[:-2] - if not raw_size or any(byte not in b"0123456789abcdefABCDEF" for byte in raw_size): - raise OSError("invalid_http_response") - try: - size = int(raw_size, 16) - except ValueError as exc: - raise OSError("invalid_http_response") from exc - if size == 0: - while True: - trailer = await _read_http_line(reader) - framing_bytes += len(trailer) - if framing_bytes > MAX_HTTP_HEADER_BYTES: - raise OSError("invalid_http_response") - if trailer == b"\r\n": - return bytes(body) - if not trailer: - raise OSError("invalid_http_response") - key, _ = _parse_header_line(trailer) - if key in {"content-length", "transfer-encoding"}: - raise OSError("invalid_http_response") - if size < 0 or len(body) + size > MAX_HTTP_RESPONSE_BYTES: - raise OSError("response_too_large") - try: - body.extend(await reader.readexactly(size)) - if await reader.readexactly(2) != b"\r\n": + if any( + ord(char) < 32 and char != "\t" or ord(char) == 127 + for char in response.reason or "" + ): raise OSError("invalid_http_response") - except asyncio.IncompleteReadError as exc: - raise OSError("invalid_http_response") from exc - - -async def _read_http_line(reader: asyncio.StreamReader) -> bytes: - try: - line = await reader.readline() - except ValueError as exc: - raise OSError("invalid_http_response") from exc - if len(line) > MAX_HTTP_HEADER_BYTES: - raise OSError("invalid_http_response") - return line - - -def _parse_header_line(line: bytes) -> tuple[str, str]: - if not line.endswith(b"\r\n"): - raise OSError("invalid_http_response") - try: - raw_name, raw_value = line[:-2].split(b":", 1) - except ValueError as exc: + if response.headers.get("Transfer-Encoding", "chunked").lower() != "chunked": + raise OSError("invalid_http_response") + # Check aggregate size too; parser limits bound each field/count. + if sum(len(k) + len(v) + 4 for k, v in response.raw_headers) > MAX_HTTP_HEADER_BYTES: + raise OSError("invalid_http_response") + if response.headers.get("Content-Encoding", "identity").lower() != "identity": + raise OSError("unsupported_content_encoding") + if ( + response.content_length is not None + and response.content_length > MAX_HTTP_RESPONSE_BYTES + ): + raise OSError("response_too_large") + data = bytearray() + async for chunk in response.content.iter_chunked(16_384): + if len(data) + len(chunk) > MAX_HTTP_RESPONSE_BYTES: + raise OSError("response_too_large") + data.extend(chunk) + return response.status, data.decode("utf-8", errors="replace") + except aiohttp.ClientError as exc: raise OSError("invalid_http_response") from exc - if not raw_name or any(byte not in HTTP_TOKEN_BYTES for byte in raw_name): - raise OSError("invalid_http_response") - if any((byte < 32 and byte != 9) or byte == 127 for byte in raw_value): - raise OSError("invalid_http_response") - return raw_name.decode("ascii").lower(), raw_value.decode("iso-8859-1").strip() + finally: + await resolver.close() diff --git a/meshcore_nomad_bridge/text.py b/meshcore_nomad_bridge/text.py index 3ffb4ab..d06df2c 100644 --- a/meshcore_nomad_bridge/text.py +++ b/meshcore_nomad_bridge/text.py @@ -4,7 +4,6 @@ import re - _CODE_FENCE_RE = re.compile(r"```.*?```", re.DOTALL) _HEADING_RE = re.compile(r"^#{1,6}\s*", re.MULTILINE) _TABLE_LINE_RE = re.compile(r"^\s*\|.*\|\s*$", re.MULTILINE) diff --git a/pyproject.toml b/pyproject.toml index 4cd9cad..4c1fde2 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -12,7 +12,9 @@ authors = [ { name = "OpenHop Community" } ] dependencies = [ - "openhop-core==1.1.1" + "openhop-core==1.1.1", + "aiohttp>=3.14.3,<4", + "aiodns>=3.2,<5" ] [project.optional-dependencies] diff --git a/scripts/release_automation.py b/scripts/release_automation.py index e06ed4d..36616d6 100644 --- a/scripts/release_automation.py +++ b/scripts/release_automation.py @@ -17,11 +17,15 @@ from pathlib import Path import re import subprocess -import tomllib import urllib.error import urllib.request import zipfile +try: + import tomllib +except ModuleNotFoundError: # Python 3.10 test/development environments + import tomli as tomllib + REPOSITORY = "openhop-dev/openhop-nomad-plugin" CATALOGUE = "openhop-dev/openhop-plugin-catalogue" PLUGIN = "openhop.nomad" diff --git a/tests/test_config.py b/tests/test_config.py index bbcb455..8f658e6 100644 --- a/tests/test_config.py +++ b/tests/test_config.py @@ -79,7 +79,9 @@ def test_settings_load_plugin_owned_config(monkeypatch: pytest.MonkeyPatch, tmp_ assert settings.allowed_sender_prefixes == ("010203040506", "aabbccddeeff") -def test_environment_overrides_plugin_config(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> None: +def test_environment_overrides_plugin_config( + monkeypatch: pytest.MonkeyPatch, tmp_path: Path +) -> None: _clean_env(monkeypatch) data_dir = tmp_path / "plugin-data" _write_config( @@ -100,7 +102,9 @@ def test_environment_overrides_plugin_config(monkeypatch: pytest.MonkeyPatch, tm assert settings.meshcore_port == 6001 -def test_session_map_defaults_to_plugin_data(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> None: +def test_session_map_defaults_to_plugin_data( + monkeypatch: pytest.MonkeyPatch, tmp_path: Path +) -> None: _clean_env(monkeypatch) data_dir = tmp_path / "plugin-data" _write_config(data_dir, nomad_url="http://10.5.30.7", nomad_model="test-model") @@ -164,16 +168,13 @@ def test_nonfinite_numeric_settings_are_rejected( @pytest.mark.parametrize( "url", [ - "http://nomad.local:8080", "http://10.5.30.7:bad", "http://10.5.30.7:70000", "http://10.5.30.7/base", "http://10.5.30.7/?query=1", ], ) -def test_nomad_url_requires_valid_ip_literal( - monkeypatch: pytest.MonkeyPatch, url: str -) -> None: +def test_nomad_url_rejects_invalid_origin(monkeypatch: pytest.MonkeyPatch, url: str) -> None: _clean_env(monkeypatch) monkeypatch.setenv("NOMAD_URL", url) monkeypatch.setenv("NOMAD_MODEL", "test-model") @@ -229,3 +230,19 @@ def test_invalid_plugin_config_is_reported(monkeypatch: pytest.MonkeyPatch, tmp_ with pytest.raises(ConfigError, match="config.json"): Settings.from_env() + + +@pytest.mark.parametrize( + "url", + [ + "http://nomad.local:8080", + "http://nomad_admin:8080", + "https://localhost", + "http://[::1]:8080", + ], +) +def test_nomad_url_accepts_hostnames_and_ip_literals(monkeypatch, url): + _clean_env(monkeypatch) + monkeypatch.setenv("NOMAD_URL", url) + monkeypatch.setenv("NOMAD_MODEL", "test") + assert Settings.from_env().nomad_url == url diff --git a/tests/test_meshcore_client.py b/tests/test_meshcore_client.py index 80b201a..9cb379e 100644 --- a/tests/test_meshcore_client.py +++ b/tests/test_meshcore_client.py @@ -11,7 +11,8 @@ @pytest.mark.asyncio -async def test_send_text_retries_after_timeout_and_logs_ack() -> None: +async def test_send_text_retries_after_timeout_and_logs_acceptance(caplog) -> None: + caplog.set_level("INFO") client = MeshCoreClient(host="127.0.0.1", port=5001) call_count = 0 @@ -28,6 +29,8 @@ async def fake_send_command_expect(*args, **kwargs): # type: ignore[no-untyped- assert result is True assert call_count == 2 + assert "accepted" in caplog.text + assert "DM ACK" not in caplog.text @pytest.mark.asyncio diff --git a/tests/test_message_handler.py b/tests/test_message_handler.py index ea1bd89..479a2b1 100644 --- a/tests/test_message_handler.py +++ b/tests/test_message_handler.py @@ -1,5 +1,6 @@ import asyncio from dataclasses import replace +from itertools import pairwise import pytest @@ -103,7 +104,7 @@ async def record_sleep(delay: float) -> None: await service._handle_message(_msg("chunk this", ts=701)) assert len(mesh.sent) == 3 - assert delays == [2.0, 2.0] + assert delays == pytest.approx([2.0, 2.0], abs=0.01) def _msg(text: str, ts: int = 1, sender: bytes = b"\x01\x02\x03\x04\x05\x06") -> IncomingMessage: @@ -301,12 +302,8 @@ async def test_dispatch_silently_enforces_global_rate_limit() -> None: nomad = SlowNomad(delay=0) service = BridgeService(settings=settings, meshcore=mesh, nomad=nomad) - await service._dispatch_message( - _msg("first", ts=401, sender=b"\x01\x02\x03\x04\x05\x06") - ) - await service._dispatch_message( - _msg("second", ts=402, sender=b"\x11\x12\x13\x14\x15\x16") - ) + await service._dispatch_message(_msg("first", ts=401, sender=b"\x01\x02\x03\x04\x05\x06")) + await service._dispatch_message(_msg("second", ts=402, sender=b"\x11\x12\x13\x14\x15\x16")) await asyncio.gather(*list(service._inflight)) assert nomad.calls == 1 @@ -387,3 +384,119 @@ async def test_reset_command_starts_new_session_in_persistent_mode() -> None: assert nomad.calls == 0 assert nomad.resets == 1 assert mesh.sent[-1][1] == NOMAD_RESET_MESSAGE + + +@pytest.mark.asyncio +async def test_failed_chunk_stops_remaining_reply(caplog): + class RejectMesh(FakeMeshCore): + async def send_text(self, recipient_prefix, text): + self.sent.append((recipient_prefix, text)) + return False + + mesh = RejectMesh() + service = BridgeService( + replace(_settings(), max_chunk_bytes=40, reply_chunk_delay_seconds=0), mesh, SlowNomad(0) + ) + service._nomad.ask_for_sender = _long_answer + await service._handle_message(_msg("question")) + assert len(mesh.sent) == 1 + assert "Stopping reply" in caplog.text + + +async def _long_answer(*args): + return "x" * 95 + + +@pytest.mark.asyncio +async def test_empty_allowlist_warns_at_startup(caplog): + mesh = FakeMeshCore() + + async def run(*args): + return None + + mesh.run = run + service = BridgeService(replace(_settings(), allowed_sender_prefixes=()), mesh, SlowNomad(0)) + service._register_signals = lambda: None + await service.run() + assert "allowed_sender_prefixes is empty; all senders are denied" in caplog.text + + +@pytest.mark.asyncio +async def test_outbound_pacing_is_global_including_error_replies(): + times = [] + + class TimedMesh(FakeMeshCore): + async def send_text(self, recipient_prefix, text): + times.append(asyncio.get_running_loop().time()) + return await super().send_text(recipient_prefix, text) + + service = BridgeService( + replace(_settings(), reply_chunk_delay_seconds=0.04, max_prompt_bytes=5), + TimedMesh(), + SlowNomad(0), + ) + await asyncio.gather(*(service._handle_message(_msg("too long", ts=i)) for i in range(3))) + assert len(times) == 3 + assert all(b - a >= 0.035 for a, b in pairwise(times)) + + +@pytest.mark.asyncio +async def test_cancelled_pacing_waiter_does_not_block_next_reply(): + mesh = FakeMeshCore() + service = BridgeService( + replace(_settings(), reply_chunk_delay_seconds=0.05, max_prompt_bytes=5), mesh, SlowNomad(0) + ) + await service._handle_message(_msg("too long", ts=1)) + task = asyncio.create_task(service._handle_message(_msg("too long", ts=2))) + await asyncio.sleep(0.005) + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + await asyncio.wait_for(service._handle_message(_msg("too long", ts=3)), 0.2) + assert len(mesh.sent) == 2 + + +@pytest.mark.asyncio +async def test_concurrent_chunked_answers_share_pacing(): + times = [] + + class TimedMesh(FakeMeshCore): + async def send_text(self, recipient_prefix, text): + times.append(asyncio.get_running_loop().time()) + return await super().send_text(recipient_prefix, text) + + mesh = TimedMesh() + service = BridgeService( + replace(_settings(), max_chunk_bytes=40, reply_chunk_delay_seconds=0.02), mesh, SlowNomad(0) + ) + service._nomad.ask_for_sender = _long_answer + await asyncio.gather( + service._handle_message(_msg("a", ts=1)), service._handle_message(_msg("b", ts=2)) + ) + assert len(mesh.sent) == 6 + assert all(b - a >= 0.018 for a, b in pairwise(times)) + + +@pytest.mark.asyncio +async def test_cancelling_active_send_releases_pacing_lock(): + entered = asyncio.Event() + + class BlockingMesh(FakeMeshCore): + async def send_text(self, recipient_prefix, text): + await super().send_text(recipient_prefix, text) + if len(self.sent) == 1: + entered.set() + await asyncio.Event().wait() + return True + + mesh = BlockingMesh() + service = BridgeService( + replace(_settings(), reply_chunk_delay_seconds=0.02), mesh, SlowNomad(0) + ) + task = asyncio.create_task(service._handle_message(_msg("first", ts=1))) + await entered.wait() + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + await asyncio.wait_for(service._handle_message(_msg("next", ts=2)), 0.2) + assert len(mesh.sent) == 2 diff --git a/tests/test_nomad_client.py b/tests/test_nomad_client.py index 8ef6b1a..2da8e0c 100644 --- a/tests/test_nomad_client.py +++ b/tests/test_nomad_client.py @@ -47,12 +47,14 @@ def log_message(self, format: str, *args: object) -> None: return None -def test_http_request_rejects_control_characters_in_target() -> None: - parsed = nomad_client.urlsplit("http://127.0.0.1") - +@pytest.mark.asyncio +async def test_http_request_rejects_control_characters_in_target() -> None: with pytest.raises(OSError, match="invalid_url"): - nomad_client._build_http_request( - method="GET", parsed=parsed, path="/ok\r\nX-Injected: yes", body=None + await nomad_client._default_http_request( + method="GET", + url="http://127.0.0.1/ok\r\nX-Injected: yes", + payload=None, + timeout_seconds=2, ) @@ -162,8 +164,7 @@ async def test_oversized_http_header_is_a_controlled_network_error() -> None: @pytest.mark.asyncio async def test_chunked_response_is_decoded() -> None: status, body = await _request_raw_response( - b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n" - b"3\r\nabc\r\n3\r\ndef\r\n0\r\n\r\n" + b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n3\r\nabc\r\n3\r\ndef\r\n0\r\n\r\n" ) assert status == 200 @@ -182,18 +183,9 @@ async def test_chunked_response_is_decoded() -> None: b"2\r\n{}\r\n0\r\n\r\n" ), b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nContent-Length: 2\r\n\r\n{}", - ( - b"HTTP/1.1 200 OK\r\nTransfer-Encoding: notchunked\r\n\r\n" - b"2\r\n{}\r\n0\r\n\r\n" - ), - ( - b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n" - b"2\r\n{}\r\n0\r\ngarbage\r\n\r\n" - ), - ( - b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n" - b"2;bad=\x00\r\n{}\r\n0\r\n\r\n" - ), + (b"HTTP/1.1 200 OK\r\nTransfer-Encoding: notchunked\r\n\r\n2\r\n{}\r\n0\r\n\r\n"), + (b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n2\r\n{}\r\n0\r\ngarbage\r\n\r\n"), + (b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n2;bad=\x00\r\n{}\r\n0\r\n\r\n"), ], ) async def test_malformed_http_framing_is_rejected(raw_response: bytes) -> None: @@ -203,12 +195,10 @@ async def test_malformed_http_framing_is_rejected(raw_response: bytes) -> None: @pytest.mark.asyncio async def test_chunked_response_rejects_missing_terminal_line() -> None: - reader = asyncio.StreamReader() - reader.feed_data(b"1\r\nx\r\n0\r\n") - reader.feed_eof() - with pytest.raises(OSError, match="invalid_http_response"): - await nomad_client._read_chunked_body(reader) + await _request_raw_response( + b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n1\r\nx\r\n0\r\n" + ) @pytest.mark.asyncio @@ -231,17 +221,13 @@ async def test_default_http_request_converts_malformed_http_to_oserror() -> None @pytest.mark.asyncio async def test_default_http_request_rejects_oversized_success_body() -> None: with pytest.raises(OSError, match="response_too_large"): - await _request_raw_response( - b"HTTP/1.1 200 OK\r\nContent-Length: 262145\r\n\r\n" - ) + await _request_raw_response(b"HTTP/1.1 200 OK\r\nContent-Length: 262145\r\n\r\n") @pytest.mark.asyncio async def test_default_http_request_rejects_oversized_error_body() -> None: with pytest.raises(OSError, match="response_too_large"): - await _request_raw_response( - b"HTTP/1.1 500 Error\r\nContent-Length: 262145\r\n\r\n" - ) + await _request_raw_response(b"HTTP/1.1 500 Error\r\nContent-Length: 262145\r\n\r\n") @pytest.mark.asyncio @@ -250,7 +236,10 @@ async def fake_post(*, url: str, payload: dict[str, object], timeout_seconds: fl assert url == "http://nomad.local/api/ollama/chat" assert payload["stream"] is False assert timeout_seconds == 5 - return 200, '{"message":{"role":"assistant","content":"Hello from NOMAD"},"done":true,"model":"test"}' + return ( + 200, + '{"message":{"role":"assistant","content":"Hello from NOMAD"},"done":true,"model":"test"}', + ) client = NomadClient( base_url="http://nomad.local", @@ -325,7 +314,10 @@ async def fake_request( if method == "POST" and url == "http://nomad.local/api/chat/sessions": return 201, '{"id":"42","title":"MeshCore sender-a","model":"test-model"}' if method == "GET" and url == "http://nomad.local/api/chat/sessions/42": - return 200, '{"id":"42","messages":[{"role":"user","content":"old q"},{"role":"assistant","content":"old a"}]}' + return ( + 200, + '{"id":"42","messages":[{"role":"user","content":"old q"},{"role":"assistant","content":"old a"}]}', + ) if method == "POST" and url == "http://nomad.local/api/ollama/chat": assert payload is not None assert payload["sessionId"] == 42 @@ -352,7 +344,11 @@ async def fake_request( assert result == "new a" assert session_map.read_text(encoding="utf-8") == '{"sender-a": 42}' - assert calls[0] == ("POST", "http://nomad.local/api/chat/sessions", {"title": "MeshCore sender-a", "model": "test-model"}) + assert calls[0] == ( + "POST", + "http://nomad.local/api/chat/sessions", + {"title": "MeshCore sender-a", "model": "test-model"}, + ) @pytest.mark.asyncio @@ -472,3 +468,169 @@ def test_persistent_history_is_limited_by_utf8_bytes() -> None: history = nomad_client._limit_session_history(messages) assert history == [messages[-1]] + + +@pytest.mark.asyncio +async def test_hostname_http_request_with_repeated_standard_headers(): + async def handle(reader, writer): + await reader.readuntil(b"\r\n\r\n") + writer.write( + b"HTTP/1.1 200 OK\r\nSet-Cookie: a=1\r\nSet-Cookie: b=2\r\nContent-Length: 2\r\n\r\n{}" + ) + await writer.drain() + writer.close() + await writer.wait_closed() + + server = await asyncio.start_server(handle, "127.0.0.1", 0) + async with server: + result = await nomad_client._default_http_request( + method="GET", + url=f"http://localhost:{server.sockets[0].getsockname()[1]}/", + payload=None, + timeout_seconds=2, + ) + assert result == (200, "{}") + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "raw", + [ + b"HTTP/1.1 200 OK\r\n" + b"X: a\r\n" * 129 + b"\r\n", + b"HTTP/1.1 200 OK\r\n" + b"X: " + b"a" * 8191 + b"\r\n\r\n", + b"HTTP/1.1 200 OK\r\n" + (b"X: " + b"a" * 8000 + b"\r\n") * 9 + b"\r\n", + ], +) +async def test_header_limits(raw): + with pytest.raises(OSError, match="invalid_http_response"): + await _request_raw_response(raw) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("framing", [b"Connection: close\r\n", b"Transfer-Encoding: chunked\r\n"]) +async def test_streaming_body_limit_without_content_length(framing): + body = b"x" * 262145 + if b"chunked" in framing: + body = b"40001\r\n" + body + b"\r\n0\r\n\r\n" + with pytest.raises(OSError, match="response_too_large"): + await _request_raw_response(b"HTTP/1.1 200 OK\r\n" + framing + b"\r\n" + body) + + +@pytest.mark.asyncio +async def test_compressed_response_is_rejected_without_decompression(): + with pytest.raises(OSError, match="unsupported_content_encoding"): + await _request_raw_response( + b"HTTP/1.1 200 OK\r\nContent-Encoding: gzip\r\nContent-Length: 0\r\n\r\n" + ) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("cancel", [False, True]) +async def test_inflight_http_timeout_or_cancellation_closes_socket(cancel): + entered, closed = asyncio.Event(), asyncio.Event() + + async def handle(reader, writer): + try: + await reader.readuntil(b"\r\n\r\n") + entered.set() + await reader.read() + closed.set() + finally: + writer.close() + await writer.wait_closed() + + server = await asyncio.start_server(handle, "127.0.0.1", 0) + async with server: + task = asyncio.create_task( + nomad_client._default_http_request( + method="GET", + url=f"http://localhost:{server.sockets[0].getsockname()[1]}/", + payload=None, + timeout_seconds=0.15 if not cancel else 5, + ) + ) + await asyncio.wait_for(entered.wait(), 2) + if cancel: + task.cancel() + with pytest.raises(asyncio.CancelledError if cancel else TimeoutError): + await task + await asyncio.wait_for(closed.wait(), 1) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("respond", [True, False]) +async def test_cares_dns_resolution_and_deadline_without_executor(monkeypatch, respond): + import functools + import struct + + queried = asyncio.Event() + + class DNS(asyncio.DatagramProtocol): + def connection_made(self, transport): + self.transport = transport + + def datagram_received(self, data, address): + queried.set() + if not respond: + return + end = 12 + while data[end]: + end += data[end] + 1 + end += 1 + qtype = struct.unpack("!H", data[end : end + 2])[0] + question = data[12 : end + 4] + answer = ( + (b"\xc0\x0c" + struct.pack("!HHIH", 1, 1, 10, 4) + b"\x7f\x00\x00\x01") + if qtype == 1 + else b"" + ) + self.transport.sendto( + data[:2] + struct.pack("!HHHHH", 0x8180, 1, bool(answer), 0, 0) + question + answer, + address, + ) + + loop = asyncio.get_running_loop() + transport, _ = await loop.create_datagram_endpoint(DNS, local_addr=("127.0.0.1", 0)) + real_resolver = nomad_client.aiohttp.AsyncResolver + monkeypatch.setattr( + nomad_client.aiohttp, + "AsyncResolver", + functools.partial( + real_resolver, + nameservers=["127.0.0.1"], + udp_port=transport.get_extra_info("sockname")[1], + ), + ) + + def forbid_executor(*args, **kwargs): + raise AssertionError("DNS must not use blocking getaddrinfo workers") + + monkeypatch.setattr(loop, "run_in_executor", forbid_executor) + + async def handle(reader, writer): + await reader.readuntil(b"\r\n\r\n") + writer.write(b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\n\r\n{}") + await writer.drain() + writer.close() + await writer.wait_closed() + + server = await asyncio.start_server(handle, "127.0.0.1", 0) + before = asyncio.all_tasks() + try: + async with server: + request = nomad_client._default_http_request( + method="GET", + url=f"http://nomad_admin.test:{server.sockets[0].getsockname()[1]}/", + payload=None, + timeout_seconds=1 if respond else 0.1, + ) + if respond: + assert await request == (200, "{}") + else: + with pytest.raises(TimeoutError): + await request + assert queried.is_set() + await asyncio.sleep(0) + assert not (asyncio.all_tasks() - before) + finally: + transport.close() diff --git a/tests/test_plugin_package.py b/tests/test_plugin_package.py index 8e1d59b..4bc2c49 100644 --- a/tests/test_plugin_package.py +++ b/tests/test_plugin_package.py @@ -1,7 +1,10 @@ import json from pathlib import Path -import tomllib +try: + import tomllib +except ModuleNotFoundError: # Python 3.10 + import tomli as tomllib ROOT = Path(__file__).resolve().parents[1] From f392980999fb7dfa07792b784355783fe1e4b4be Mon Sep 17 00:00:00 2001 From: yellowcooln <12516003+yellowcooln@users.noreply.github.com> Date: Fri, 11 Sep 2026 10:54:24 -0400 Subject: [PATCH 04/27] fix: default NOMAD Docker installs to container DNS endpoint --- README.md | 34 ++++++++++++++++++++++++++++++---- config.default.json | 2 +- openhop-plugin.json | 2 +- tests/test_config.py | 41 +++++++++++++++++++++++++++++++++++++++++ ui/app.js | 2 +- ui/index.html | 4 ++-- 6 files changed, 76 insertions(+), 9 deletions(-) diff --git a/README.md b/README.md index 781e380..a2fb2f5 100644 --- a/README.md +++ b/README.md @@ -56,6 +56,26 @@ for configuration compatibility. When `OPENHOP_PLUGIN_DATA` is not set, existing standalone behaviour is preserved and the session map defaults to `./data/nomad_sessions.json`. +## NOMAD Docker networking + +New plugin installations default to `http://nomad_admin:8080`. This is the NOMAD +container's DNS name and internal HTTP port, not the Repeater container's loopback. +The openHop Repeater container (which runs the plugin) and `nomad_admin` **must share +a user-defined Docker network** with the `nomad_admin` name/alias available there. +Docker's default bridge network does not provide this name-resolution contract. + +Persist both containers' network attachments in their Compose/deployment configuration +so they survive container recreation and upgrades. For separate Compose projects, use +the same externally managed user-defined network in both configurations. A one-off +`docker network connect` is not a durable replacement for that configuration. + +For standalone or remote deployments where `nomad_admin` is not resolvable, explicitly +set `nomad_url` in `config.json`, or override it with `NOMAD_URL`, to a reachable origin +such as `http://192.0.2.10:8080` or `https://nomad.example.org`. Use +`http://127.0.0.1:8080` only when NOMAD actually runs in the plugin's own network +namespace (for example standalone processes on the same host). +The Companion default remains `127.0.0.1:5001`; this change affects only NOMAD HTTP. + ## config.json With `OPENHOP_PLUGIN_DATA` set, the plugin reads `$OPENHOP_PLUGIN_DATA/config.json` if it exists. @@ -64,7 +84,7 @@ Minimum configuration: ```json { - "nomad_url": "http://192.168.0.170:8080", + "nomad_url": "http://nomad_admin:8080", "nomad_model": "qwen2.5:3b-instruct", "allowed_sender_prefixes": ["001122334455"] } @@ -76,7 +96,7 @@ Typical configuration: { "meshcore_host": "127.0.0.1", "meshcore_port": 5001, - "nomad_url": "http://192.168.0.170:8080", + "nomad_url": "http://nomad_admin:8080", "nomad_model": "qwen2.5:3b-instruct", "nomad_collection": null, "nomad_timeout_seconds": 120, @@ -177,7 +197,12 @@ startup warning; there is no public/open mode. Old persistent-session maps are n but persistent mode is rejected at startup. Existing hostname/IP URLs remain usable subject to the origin rules above. Review the conservative request limits and global reply pacing; these intentionally restrict traffic compared with earlier versions. No automatic config -migration or deployment is performed. +migration or deployment is performed. The new Docker URL is an installation default, +not an upgrade migration: preserve existing `config.json` files and environment +overrides; do not replace them with `config.default.json` during upgrades. Existing +explicit URLs, including loopback, remain unchanged. If an older loopback setting +is wrong for your Docker setup, change it deliberately after configuring the shared +network described above. ## Plugin manifest @@ -204,7 +229,8 @@ The plugin remains a lightweight service package and does not add a custom permi ## Standalone development -The plugin remains runnable without the plugin manager: +The plugin remains runnable without the plugin manager. This example assumes NOMAD +runs on the same host/network namespace; otherwise set a reachable remote URL: ```bash python3 -m venv .venv diff --git a/config.default.json b/config.default.json index 17fe17f..3a0a721 100644 --- a/config.default.json +++ b/config.default.json @@ -1,7 +1,7 @@ { "meshcore_host": "127.0.0.1", "meshcore_port": 5001, - "nomad_url": "http://127.0.0.1:8080", + "nomad_url": "http://nomad_admin:8080", "nomad_model": "qwen2.5:3b-instruct", "nomad_collection": null, "nomad_timeout_seconds": 120, diff --git a/openhop-plugin.json b/openhop-plugin.json index 69a067a..5239a2a 100644 --- a/openhop-plugin.json +++ b/openhop-plugin.json @@ -16,7 +16,7 @@ "defaults": { "meshcore_host": "127.0.0.1", "meshcore_port": 5001, - "nomad_url": "http://127.0.0.1:8080", + "nomad_url": "http://nomad_admin:8080", "nomad_model": "qwen2.5:3b-instruct", "nomad_collection": null, "nomad_timeout_seconds": 120, diff --git a/tests/test_config.py b/tests/test_config.py index 8f658e6..4ef7076 100644 --- a/tests/test_config.py +++ b/tests/test_config.py @@ -43,6 +43,47 @@ def _write_config(path: Path, **values: object) -> None: (path / "config.json").write_text(json.dumps(values), encoding="utf-8") +@pytest.mark.parametrize("source", ["config.default.json", "openhop-plugin.json"]) +def test_packaged_defaults_use_nomad_docker_endpoint(monkeypatch, tmp_path, source): + _clean_env(monkeypatch) + root = Path(__file__).resolve().parents[1] + defaults = json.loads((root / source).read_text(encoding="utf-8")) + if source == "openhop-plugin.json": + defaults = defaults["config"]["defaults"] + _write_config(tmp_path, **defaults) + monkeypatch.setenv("OPENHOP_PLUGIN_DATA", str(tmp_path)) + + settings = Settings.from_env() + + assert settings.nomad_url == "http://nomad_admin:8080" + assert (settings.meshcore_host, settings.meshcore_port) == ("127.0.0.1", 5001) + + +@pytest.mark.parametrize("source", ["ui/app.js", "ui/index.html"]) +def test_ui_endpoint_defaults_match_docker_install(source): + root = Path(__file__).resolve().parents[1] + content = (root / source).read_text(encoding="utf-8") + assert "http://nomad_admin:8080" in content + assert "http://127.0.0.1:8080" not in content + + +@pytest.mark.parametrize("url", ["http://127.0.0.1:8080", "https://remote.example:8443"]) +@pytest.mark.parametrize("override", [None, "http://other.example:8080"]) +def test_existing_endpoint_is_preserved_without_rewriting_config( + monkeypatch, tmp_path, url, override +): + _clean_env(monkeypatch) + _write_config(tmp_path, nomad_url=url, nomad_model="existing-model") + config_path = tmp_path / "config.json" + before = config_path.read_bytes() + monkeypatch.setenv("OPENHOP_PLUGIN_DATA", str(tmp_path)) + if override: + monkeypatch.setenv("NOMAD_URL", override) + + assert Settings.from_env().nomad_url == (override or url) + assert config_path.read_bytes() == before + + def test_settings_load_plugin_owned_config(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> None: _clean_env(monkeypatch) data_dir = tmp_path / "plugin-data" diff --git a/ui/app.js b/ui/app.js index cef1b94..6487f10 100644 --- a/ui/app.js +++ b/ui/app.js @@ -7,7 +7,7 @@ const defaults = { meshcore_host: "127.0.0.1", meshcore_port: 5001, - nomad_url: "http://127.0.0.1:8080", + nomad_url: "http://nomad_admin:8080", nomad_model: "qwen2.5:3b-instruct", nomad_collection: null, nomad_timeout_seconds: 120, diff --git a/ui/index.html b/ui/index.html index caef055..77edbd8 100644 --- a/ui/index.html +++ b/ui/index.html @@ -29,7 +29,7 @@