M5: push (FCM + ntfy) + offline catch-up + background notifications

Server (gateway-plugin):
- push.py: FCM (HTTP v1 service-account / legacy key) + ntfy backends; NtfyBackend.server_url for app discovery
- adapter.py: push on offline broadcast + high-priority push when live; fcm.register; server_caps.push_ntfy_server; outbox now ALWAYS appends so a reconnecting app catches up on live-delivered frames (fixes empty chat after notification tap / activity recreation)
- protocol.py / ws_server.py / outbox.py / plugin.yaml: M5 frames + env vars

App (Kotlin CMP):
- Protocol.kt: notification / fcm.register / sync frames + push caps
- GatewayClient.kt: hello carries push creds; auto-sync on reconnect
- IrisController.kt: banners, deep-link, notifyMessageIfBackgrounded (system notification on a regular reply when backgrounded)
- Android: PlatformPush, AppBridge, IrisNotifications, NtfyListenerService, IrisFirebaseMessagingService, AndroidPush, AndroidSecureStore
- Desktop: DesktopPush, DesktopSecureStore
- build files + manifest (permissions, services, deep-links)

Tests: 35-test tests/gateway/test_android.py suite passes (incl. new regression test_live_delivered_frame_still_parked_for_sync). E2E verified on device: push fire, reconnect sync, background notification, and message replay after ChatStore reset.
This commit is contained in:
ARIA committed 2026-08-19 19:20:55 +02:00
1 parent 913ee91024
commit 2ecfe1c05c
28 files changed
+1739 -51

No files matched your search

+418 -10
View File
@@ -30,6 +30,14 @@ register the (delivery-validated) file in the media registry and emit
``media.offer``; ``media.pull`` streams the file back as chunked binary
frames, re-checking ``validate_media_delivery_path`` at pull time.
Milestone M5: push + offline. Frames with no live subscriber are parked in
the outbox (M3) AND wake the device via the push backend (``push.py``: FCM
HTTP v1 primary, ntfy fallback, selected by ``ANDROID_PUSH_BACKEND``).
``notification`` frames render in-app banners and mirror to push (channel
events, cron deliveries, approvals, clarifies); high-priority kinds push even
when a device is live. ``fcm.register`` rotates push tokens (registry + live
connection). The outbox enforces a row cap with a throttled prune notice.
Configuration in config.yaml::
gateway:
@@ -106,6 +114,7 @@ from . import protocol # noqa: E402
from . import search as search_bridge # noqa: E402
from .channels import get_directory # noqa: E402
from .outbox import Outbox # noqa: E402
from .push import NtfyBackend, PushBackend, build_push_backend # noqa: E402
from .pairing import ( # noqa: E402
DeviceRegistry,
generate_token,
@@ -280,6 +289,42 @@ def _strip_streaming_cursor(text: str) -> str:
return text
def _push_preview(text: Any, limit: int = 120) -> str:
"""Short single-line preview for push bodies (lock-screen privacy: no
secrets, no full bodies -- full content arrives via ``sync``)."""
s = " ".join(str(text or "").split())
if len(s) > limit:
s = s[: limit - 1] + "…"
return s
# Cron delivery wrap (cron/scheduler.py ``_deliver_result``,
# cron.wrap_response: true):
# "Cronjob Response: <name>\n(job_id: <id>)\n-------------\n\n<content>\n\n
# To stop or manage this job, send me a new message (e.g. ...)."
_CRON_WRAP_RE = re.compile(
r"^Cronjob Response: (.+?)\n\(job_id: [^)]*\)\n-+\n\n"
)
_CRON_FOOTER = "\n\nTo stop or manage this job"
def _cron_brief(content: str, job_id: str) -> Tuple[str, str]:
"""Parse a cron delivery into ``(job_name, inner_text)``.
Falls back to ``(job_id, content)`` when the wrap is disabled
(``cron.wrap_response: false``) or unrecognised.
"""
m = _CRON_WRAP_RE.match(content or "")
if not m:
return str(job_id or "cron"), (content or "").strip()
name = m.group(1).strip()
body = content[m.end():]
idx = body.rfind(_CRON_FOOTER)
if idx != -1:
body = body[:idx]
return name, body.strip()
def _split_reasoning(text: str) -> Tuple[Optional[str], str]:
"""Split a code-style reasoning prefix off the front of *text*.
@@ -693,6 +738,17 @@ class AndroidAdapter(BasePlatformAdapter):
# last finalized assistant message id per chat (offer association).
self._media = media_bridge.MediaStore(get_hermes_home())
self._last_message_id: Dict[str, str] = {}
# M5: push backend (FCM primary, ntfy fallback) + the throttle for
# the outbox-prune banner.
self._push: PushBackend = build_push_backend(
self.push_backend,
fcm_service_account=_get_scoped_secret("ANDROID_FCM_SERVICE_ACCOUNT"),
fcm_server_key=_get_scoped_secret("ANDROID_FCM_SERVER_KEY"),
ntfy_topic=_get_scoped_secret("NTFY_TOPIC"),
ntfy_server_url=os.getenv("NTFY_SERVER_URL", "").strip() or None,
ntfy_auth_token=_get_scoped_secret("NTFY_AUTH_TOKEN"),
)
self._prune_notified_at = 0.0
def _turn_state(self, chat_id: str) -> _TurnState:
st = self._turns.get(chat_id)
@@ -747,6 +803,16 @@ class AndroidAdapter(BasePlatformAdapter):
except Exception:
logger.warning("android: ensure_default failed", exc_info=True)
# M5: push backend status (degrade gracefully when unconfigured).
if not self._push.configured():
logger.warning(
"android: push backend %r not configured (no credentials) -- "
"offline devices will not be woken; outbox + sync still apply",
self.push_backend,
)
else:
logger.info("android: push backend: %s", self._push.name)
self._connected = True
self._mark_connected()
logger.info("android: connected; WS server on %s:%s", self.host, self.port)
@@ -817,7 +883,45 @@ class AndroidAdapter(BasePlatformAdapter):
)
return SendResult(success=True, message_id=message_id)
# 2. Final message (non-streaming final, or streaming fallback final).
# 2. Cron delivery (the cron scheduler passes ``job_id`` in metadata):
# a final assistant message + a high-priority notification banner
# (pushed even when a device is live, docs/08 §8.1).
if meta.get("job_id"):
reasoning, _ = _split_reasoning(content)
if reasoning:
_reset_reasoning()
else:
await _wait_for_reasoning_flushed()
reasoning = _take_reasoning() or None
name, inner = _cron_brief(content, str(meta.get("job_id")))
message_id = _mint_message_id()
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_CRON,
f"Cron: {name}",
_push_preview(inner),
thread_id=thread_id,
),
)
await self._broadcast_or_log(
chat_id,
protocol.message(
chat_id=chat_id,
message_id=message_id,
role=protocol.ROLE_ASSISTANT,
text=inner,
thread_id=thread_id,
reasoning=reasoning,
ts=int(time.time() * 1000),
),
)
self._last_message_id[chat_id] = message_id
state.active = False
return SendResult(success=True, message_id=message_id)
# 3. Final message (non-streaming final, or streaming fallback final).
if meta.get("notify") is True:
reasoning, body = _split_reasoning(content)
# Non-streaming: reasoning is prepended to content (split above).
@@ -863,11 +967,11 @@ class AndroidAdapter(BasePlatformAdapter):
state.active = False
return SendResult(success=True, message_id=message_id)
# 3. Tool progress (first tool bubble of an editable line buffer).
# 4. Tool progress (first tool bubble of an editable line buffer).
if _is_tool_progress(content):
return await self._emit_tool_lines(chat_id, content, state, thread_id, is_edit=False)
# 4. Commentary (interim assistant beat).
# 5. Commentary (interim assistant beat).
message_id = _mint_message_id()
state.active = True
await self._broadcast_or_log(
@@ -1043,17 +1147,151 @@ class AndroidAdapter(BasePlatformAdapter):
async def _broadcast_or_log(self, chat_id: str, frame: "protocol.Frame") -> None:
delivered = await self._ws_server.broadcast(frame)
# M3/M5: always append to the outbox so a reconnecting app can catch
# up on *all* recent frames, not just the ones that were parked. This
# covers the case where the app's in-memory ChatStore is reset (e.g.
# activity recreation / composition recompose) while the process stays
# alive: the recreated controller re-syncs from its stale cursor and
# replays the frames it missed. The push below still only fires when
# there is no live subscriber (or for high-priority events).
try:
cursor = self._outbox.append(chat_id, frame.to_json())
except Exception:
logger.warning("android: outbox append failed", exc_info=True)
return
if delivered == 0:
# M3: no live device -- park the frame in the outbox so a
# reconnecting app can `sync` it (M5 adds push to wake the device).
logger.info(
"android: no live devices for %s; %s frame parked in outbox (cursor=%s)",
chat_id, frame.type, cursor,
)
# M5: wake the offline device(s) via the push backend.
await self._maybe_push(chat_id, frame, cursor)
await self._maybe_notify_outbox_prune(chat_id)
elif (
frame.type == protocol.TYPE_NOTIFICATION
and frame.payload.get("kind") in protocol.HIGH_PRIORITY_NOTIF_KINDS
):
# M5: high-priority events (approval/clarify/cron) push even when
# a device is live -- the app may be backgrounded and decides
# whether to also show an in-app banner (docs/08 §8.1).
await self._maybe_push(chat_id, frame, cursor)
# ── M5: push ───────────────────────────────────────────────────────────
def _push_summary(
self, frame: "protocol.Frame"
) -> Optional[Tuple[str, str, str, str]]:
"""``(title, body, kind, priority)`` for a pushable frame, else None.
Only terminal/interesting frames wake a device: intermediate
streaming and tool frames are replayed by ``sync`` without a push
(no notification spam per turn).
"""
t = frame.type
p = frame.payload
if t == protocol.TYPE_MESSAGE:
return (
self._channel_name(frame.chat_id or ""),
_push_preview(p.get("text")),
"message",
"normal",
)
if t == protocol.TYPE_MESSAGE_STOP:
return (
self._channel_name(frame.chat_id or ""),
_push_preview(p.get("final_text")),
"message",
"normal",
)
if t == protocol.TYPE_NOTIFICATION:
kind = str(p.get("kind") or protocol.NOTIF_GENERIC)
priority = "high" if kind in protocol.HIGH_PRIORITY_NOTIF_KINDS else "normal"
return (
str(p.get("title") or "Iris"),
str(p.get("body") or ""),
kind,
priority,
)
if t == protocol.TYPE_MEDIA_OFFER:
return (
self._channel_name(frame.chat_id or ""),
f"New {p.get('kind') or 'media'}: {p.get('filename') or ''}".strip(),
"media",
"normal",
)
return None
async def _maybe_push(self, chat_id: str, frame: "protocol.Frame", cursor: int) -> None:
"""Fire the configured push backend for a parked (or high-priority)
frame. Best-effort: failures are logged, never raised."""
summary = self._push_summary(frame)
if summary is None:
return
title, body, kind, priority = summary
backend = self._push
if backend is None or not backend.token_field:
return
devices = self._devices.list()
if not backend.configured() and not any(
d.get(backend.token_field) for d in devices
):
return
data: Dict[str, Any] = {"chat_id": chat_id, "kind": kind, "cursor": str(cursor)}
if frame.thread_id:
data["thread_id"] = frame.thread_id
message_id = frame.payload.get("message_id")
if isinstance(message_id, str) and message_id:
data["message_id"] = message_id
for device in devices:
device_id = device.get("device_id")
if not device_id:
continue
# Prefer the live connection's token (fcm.register refreshes it
# in memory) over the possibly-stale registry row.
conn = self._ws_server.connection(device_id)
token = getattr(conn, backend.token_field, None) if conn is not None else None
if not token:
token = device.get(backend.token_field)
if not token:
continue
try:
cursor = self._outbox.append(chat_id, frame.to_json())
logger.info(
"android: no live devices for %s; %s frame parked in outbox (cursor=%s)",
chat_id, frame.type, cursor,
ok = await backend.send(
device_id=device_id,
chat_id=chat_id,
title=title,
body=body,
data=data,
token=token,
priority=priority,
)
except Exception:
logger.warning("android: outbox append failed", exc_info=True)
logger.warning("android: push via %s failed", backend.name, exc_info=True)
continue
if ok:
logger.info(
"android: push via %s -> %s (%s, chat=%s)",
backend.name, device_id, frame.type, chat_id,
)
async def _maybe_notify_outbox_prune(self, chat_id: str) -> None:
"""When the outbox row cap pruned old frames, tell the app (throttled
to once per hour so a full box doesn't banner per frame)."""
pruned = self._outbox.take_overflow_pruned()
if pruned <= 0:
return
now = time.time()
if now - self._prune_notified_at < 3600.0:
return
self._prune_notified_at = now
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_GENERIC,
"Outbox",
f"{pruned} older message(s) pruned",
),
)
async def send_typing(self, chat_id: str, metadata: Optional[Dict[str, Any]] = None) -> None:
"""Send a typing indicator (``typing`` frame, on=true)."""
@@ -1464,6 +1702,16 @@ class AndroidAdapter(BasePlatformAdapter):
resp = protocol.channel_created(entry)
resp.id = frame.id
await self._ws_server.broadcast(resp)
# M5: banner + push mirror (parked in the outbox when offline).
await self._broadcast_or_log(
entry["chat_id"],
protocol.notification(
entry["chat_id"],
protocol.NOTIF_CHANNEL_CREATED,
"Channels",
f"New channel: {name}",
),
)
async def on_channel_rename(self, frame: protocol.Frame, device_id: str) -> None:
chat_id = frame.chat_id or frame.payload.get("chat_id")
@@ -1496,6 +1744,16 @@ class AndroidAdapter(BasePlatformAdapter):
resp = protocol.channel_renamed(entry)
resp.id = frame.id
await self._ws_server.broadcast(resp)
# M5: banner + push mirror (parked in the outbox when offline).
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_CHANNEL_RENAMED,
"Channels",
f"Renamed to {name}",
),
)
async def on_channel_set_default(self, frame: protocol.Frame, device_id: str) -> None:
chat_id = frame.chat_id or frame.payload.get("chat_id")
@@ -1536,6 +1794,16 @@ class AndroidAdapter(BasePlatformAdapter):
resp = protocol.channel_deleted(chat_id)
resp.id = frame.id
await self._ws_server.broadcast(resp)
# M5: banner + push mirror (parked in the outbox when offline).
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_CHANNEL_DELETED,
"Channels",
f"{entry.get('name') or chat_id} deleted",
),
)
async def on_channel_list(self, frame: protocol.Frame, device_id: str) -> None:
channels = self._channels.list(include_archived=False)
@@ -1596,6 +1864,140 @@ class AndroidAdapter(BasePlatformAdapter):
done = protocol.sync_done(self._outbox.latest_cursor(), id=frame.id)
await self._ws_server.send_to(device_id, done)
# ── M5: push token registration ───────────────────────────────────────
async def on_fcm_register(self, frame: protocol.Frame, device_id: str) -> None:
"""Update the device's push tokens (FCM rotation / ntfy topic).
Persists to the device registry AND refreshes the live connection so
the next push targets the current token without a stale read.
"""
fcm_token = frame.payload.get("fcm_token")
ntfy_topic = frame.payload.get("ntfy_topic")
fcm_token = fcm_token if isinstance(fcm_token, str) and fcm_token else None
ntfy_topic = ntfy_topic if isinstance(ntfy_topic, str) and ntfy_topic else None
if fcm_token is None and ntfy_topic is None:
return
try:
self._devices.update_push_tokens(
device_id, fcm_token=fcm_token, ntfy_topic=ntfy_topic
)
except Exception:
logger.warning("android: fcm.register update failed", exc_info=True)
return
conn = self._ws_server.connection(device_id)
if conn is not None:
if fcm_token is not None:
conn.fcm_token = fcm_token
if ntfy_topic is not None:
conn.ntfy_topic = ntfy_topic
logger.info("android: push tokens updated for %s", device_id)
# ── M5: approval / clarify banners ────────────────────────────────────
async def send_slash_confirm(
self,
chat_id: str,
title: str,
message: str,
session_key: str,
confirm_id: str,
metadata: Optional[Dict[str, Any]] = None,
) -> SendResult:
"""Banner + push for a slash-command approval prompt.
The gateway's text fallback still renders the actionable prompt (the
app has no inline buttons yet); the notification is the push-visible
signal (high priority: pushed even when a device is live).
"""
thread_id = _thread_id_from_metadata(metadata)
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_APPROVAL,
title or "Approval needed",
_push_preview(message),
thread_id=thread_id,
),
)
return await super().send_slash_confirm(
chat_id, title, message, session_key, confirm_id, metadata=metadata
)
async def send_clarify(
self,
chat_id: str,
question: str,
choices: Optional[list],
clarify_id: str,
session_key: str,
metadata: Optional[Dict[str, Any]] = None,
) -> SendResult:
"""Banner + push for a clarify prompt.
Renders the prompt as a proper ``message`` frame (the base text
fallback would route through ``send()`` and be misclassified as
commentary/tool progress) and keeps the gateway's text intercept
working via ``mark_awaiting_text``.
"""
thread_id = _thread_id_from_metadata(metadata)
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_CLARIFY,
"Question",
_push_preview(question),
thread_id=thread_id,
),
)
if choices:
# Multi-select clarifies register their flag on the pending entry;
# look it up by id (mirrors the base text fallback).
_is_multi = False
try:
from tools import clarify_gateway as _cg
with _cg._lock:
_entry = _cg._entries.get(clarify_id)
_is_multi = bool(_entry and getattr(_entry, "multi_select", False))
except Exception:
_is_multi = False
lines = [f"❓ {question}", ""]
for i, choice in enumerate(choices, start=1):
lines.append(f" {i}. {choice}")
lines.append("")
if _is_multi:
lines.append(
"Multiple selections allowed — reply with the numbers "
"separated by commas or spaces (e.g. \"1, 3\"), the option "
"text, or your own answer."
)
else:
lines.append("Reply with the number, the option text, or your own answer.")
text = "\n".join(lines)
# Text fallback: enable text-capture so the gateway intercept
# picks up the user's typed reply (e.g. "2" or choice text).
from tools.clarify_gateway import mark_awaiting_text
mark_awaiting_text(clarify_id)
else:
text = f"❓ {question}"
message_id = _mint_message_id()
await self._broadcast_or_log(
chat_id,
protocol.message(
chat_id=chat_id,
message_id=message_id,
role=protocol.ROLE_ASSISTANT,
text=text,
thread_id=thread_id,
ts=int(time.time() * 1000),
),
)
return SendResult(success=True, message_id=message_id)
# ── Chat info ─────────────────────────────────────────────────────────
def _channel_name(self, chat_id: str) -> str:
@@ -1630,6 +2032,12 @@ class AndroidAdapter(BasePlatformAdapter):
"media": True, # M4: media.upload/offer/pull
"search": True, # M3: search frame
"push": self.push_backend,
# M5: ntfy server URL (app listener discovery; "" when not ntfy).
"push_ntfy_server": (
self._push.server_url
if isinstance(self._push, NtfyBackend)
else ""
),
"pickers": False, # M2+
}
+39 -1
View File
@@ -28,6 +28,10 @@ logger = logging.getLogger(__name__)
DEFAULT_RETENTION_HOURS = 72
_REPLAY_LIMIT = 1000
_PRUNE_INTERVAL_S = 3600.0
# Row cap (docs/08 §8.3): a very busy offline period must not grow the outbox
# without bound; the oldest rows beyond the cap are pruned and the app is
# told (generic notification, throttled in the adapter).
DEFAULT_MAX_ROWS = 5000
class Outbox:
@@ -38,12 +42,19 @@ class Outbox:
``DeviceRegistry`` / ``ChannelDirectory``).
"""
def __init__(self, db_path: Path, retention_hours: int = DEFAULT_RETENTION_HOURS):
def __init__(
self,
db_path: Path,
retention_hours: int = DEFAULT_RETENTION_HOURS,
max_rows: int = DEFAULT_MAX_ROWS,
):
self._db_path = Path(db_path)
self._db_path.parent.mkdir(parents=True, exist_ok=True)
self._retention_hours = max(1, int(retention_hours))
self._max_rows = max(1, int(max_rows))
self._lock = threading.Lock()
self._last_prune = 0.0
self._overflow_pruned = 0
self._conn = sqlite3.connect(str(self._db_path), check_same_thread=False)
self._conn.row_factory = sqlite3.Row
with self._lock:
@@ -90,10 +101,37 @@ class Outbox:
"VALUES (?, ?, ?, ?)",
(cursor, chat_id, frame_json, now),
)
self._enforce_row_cap()
self._conn.commit()
self._maybe_prune()
return cursor
def take_overflow_pruned(self) -> int:
"""Rows pruned by the row cap since the last call (and reset to 0).
The adapter turns a non-zero count into a (throttled) generic
notification so the app knows older frames are gone.
"""
with self._lock:
n = self._overflow_pruned
self._overflow_pruned = 0
return n
def _enforce_row_cap(self) -> None:
"""Drop the oldest rows beyond ``max_rows`` (caller holds the lock)."""
row = self._conn.execute("SELECT COUNT(*) AS n FROM outbox").fetchone()
n = int(row["n"]) if row else 0
excess = n - self._max_rows
if excess <= 0:
return
self._conn.execute(
"DELETE FROM outbox WHERE cursor IN ("
" SELECT cursor FROM outbox ORDER BY cursor ASC LIMIT ?)",
(excess,),
)
self._overflow_pruned += excess
logger.info("android outbox: row cap pruned %s oldest row(s)", excess)
def latest_cursor(self) -> int:
"""The high-water cursor (0 when nothing has been appended)."""
with self._lock:
+4
View File
@@ -57,6 +57,10 @@ optional_env:
description: "ntfy server URL (default https://ntfy.sh)"
prompt: "ntfy server URL"
password: false
- name: NTFY_AUTH_TOKEN
description: "ntfy auth token for a private topic (trust boundary)"
prompt: "ntfy auth token"
password: true
- name: ANDROID_WS_CERT
description: "TLS cert path for WSS (optional)"
prompt: "WSS cert"
+55
View File
@@ -76,6 +76,10 @@ TYPE_MEDIA_OFFER = "media.offer"
TYPE_MEDIA_PULL = "media.pull"
TYPE_MEDIA_PULL_END = "media.pull.end"
# Push / notifications (M5)
TYPE_NOTIFICATION = "notification"
TYPE_FCM_REGISTER = "fcm.register"
# ---------------------------------------------------------------------------
# Error codes (``error`` frame payload.code)
# ---------------------------------------------------------------------------
@@ -96,6 +100,22 @@ ROLE_ASSISTANT = "assistant"
ROLE_SYSTEM = "system"
ROLE_CRON = "cron"
# ---------------------------------------------------------------------------
# Notification kinds (``notification`` frame payload.kind, M5)
# ---------------------------------------------------------------------------
NOTIF_CHANNEL_CREATED = "channel_created"
NOTIF_CHANNEL_RENAMED = "channel_renamed"
NOTIF_CHANNEL_DELETED = "channel_deleted"
NOTIF_CRON = "cron"
NOTIF_APPROVAL = "approval"
NOTIF_CLARIFY = "clarify"
NOTIF_GENERIC = "generic"
# Kinds that push even when a device is live (the app may be backgrounded;
# it decides whether to also show an in-app banner).
HIGH_PRIORITY_NOTIF_KINDS = frozenset({NOTIF_APPROVAL, NOTIF_CLARIFY, NOTIF_CRON})
# ---------------------------------------------------------------------------
# Envelope
@@ -452,6 +472,41 @@ def sync_done(cursor: int, *, id: Optional[int] = None) -> Frame:
return Frame(type=TYPE_SYNC_DONE, id=id, payload={"cursor": cursor})
# ---------------------------------------------------------------------------
# Push / notification frames (M5)
# ---------------------------------------------------------------------------
def notification(
chat_id: str,
kind: str,
title: str,
body: str,
*,
thread_id: Optional[str] = None,
ts: Optional[int] = None,
) -> Frame:
"""Event: a transient in-app banner (and a push mirror when the device is
offline). ``kind`` is one of the ``NOTIF_*`` constants."""
payload: Dict[str, Any] = {"kind": kind, "title": title, "body": body}
if ts is not None:
payload["ts"] = ts
return Frame(
type=TYPE_NOTIFICATION, chat_id=chat_id, thread_id=thread_id, payload=payload
)
def fcm_register(
fcm_token: Optional[str] = None, ntfy_topic: Optional[str] = None
) -> Frame:
"""Request: update the device's push tokens (FCM rotation / ntfy topic)."""
payload: Dict[str, Any] = {}
if fcm_token:
payload["fcm_token"] = fcm_token
if ntfy_topic:
payload["ntfy_topic"] = ntfy_topic
return Frame(type=TYPE_FCM_REGISTER, payload=payload)
# ---------------------------------------------------------------------------
# Media frames (M4)
# ---------------------------------------------------------------------------
+322 -6
View File
@@ -4,12 +4,328 @@
- ``FcmBackend``: FCM HTTP v1 via ``httpx`` + a Firebase service account
(``ANDROID_FCM_SERVICE_ACCOUNT``), or a legacy server key
(``ANDROID_FCM_SERVER_KEY``).
- ``NtfyBackend``: reuses hermes ntfy publish (``NTFY_TOPIC`` /
``NTFY_SERVER_URL``).
- ``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 ``ANDROID_PUSH_BACKEND`` (``fcm`` default, ``ntfy`` fallback).
Fired when a frame has no live subscriber; data payload drives a silent sync
on the device.
Fired when a frame has no live subscriber; the data payload drives a silent
sync on the device (docs/08-push.md).
Milestone M5.
"""
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, Dict, Optional
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
_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.
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"
token_field = "fcm_token"
def __init__(
self,
service_account: Optional[str] = None,
server_key: Optional[str] = None,
):
self._sa_path = (service_account or "").strip() or None
self._server_key = (server_key or "").strip() or None
self._sa: Optional[Dict[str, Any]] = None
self._sa_failed = False
self._access_token: Optional[str] = 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) -> Optional[Dict[str, Any]]:
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, "r", 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("android: FCM service account unreadable: %s", self._sa_path)
self._sa_failed = True
return None
async def _authorization(self, client: httpx.AsyncClient) -> Optional[str]:
"""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("android: 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("android: FCM token exchange failed", exc_info=True)
return None
if resp.status_code != 200:
logger.warning(
"android: 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("android: FCM send failed (network)", exc_info=True)
return False
if resp.status_code >= 300:
# 404 NOT_FOUND = stale/invalid registration token.
logger.warning(
"android: 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"
token_field = "ntfy_topic"
def __init__(
self,
topic: Optional[str] = None,
server_url: Optional[str] = None,
auth_token: Optional[str] = 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("android: ntfy publish failed (network)", exc_info=True)
return False
if resp.status_code >= 300:
logger.warning(
"android: ntfy publish HTTP %s: %s", resp.status_code, resp.text[:200]
)
return False
return True
def build_push_backend(
name: Optional[str],
*,
fcm_service_account: Optional[str] = None,
fcm_server_key: Optional[str] = None,
ntfy_topic: Optional[str] = None,
ntfy_server_url: Optional[str] = None,
ntfy_auth_token: Optional[str] = None,
) -> PushBackend:
"""Select the backend by name (``ANDROID_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)
+51 -3
View File
@@ -18,8 +18,12 @@ Options:
--send TEXT send this message after pairing (default: "hello")
--upload F M4: upload F (chunked media.upload) and attach it to the
message.send via media_refs
--pull-offer M4: when a media.offer arrives during the turn, pull the
media (chunked) and verify the byte count
--pull-offer M4: when a media.offer arrives during the turn, pull the
media (chunked) and verify the byte count
--sync C M5: after pairing, send sync {cursor: C} and print the
replay + sync.done (no turn is driven)
--fcm-token M5: attach this FCM token to the hello payload
--fcm-reg M5: after pairing, send fcm.register with --fcm-token
--timeout S seconds to wait for the final reply (default 120)
--authfail expect an auth rejection (wrong token) and exit 0 on it
"""
@@ -94,6 +98,13 @@ def _print_frame(raw):
extra = f" ok={payload.get('ok')} ref={payload.get('media_ref')}"
elif ftype == "media.pull.end":
extra = f" ok={payload.get('ok')}"
elif ftype == "notification":
extra = (f" kind={payload.get('kind')} title={payload.get('title')!r} "
f"body={(payload.get('body') or '')[:100]!r}")
elif ftype == "sync":
extra = f" cursor={payload.get('cursor')}"
elif ftype == "sync.done":
extra = f" cursor={payload.get('cursor')}"
scope = f" chat={chat}" if chat else ""
idpart = f" id={fid}" if fid is not None else ""
print(f" <- {ftype}{idpart}{scope}{extra}")
@@ -200,8 +211,10 @@ async def run(args) -> int:
"caps": {"min_protocol": 1},
},
}
if args.fcm_token:
hello["payload"]["fcm_token"] = args.fcm_token
await ws.send(json.dumps(hello))
print(" -> hello")
print(" -> hello" + (f" fcm_token={args.fcm_token[:12]}…" if args.fcm_token else ""))
# First response must be hello.ack (or an auth error).
try:
@@ -224,6 +237,35 @@ async def run(args) -> int:
await ws.close()
return 5
# M5: optional fcm.register after pairing.
if args.fcm_reg:
reg_token = args.fcm_token or f"probe-{uuid.uuid4().hex[:12]}"
await ws.send(json.dumps({
"v": 1, "type": "fcm.register",
"payload": {"fcm_token": reg_token},
}))
print(f" -> fcm.register fcm_token={reg_token[:12]}…")
# M5: sync catch-up mode (no turn driven).
if args.sync is not None:
await ws.send(json.dumps({
"v": 1, "id": 1, "type": "sync", "payload": {"cursor": args.sync},
}))
print(f" -> sync cursor={args.sync}")
while True:
raw = await asyncio.wait_for(ws.recv(), timeout=30)
data = _print_frame(raw)
if data is None:
continue
if data.get("type") == "sync.done":
print(f"== sync done at cursor {data['payload'].get('cursor')}")
await ws.close()
return 0
if data.get("type") == "error":
print(f"!! sync failed: {data['payload']}")
await ws.close()
return 8
if not args.send and not args.upload:
print("== paired OK (no --send/--upload; exiting)")
await ws.close()
@@ -311,6 +353,12 @@ def main() -> int:
help="M4: file to upload (chunked) and attach via media_refs")
p.add_argument("--pull-offer", action="store_true",
help="M4: pull any media.offer that arrives during the turn")
p.add_argument("--sync", type=int, default=None,
help="M5: send sync {cursor} after pairing, print replay, exit")
p.add_argument("--fcm-token", default="",
help="M5: FCM token to attach to the hello payload")
p.add_argument("--fcm-reg", action="store_true",
help="M5: send fcm.register after pairing (uses --fcm-token)")
p.add_argument("--timeout", type=float, default=120.0)
p.add_argument("--authfail", action="store_true",
help="expect an auth rejection (wrong token)")
+2 -12
View File
@@ -316,18 +316,8 @@ class WsServer:
await self._adapter.on_media_upload_end(frame, device_id)
elif frame.type == protocol.TYPE_MEDIA_PULL:
await self._adapter.on_media_pull(frame, device_id)
elif frame.type == "fcm.register":
fcm_token = frame.payload.get("fcm_token")
ntfy_topic = frame.payload.get("ntfy_topic")
if isinstance(fcm_token, str) or isinstance(ntfy_topic, str):
try:
self._devices.update_push_tokens(
device_id,
fcm_token=fcm_token if isinstance(fcm_token, str) else None,
ntfy_topic=ntfy_topic if isinstance(ntfy_topic, str) else None,
)
except Exception:
logger.warning("android: fcm.register update failed", exc_info=True)
elif frame.type == protocol.TYPE_FCM_REGISTER:
await self._adapter.on_fcm_register(frame, device_id)
# Unknown types are ignored (forward-compat).
# ── Helpers ───────────────────────────────────────────────────────────