Files
iris_x_hermes/gateway-plugin/push.py
T
ARIA 746d809d48
CI / Gateway plugin tests (push) Successful in 5m3s
CI / Kotlin tests (android host + desktop) (push) Successful in 7m6s
Default push backend to ntfy; FCM opt-in with privacy warning (issue #10)
- IRIS_PUSH_BACKEND now defaults to ntfy (keeps push metadata on your own
  infrastructure); FCM is opt-in via IRIS_PUSH_BACKEND=fcm
- build_push_backend(): ntfy for empty/unknown names, FCM only on explicit 'fcm'
- gateway setup: warn when FCM is chosen (metadata routed via Google's servers)
- README: privacy note + dedicated push section; new docs/playstore-listing.md
  with the FCM/ntfy privacy note for the Play Store listing
- docs: 00/02/03/08/12/16 + setup.md updated to ntfy-default wording
- tests: default-backend assertion updated (86/86 pass)
2026-08-24 19:04:40 +02:00

348 lines
12 KiB
Python

"""Push backends: ntfy (default) + FCM (optional).
``PushBackend`` interface with two implementations:
- ``FcmBackend``: FCM HTTP v1 via ``httpx`` + a Firebase service account
(``IRIS_FCM_SERVICE_ACCOUNT``), or a legacy server key
(``IRIS_FCM_SERVER_KEY``).
- ``NtfyBackend``: publishes to ``NTFY_TOPIC`` on ``NTFY_SERVER_URL``
(default ``https://ntfy.sh``) via ``httpx``; the app's listener
subscribes to the topic.
Selected by ``IRIS_PUSH_BACKEND`` (``ntfy`` default, ``fcm`` optional).
Fired when a frame has no live subscriber; the data payload drives a silent
sync on the device (docs/08-push.md).
Zero new dependencies: ``httpx`` is a hermes core dep. ``google.auth`` is NOT
installed, so the FCM service-account OAuth2 access token is minted directly
with PyJWT + cryptography (both core deps).
Payloads carry no secrets and only a short preview (lock-screen privacy);
full content is fetched via ``sync`` over the authenticated WS.
"""
import json
import logging
import threading
import time
from pathlib import Path
from typing import Any
from urllib.parse import quote
import httpx
logger = logging.getLogger(__name__)
FCM_SCOPE = "https://www.googleapis.com/auth/firebase.messaging"
# Not a secret: the well-known Google OAuth2 token endpoint.
# pi-lens-ignore: S105
FCM_TOKEN_URL = "https://oauth2.googleapis.com/token" # noqa: S105
FCM_V1_SEND_URL = "https://fcm.googleapis.com/v1/projects/{project_id}/messages:send"
FCM_LEGACY_SEND_URL = "https://fcm.googleapis.com/fcm/send"
# Refresh the cached access token this long before its expiry.
_TOKEN_REFRESH_MARGIN_S = 600.0
_DEFAULT_NTFY_SERVER = "https://ntfy.sh"
_NTFY_BODY_LIMIT = 4096
_HTTP_TIMEOUT_S = 15.0
# HTTP status boundaries: 200 == success; >= 300 == redirect/error range.
_HTTP_OK = 200
_HTTP_ERROR_MIN = 300
_NTFY_PRIORITY = {"high": "5", "normal": "3", "low": "1"}
class PushBackend:
"""Interface: wake one device (FCM token or ntfy topic)."""
name: str = "push"
# DeviceRegistry column that carries this backend's target token.
# Not a secret: a DB column name (string literal), not a credential.
# pi-lens-ignore: S105, python-hardcoded-secrets
token_field: str = ""
def configured(self) -> bool:
"""True when the backend has credentials to send with."""
raise NotImplementedError
async def send(
self,
*,
device_id: str,
chat_id: str,
title: str,
body: str,
data: dict[str, Any],
token: str,
priority: str = "normal",
data_only: bool = False,
) -> bool:
"""Deliver one push to *token*. Returns True on success.
``data`` is the silent-sync payload (``chat_id``, ``kind``,
``cursor``, optional ``thread_id``/``message_id``) -- string values
only on the wire.
"""
raise NotImplementedError
class FcmBackend(PushBackend):
"""FCM HTTP v1 (service account) or legacy ``/fcm/send`` (server key)."""
name = "fcm"
# Not a secret: a DB column name (string literal), not a credential.
# pi-lens-ignore: S105
token_field = "fcm_token" # noqa: S105
def __init__(
self,
service_account: str | None = None,
server_key: str | None = None,
):
self._sa_path = (service_account or "").strip() or None
self._server_key = (server_key or "").strip() or None
self._sa: dict[str, Any] | None = None
self._sa_failed = False
self._access_token: str | None = None
self._token_expiry = 0.0
self._lock = threading.Lock()
def configured(self) -> bool:
if self._server_key:
return True
return bool(self._sa_path and Path(self._sa_path).is_file())
def _load_sa(self) -> dict[str, Any] | None:
if self._sa is not None:
return self._sa
if not self._sa_path or self._sa_failed:
return None
try:
with open(self._sa_path, encoding="utf-8") as f:
sa = json.load(f)
if isinstance(sa, dict) and sa.get("client_email") and sa.get("private_key"):
self._sa = sa
return sa
except Exception:
logger.warning("iris: FCM service account unreadable: %s", self._sa_path)
self._sa_failed = True
return None
async def _authorization(self, client: httpx.AsyncClient) -> str | None:
"""Bearer token: the legacy server key, or a cached service-account
OAuth2 access token (JWT-bearer grant, minted with PyJWT)."""
if self._server_key:
return self._server_key
sa = self._load_sa()
if sa is None:
return None
now = time.time()
with self._lock:
if self._access_token and now < self._token_expiry - _TOKEN_REFRESH_MARGIN_S:
return self._access_token
import jwt # PyJWT (core dep)
claims = {
"iss": sa["client_email"],
"scope": FCM_SCOPE,
"aud": FCM_TOKEN_URL,
"iat": int(now),
"exp": int(now) + 3600,
}
headers = {"kid": sa["private_key_id"]} if sa.get("private_key_id") else None
try:
assertion = jwt.encode(
claims, sa["private_key"], algorithm="RS256", headers=headers
)
except Exception:
logger.warning("iris: FCM JWT mint failed", exc_info=True)
return None
try:
resp = await client.post(
FCM_TOKEN_URL,
data={
"grant_type": "urn:ietf:params:oauth:grant-type:jwt-bearer",
"assertion": assertion,
},
timeout=_HTTP_TIMEOUT_S,
)
except Exception:
logger.warning("iris: FCM token exchange failed", exc_info=True)
return None
if resp.status_code != _HTTP_OK:
logger.warning(
"iris: FCM token exchange HTTP %s: %s",
resp.status_code, resp.text[:200],
)
return None
try:
data = resp.json()
except (json.JSONDecodeError, ValueError):
return None
token = data.get("access_token")
if not isinstance(token, str) or not token:
return None
with self._lock:
self._access_token = token
try:
self._token_expiry = now + float(data.get("expires_in", 3600))
except (TypeError, ValueError):
self._token_expiry = now + 3600.0
return token
async def send(
self,
*,
device_id: str,
chat_id: str,
title: str,
body: str,
data: dict[str, Any],
token: str,
priority: str = "normal",
data_only: bool = False,
) -> bool:
if not token:
return False
data = {str(k): str(v) for k, v in (data or {}).items()}
notification = None if data_only else {"title": title or "Iris", "body": body or ""}
async with httpx.AsyncClient(timeout=_HTTP_TIMEOUT_S) as client:
if self._server_key:
payload: dict[str, Any] = {"to": token}
if notification:
payload["notification"] = notification
if data:
payload["data"] = data
auth = self._server_key
url = FCM_LEGACY_SEND_URL
else:
sa = self._load_sa()
project_id = (sa or {}).get("project_id")
if not project_id:
return False
message: dict[str, Any] = {"token": token}
if notification:
message["notification"] = notification
if data:
message["data"] = data
message["android"] = {
"priority": "high" if priority == "high" else "normal"
}
payload = {"message": message}
auth = await self._authorization(client)
if auth is None:
return False
url = FCM_V1_SEND_URL.format(project_id=project_id)
try:
resp = await client.post(
url,
json=payload,
headers={"Authorization": f"Bearer {auth}"},
)
except Exception:
logger.warning("iris: FCM send failed (network)", exc_info=True)
return False
if resp.status_code >= _HTTP_ERROR_MIN:
# 404 NOT_FOUND = stale/invalid registration token.
logger.warning(
"iris: FCM send HTTP %s: %s", resp.status_code, resp.text[:200]
)
return False
return True
class NtfyBackend(PushBackend):
"""ntfy publish (self-host friendly; zero Firebase).
The structured payload rides in an ``X-Data`` header (JSON) the app's
listener parses; the message body is the short preview. A private topic
+ ``NTFY_AUTH_TOKEN`` provides the trust boundary (docs/08 §8.7).
"""
name = "ntfy"
# Not a secret: a DB column name (string literal), not a credential.
# pi-lens-ignore: S105
token_field = "ntfy_topic" # noqa: S105
def __init__(
self,
topic: str | None = None,
server_url: str | None = None,
auth_token: str | None = None,
):
self._topic = (topic or "").strip() or None
self._server = (
(server_url or _DEFAULT_NTFY_SERVER).strip().rstrip("/")
or _DEFAULT_NTFY_SERVER
)
self._auth_token = (auth_token or "").strip() or None
@property
def server_url(self) -> str:
"""The ntfy server this backend publishes to (for app discovery)."""
return self._server
def configured(self) -> bool:
return bool(self._topic)
async def send(
self,
*,
device_id: str,
chat_id: str,
title: str,
body: str,
data: dict[str, Any],
token: str,
priority: str = "normal",
data_only: bool = False,
) -> bool:
topic = (token or self._topic or "").strip()
if not topic:
return False
headers = {
"Content-Type": "text/plain; charset=utf-8",
"X-Title": (title or "Iris")[:512],
"X-Priority": _NTFY_PRIORITY.get(priority, "3"),
"X-Tag": "bell",
"X-Data": json.dumps(data or {}, separators=(",", ":")),
}
if self._auth_token:
headers["Authorization"] = f"Bearer {self._auth_token}"
text = (body or "")[:_NTFY_BODY_LIMIT]
url = f"{self._server}/{quote(topic, safe='')}"
try:
async with httpx.AsyncClient(timeout=_HTTP_TIMEOUT_S) as client:
resp = await client.post(
url, content=text.encode("utf-8"), headers=headers
)
except Exception:
logger.warning("iris: ntfy publish failed (network)", exc_info=True)
return False
if resp.status_code >= _HTTP_ERROR_MIN:
logger.warning(
"iris: ntfy publish HTTP %s: %s", resp.status_code, resp.text[:200]
)
return False
return True
def build_push_backend(
name: str | None,
*,
fcm_service_account: str | None = None,
fcm_server_key: str | None = None,
ntfy_topic: str | None = None,
ntfy_server_url: str | None = None,
ntfy_auth_token: str | None = None,
) -> PushBackend:
"""Select the backend by name (``IRIS_PUSH_BACKEND``; ntfy default).
ntfy is the default: it keeps push metadata on your own infrastructure.
FCM is opt-in (``IRIS_PUSH_BACKEND=fcm``) — its metadata (title, device
token) is routed through Google's servers.
"""
if (name or "").strip().lower() == "fcm":
return FcmBackend(service_account=fcm_service_account, server_key=fcm_server_key)
return NtfyBackend(
topic=ntfy_topic, server_url=ntfy_server_url, auth_token=ntfy_auth_token
)