Files
iris_x_hermes/gateway-plugin/outbox.py
T
ARIA 2ecfe1c05c 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.
2026-08-19 19:20:55 +02:00

197 lines
7.5 KiB
Python

"""SQLite offline outbox + monotonic sync cursor (M3).
Undelivered frames (sent while no device is live) are appended with a
**monotonic** cursor so a reconnecting app can ``sync {cursor}`` the delta
without re-reading full history. The cursor is a separate high-water counter
that only ever increases -- pruning old rows never resets it, so a late
reconnect can't be handed a cursor lower than one it already saw.
Retention prunes rows older than ``outbox_retention_hours`` (default 72h). A
device offline longer than the window misses those frames; it recovers full
context via ``history`` (M5 wires push so the device is woken to sync).
Storage: ``get_hermes_home()/"android"/outbox.db``.
Milestone M3 (built), extended in M5 (push integration).
"""
import json
import logging
import sqlite3
import threading
import time
from pathlib import Path
from typing import Any, Dict, List, Optional
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:
"""Persistent outbox under ``get_hermes_home()/"android"``.
Thread-safe (single connection + lock); operations are small and fast
enough to run inline on the gateway's asyncio loop (mirrors
``DeviceRegistry`` / ``ChannelDirectory``).
"""
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:
self._conn.execute("PRAGMA journal_mode=WAL")
self._conn.execute(
"""
CREATE TABLE IF NOT EXISTS outbox (
cursor INTEGER PRIMARY KEY,
chat_id TEXT,
frame TEXT NOT NULL,
created REAL NOT NULL DEFAULT 0
)
"""
)
self._conn.execute(
"CREATE INDEX IF NOT EXISTS idx_outbox_created ON outbox (created)"
)
self._conn.execute(
"""
CREATE TABLE IF NOT EXISTS counters (
name TEXT PRIMARY KEY,
value INTEGER NOT NULL DEFAULT 0
)
"""
)
self._conn.commit()
# ── append / cursor ───────────────────────────────────────────────────
def append(self, chat_id: Optional[str], frame_json: str) -> int:
"""Append a frame; returns the (monotonic) cursor assigned to it."""
now = time.time()
with self._lock:
self._conn.execute(
"INSERT INTO counters (name, value) VALUES ('cursor', 1) "
"ON CONFLICT(name) DO UPDATE SET value = value + 1"
)
row = self._conn.execute(
"SELECT value FROM counters WHERE name = 'cursor'"
).fetchone()
cursor = int(row["value"]) if row else 1
self._conn.execute(
"INSERT INTO outbox (cursor, chat_id, frame, created) "
"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:
row = self._conn.execute(
"SELECT value FROM counters WHERE name = 'cursor'"
).fetchone()
return int(row["value"]) if row else 0
# ── replay ────────────────────────────────────────────────────────────
def replay(self, cursor: int, limit: int = _REPLAY_LIMIT) -> List[Dict[str, Any]]:
"""Frames with ``cursor > `cursor```, oldest first.
Each entry: ``{cursor, chat_id, frame}`` where ``frame`` is the parsed
frame dict (the caller re-serializes / forwards it to the device).
"""
cursor = max(0, int(cursor or 0))
limit = max(1, min(int(limit or _REPLAY_LIMIT), _REPLAY_LIMIT))
with self._lock:
rows = self._conn.execute(
"SELECT cursor, chat_id, frame FROM outbox "
"WHERE cursor > ? ORDER BY cursor ASC LIMIT ?",
(cursor, limit),
).fetchall()
out: List[Dict[str, Any]] = []
for r in rows:
try:
frame = json.loads(r["frame"])
except (json.JSONDecodeError, TypeError):
continue
if not isinstance(frame, dict):
continue
out.append(
{"cursor": int(r["cursor"]), "chat_id": r["chat_id"], "frame": frame}
)
return out
# ── retention ─────────────────────────────────────────────────────────
def _maybe_prune(self) -> None:
now = time.time()
if now - self._last_prune < _PRUNE_INTERVAL_S:
return
self._last_prune = now
cutoff = now - self._retention_hours * 3600
with self._lock:
try:
self._conn.execute("DELETE FROM outbox WHERE created < ?", (cutoff,))
self._conn.commit()
except sqlite3.Error as e:
logger.debug("android outbox: prune failed: %s", e)
def prune(self) -> None:
"""Force a retention prune (ignores the interval throttle)."""
self._last_prune = 0.0
self._maybe_prune()
def close(self) -> None:
with self._lock:
try:
self._conn.close()
except Exception:
pass