Files
iris_x_hermes/gateway-plugin/push.py
T
ARIA 7faaf2aa1c
CI / Gateway plugin tests (push) Successful in 5m5s
CI / Kotlin tests (android host + desktop) (push) Successful in 6m50s
Per-device tokens with revocation (issue #11)
Auth previously used the shared IRIS_TOKEN as the security principal:
a leaked token meant access to all devices, and a compromised device
could not be isolated.

Gateway:
- pairing.py: devices.token column (in-place migration) + revoked
  denylist table; issue_token (idempotent, 64 hex), token_for,
  reissue_token, revoke/unrevoke/is_revoked/list_revoked. The token
  never leaks into device dicts (push fan-out / listings).
- http_server.py: auth accepts the shared token (bootstrap/legacy) OR
  the device's own token (both constant-time); a revoked device_id is
  rejected with 401 before either comparison. On SSE open (pairing)
  the per-device token is minted and returned in hello.ack.
- protocol.py: hello_ack(..., device_token).
- adapter.py: setup flow (hermes gateway setup -> Iris) now offers
  'Remove a paired device?' on an existing setup: numbered select
  menu (last option = exit the removal loop), confirmation, back to
  the menu for further removals.
- tools/iris_devices.py: operator CLI (list / revoke / unrevoke /
  reissue), stdlib only.

App:
- SecureStore.deviceToken (Android: EncryptedSharedPreferences;
  Desktop: second keyring slot iris-device-token / device_token.enc).
- HelloAckPayload.deviceToken; GatewayClient stores it on hello and
  presents it instead of the shared token from then on (live provider
  in HttpGateway); savePairing/clear wipe it for re-pairing.

Docs: 09 §9.3 stretch -> implemented (revocation semantics, both
control surfaces), 04 hello.ack example, frames.schema.json, M7 row 13.

Tests: 8 new Python tests (issuance, acceptance, revocation,
isolation, unrevoke, registry unit x2, setup-flow menu) - 94/94 pass;
2 new Kotlin wire tests - green. Live-verified against a running
gateway (hello.ack token matches devices.db; revoke -> 401 even with
shared token; unrevoke -> 200; setup TUI both paths).
2026-08-24 19:37:44 +02:00

336 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)