341 lines
12 KiB
Python
341 lines
12 KiB
Python
"""Push backends: FCM (primary) + ntfy (fallback).
|
|
|
|
``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`` (``fcm`` default, ``ntfy`` fallback).
|
|
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"
|
|
FCM_TOKEN_URL = "https://oauth2.googleapis.com/token"
|
|
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: 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: python-hardcoded-secrets
|
|
token_field = "fcm_token"
|
|
|
|
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: python-hardcoded-secrets
|
|
token_field = "ntfy_topic"
|
|
|
|
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``; fcm default)."""
|
|
if (name or "").strip().lower() == "ntfy":
|
|
return NtfyBackend(
|
|
topic=ntfy_topic, server_url=ntfy_server_url, auth_token=ntfy_auth_token
|
|
)
|
|
return FcmBackend(service_account=fcm_service_account, server_key=fcm_server_key)
|