From 0bc8d4e2d3e5966674bb4d2cf0952ddc4f4bbf7a Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 21:17:47 -0700 Subject: [PATCH] refactor(ntfy): fold fatal-status handling into one helper, dedupe headers/env seeding, compact layout --- plugins/platforms/ntfy/adapter.py | 185 +++++++++--------------------- 1 file changed, 52 insertions(+), 133 deletions(-) diff --git a/plugins/platforms/ntfy/adapter.py b/plugins/platforms/ntfy/adapter.py index ff8f5fab2c..a499434647 100644 --- a/plugins/platforms/ntfy/adapter.py +++ b/plugins/platforms/ntfy/adapter.py @@ -1,32 +1,14 @@ """ntfy platform adapter (Hermes plugin). -Subscribes to a topic on ntfy.sh or a self-hosted ntfy server via HTTP -streaming (``/json`` with ``poll=false``) and publishes replies via HTTP POST. -No external SDK — only httpx. +Subscribes to a topic on ntfy.sh or a self-hosted server via HTTP streaming +(``/json`` with ``poll=false``) and publishes replies via HTTP POST (httpx only). -Configuration in config.yaml:: - - platforms: - ntfy: - enabled: true - extra: - server: "https://ntfy.sh" # or self-hosted URL - topic: "hermes-in" # subscribe topic (incoming) - publish_topic: "hermes-out" # optional — defaults to topic - token: "..." # optional Bearer / Basic auth token - markdown: true # optional — enable markdown (default: false) - -Environment variables (read at adapter construct time; env wins over ``extra``): - - NTFY_TOPIC Topic to subscribe to (required) - NTFY_SERVER_URL Server URL (default: https://ntfy.sh) - NTFY_TOKEN Bearer token or 'user:pass' for Basic auth - NTFY_PUBLISH_TOPIC Reply topic (defaults to NTFY_TOPIC) - NTFY_MARKDOWN "true"/"1"/"yes" enables X-Markdown header - NTFY_ALLOWED_USERS Allowlist (on ntfy these are topic names) - NTFY_ALLOW_ALL_USERS Allow any topic — dev only - NTFY_HOME_CHANNEL Default topic for cron / notification delivery - NTFY_HOME_CHANNEL_NAME Human label for the home channel +config.yaml ``platforms.ntfy.extra``: ``server`` (default https://ntfy.sh), +``topic`` (subscribe, required), ``publish_topic`` (defaults to topic), ``token`` +(Bearer or ``user:pass`` Basic), ``markdown`` (default false). +Env (read at construct time; env wins over ``extra``): NTFY_TOPIC, NTFY_SERVER_URL, +NTFY_TOKEN, NTFY_PUBLISH_TOPIC, NTFY_MARKDOWN ("true"/"1"/"yes"), NTFY_ALLOWED_USERS +(topic names), NTFY_ALLOW_ALL_USERS (dev only), NTFY_HOME_CHANNEL, NTFY_HOME_CHANNEL_NAME. Identity model: ntfy has no authenticated user identity; ``title`` is publisher-controlled and NOT used for authorization. Each topic is a single @@ -50,16 +32,9 @@ except ImportError: httpx = None # type: ignore[assignment] from gateway.config import Platform, PlatformConfig -from gateway.platforms.base import ( - BasePlatformAdapter, - MessageEvent, - MessageType, - SendResult, -) - +from gateway.platforms.base import BasePlatformAdapter, MessageEvent, MessageType, SendResult from gateway.platforms._shared import get_scoped_secret as _get_scoped_secret - logger = logging.getLogger(__name__) @@ -75,6 +50,7 @@ RECONNECT_BACKOFF = [2, 5, 10, 30, 60] STREAM_TIMEOUT_SECONDS = 90 # ntfy keepalive default is 55s; give margin _ECHO_TAG = "hermes-agent" # tag added to outgoing messages for echo-loop prevention _MARKDOWN_TRUTHY = ("1", "true", "yes") +_BASE_HEADERS = {"Content-Type": "text/plain; charset=utf-8", "X-Tags": _ECHO_TAG} def _build_auth_header(token: str) -> Dict[str, str]: @@ -83,15 +59,12 @@ def _build_auth_header(token: str) -> Dict[str, str]: Tokens are whitespace-stripped (pasted tokens often carry newlines that would malform the header). ``user:pass`` → Basic, anything else → Bearer. """ - if not token: - return {} - token = token.strip() + token = (token or "").strip() if not token: return {} if ":" in token: import base64 - encoded = base64.b64encode(token.encode()).decode() - return {"Authorization": f"Basic {encoded}"} + return {"Authorization": f"Basic {base64.b64encode(token.encode()).decode()}"} return {"Authorization": f"Bearer {token}"} @@ -115,9 +88,7 @@ def _response_message_id(resp) -> str: def check_requirements() -> bool: """Installable and minimally configured (reads NTFY_TOPIC directly — no full config load).""" - if not HTTPX_AVAILABLE: - return False - return bool(_get_scoped_secret("NTFY_TOPIC", "").strip()) + return HTTPX_AVAILABLE and bool(_get_scoped_secret("NTFY_TOPIC", "").strip()) def validate_config(config) -> bool: @@ -139,17 +110,13 @@ class NtfyAdapter(BasePlatformAdapter): def __init__(self, config: PlatformConfig): super().__init__(config=config, platform=Platform("ntfy")) - extra = config.extra or {} - self._server: str = ( - extra.get("server") or _get_scoped_secret("NTFY_SERVER_URL", DEFAULT_SERVER) - ).rstrip("/") + self._server: str = (extra.get("server") or _get_scoped_secret("NTFY_SERVER_URL", DEFAULT_SERVER)).rstrip("/") self._topic: str = extra.get("topic") or _get_scoped_secret("NTFY_TOPIC", "") self._publish_topic: str = ( extra.get("publish_topic") or _get_scoped_secret("NTFY_PUBLISH_TOPIC", "") or self._topic ) self._token: str = extra.get("token") or _get_scoped_secret("NTFY_TOKEN", "") - self._stream_task: Optional[asyncio.Task] = None self._http_client: Optional["httpx.AsyncClient"] = None self._seen_messages: Dict[str, float] = {} # msg_id -> timestamp (dedup) @@ -164,7 +131,6 @@ class NtfyAdapter(BasePlatformAdapter): if not self._topic: logger.warning("[%s] NTFY_TOPIC not configured", self.name) return False - try: self._http_client = httpx.AsyncClient(timeout=None) self._stream_task = asyncio.create_task(self._run_stream()) @@ -182,7 +148,6 @@ class NtfyAdapter(BasePlatformAdapter): stream_start: float = 0.0 url = f"{self._server}/{self._topic}/json" headers = self._auth_headers() - while self._running: try: logger.debug("[%s] Opening stream to %s", self.name, url) @@ -197,7 +162,6 @@ class NtfyAdapter(BasePlatformAdapter): if not self._running: return logger.warning("[%s] Stream error: %s", self.name, e) - if not self._running: return # Reset backoff if stream stayed alive for at least 60s @@ -208,37 +172,34 @@ class NtfyAdapter(BasePlatformAdapter): await asyncio.sleep(delay) backoff_idx += 1 + def _fatal_status(self, status_code: int) -> None: + """401/404 are unrecoverable: log, set the fatal state and raise ``_FatalStreamError``.""" + if status_code == 401: + logger.error( + "[%s] Authentication failed (401) — stopping reconnect loop. Check NTFY_TOKEN.", self.name, + ) + self._set_fatal_error( + "ntfy_unauthorized", "ntfy server rejected auth (401). Check NTFY_TOKEN.", retryable=False, + ) + raise _FatalStreamError("401 Unauthorized") + if status_code == 404: + logger.error("[%s] Topic not found (404): %s — stopping reconnect loop.", self.name, self._topic) + self._set_fatal_error( + "ntfy_topic_not_found", + f"ntfy topic '{self._topic}' returned 404. Check NTFY_TOPIC.", + retryable=False, + ) + raise _FatalStreamError("404 Not Found") + async def _consume_stream(self, url: str, headers: Dict[str, str]) -> None: """Open an HTTP streaming connection and dispatch events.""" # poll=false keeps a persistent streaming connection alive with keepalive events async with self._http_client.stream( - "GET", - url, - headers=headers, - params={"poll": "false"}, + "GET", url, headers=headers, params={"poll": "false"}, timeout=httpx.Timeout(connect=15.0, read=STREAM_TIMEOUT_SECONDS, write=15.0, pool=15.0), ) as response: - if response.status_code == 401: - logger.error( - "[%s] Authentication failed (401) — stopping reconnect loop. Check NTFY_TOKEN.", - self.name, - ) - self._set_fatal_error( - "ntfy_unauthorized", "ntfy server rejected auth (401). Check NTFY_TOKEN.", retryable=False, - ) - raise _FatalStreamError("401 Unauthorized") - if response.status_code == 404: - logger.error( - "[%s] Topic not found (404): %s — stopping reconnect loop.", self.name, self._topic, - ) - self._set_fatal_error( - "ntfy_topic_not_found", - f"ntfy topic '{self._topic}' returned 404. Check NTFY_TOPIC.", - retryable=False, - ) - raise _FatalStreamError("404 Not Found") + self._fatal_status(response.status_code) response.raise_for_status() - async for line in response.aiter_lines(): if not self._running: return @@ -256,7 +217,6 @@ class NtfyAdapter(BasePlatformAdapter): """Disconnect from ntfy.""" self._running = False self._mark_disconnected() - if self._stream_task: self._stream_task.cancel() try: @@ -264,11 +224,9 @@ class NtfyAdapter(BasePlatformAdapter): except asyncio.CancelledError: pass self._stream_task = None - if self._http_client: await self._http_client.aclose() self._http_client = None - self._seen_messages.clear() logger.info("[%s] Disconnected", self.name) @@ -287,30 +245,20 @@ class NtfyAdapter(BasePlatformAdapter): if not text: logger.debug("[%s] Empty message body, skipping", self.name) return - # No native user identity on ntfy: the publisher-controlled title must # NOT drive authorization, so user_id is fixed to the topic name. topic = event.get("topic") or self._topic - source = self.build_source( - chat_id=topic, chat_name=topic, chat_type="dm", user_id=topic, user_name=topic, - ) - + source = self.build_source(chat_id=topic, chat_name=topic, chat_type="dm", user_id=topic, user_name=topic) unix_ts = event.get("time") try: timestamp = ( - datetime.fromtimestamp(int(unix_ts), tz=timezone.utc) - if unix_ts else datetime.now(tz=timezone.utc) + datetime.fromtimestamp(int(unix_ts), tz=timezone.utc) if unix_ts else datetime.now(tz=timezone.utc) ) except (ValueError, OSError, TypeError): timestamp = datetime.now(tz=timezone.utc) - message_event = MessageEvent( - text=text, - message_type=MessageType.TEXT, - source=source, - message_id=msg_id, - raw_message=event, - timestamp=timestamp, + text=text, message_type=MessageType.TEXT, source=source, message_id=msg_id, + raw_message=event, timestamp=timestamp, ) logger.debug("[%s] Message on topic %s: %s", self.name, topic, text[:80]) await self.handle_message(message_event) @@ -329,38 +277,25 @@ class NtfyAdapter(BasePlatformAdapter): # -- Outbound messaging ------------------------------------------------- async def send( - self, - chat_id: str, - content: str, - reply_to: Optional[str] = None, - metadata: Optional[Dict[str, Any]] = None, + self, chat_id: str, content: str, reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, ) -> SendResult: """Publish a message to the configured publish topic.""" metadata = metadata or {} publish_topic = metadata.get("publish_topic") or self._publish_topic or chat_id if not self._http_client: return SendResult(success=False, error="HTTP client not initialized") - url = f"{self._server}/{publish_topic}" - headers = { - **self._auth_headers(), - "Content-Type": "text/plain; charset=utf-8", - "X-Tags": _ECHO_TAG, - } + headers = {**self._auth_headers(), **_BASE_HEADERS} if (self.config.extra or {}).get("markdown", False): headers["X-Markdown"] = "true" - if len(content) > self.MAX_MESSAGE_LENGTH: logger.warning( "[%s] Message truncated from %d to %d chars (ntfy limit)", self.name, len(content), self.MAX_MESSAGE_LENGTH, ) body = content[:self.MAX_MESSAGE_LENGTH] - try: - resp = await self._http_client.post( - url, content=body.encode("utf-8"), headers=headers, timeout=15.0, - ) + resp = await self._http_client.post(url, content=body.encode("utf-8"), headers=headers, timeout=15.0) if resp.status_code < 300: return SendResult(success=True, message_id=_response_message_id(resp)) body_text = resp.text @@ -393,16 +328,11 @@ def _env_enablement() -> dict | None: topic = _get_scoped_secret("NTFY_TOPIC", "").strip() if not topic: return None - seed: dict = { - "topic": topic, - "server": _get_scoped_secret("NTFY_SERVER_URL", DEFAULT_SERVER).rstrip("/"), - } - publish_topic = _get_scoped_secret("NTFY_PUBLISH_TOPIC", "").strip() - if publish_topic: - seed["publish_topic"] = publish_topic - token = _get_scoped_secret("NTFY_TOKEN", "").strip() - if token: - seed["token"] = token + seed: dict = {"topic": topic, "server": _get_scoped_secret("NTFY_SERVER_URL", DEFAULT_SERVER).rstrip("/")} + for key, env in (("publish_topic", "NTFY_PUBLISH_TOPIC"), ("token", "NTFY_TOKEN")): + value = _get_scoped_secret(env, "").strip() + if value: + seed[key] = value markdown = _get_scoped_secret("NTFY_MARKDOWN", "").strip().lower() if markdown: seed["markdown"] = markdown in _MARKDOWN_TRUTHY @@ -413,13 +343,8 @@ def _env_enablement() -> dict | None: async def _standalone_send( - pconfig, - chat_id: str, - message: str, - *, - thread_id: Optional[str] = None, - media_files: Optional[List[str]] = None, - force_document: bool = False, + pconfig, chat_id: str, message: str, *, + thread_id: Optional[str] = None, media_files: Optional[List[str]] = None, force_document: bool = False, ) -> Dict[str, Any]: """Out-of-process publish for cron / send_message_tool when no gateway adapter is live. @@ -429,7 +354,6 @@ async def _standalone_send( """ if not HTTPX_AVAILABLE: return {"error": "ntfy standalone send: httpx not installed"} - extra = getattr(pconfig, "extra", {}) or {} server = (extra.get("server") or _get_scoped_secret("NTFY_SERVER_URL", DEFAULT_SERVER)).rstrip("/") publish_topic = ( @@ -441,23 +365,18 @@ async def _standalone_send( ) if not publish_topic: return {"error": "ntfy standalone send: NTFY_TOPIC not configured"} - token = extra.get("token") or _get_scoped_secret("NTFY_TOKEN", "") markdown_env = _get_scoped_secret("NTFY_MARKDOWN", "").strip().lower() - headers = {"Content-Type": "text/plain; charset=utf-8", "X-Tags": _ECHO_TAG, **_build_auth_header(token)} + headers = {**_BASE_HEADERS, **_build_auth_header(token)} if bool(extra.get("markdown")) or markdown_env in _MARKDOWN_TRUTHY: headers["X-Markdown"] = "true" - body = _truncate_body(message, context="ntfy standalone") - url = f"{server}/{publish_topic}" try: async with httpx.AsyncClient(timeout=15.0) as client: - resp = await client.post(url, content=body, headers=headers) + resp = await client.post(f"{server}/{publish_topic}", content=body, headers=headers) if resp.status_code >= 300: return {"error": f"ntfy HTTP {resp.status_code}: {resp.text[:200]}"} - return { - "success": True, "platform": "ntfy", "chat_id": publish_topic, "message_id": _response_message_id(resp), - } + return {"success": True, "platform": "ntfy", "chat_id": publish_topic, "message_id": _response_message_id(resp)} except Exception as e: return {"error": f"ntfy standalone send failed: {e}"}