refactor(ntfy): fold fatal-status handling into one helper, dedupe headers/env seeding, compact layout

This commit is contained in:
Teknium
2026-09-02 21:17:47 -07:00
parent 113f04616b
commit 0bc8d4e2d3
+52 -133
View File
@@ -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}"}