aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src/meshbay_hub/api
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-10-09 12:08:31 +0200
committerChristophe Besson <cbesson@gmail.com>2026-10-09 12:08:31 +0200
commit6832df6177ad973ad0e1b4f0a49d7a6da06c6e04 (patch)
treed9040ce0da5d82400d1b973344615ca1c6b67b3c /packages/meshbay-hub/src/meshbay_hub/api
parent2860f1de75af1d44d35292ecbf79c68f02409d19 (diff)
downloadmeshbay-6832df6177ad973ad0e1b4f0a49d7a6da06c6e04.tar.gz
feat: notifications on Android while closed, with nothing to install
The phone fetches what is new every fifteen minutes with a poll secret (POST /v1/push/poll) that reads notification lines and nothing else. When a UnifiedPush distributor is already installed, the hub also pushes at once, encrypted to the phone (RFC 8291); losing the distributor falls back to fetching. The hub now honours "disable all notifications" itself: create_notification creates nothing for that account, as it already did for a muted group, so neither switch lets anything reach a phone. The interface used to be the only reader of the account-wide switch. Push endpoints are member-supplied URLs: a send refuses non-public addresses, connects to the address it checked, and follows no redirect. Android build untested here (no SDK on this machine). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/api')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/notifications.py23
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/push.py277
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/users.py3
3 files changed, 300 insertions, 3 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/notifications.py b/packages/meshbay-hub/src/meshbay_hub/api/notifications.py
index e4bac2e..128c0d1 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/notifications.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/notifications.py
@@ -26,8 +26,9 @@ from sqlalchemy import delete, func, select
from sqlalchemy.ext.asyncio import AsyncSession
from meshbay_hub.api.deps import get_current_user
+from meshbay_hub.api.push import push_notification
from meshbay_hub.db.engine import get_db
-from meshbay_hub.db.models import GroupMember, Notification, User
+from meshbay_hub.db.models import GroupMember, Notification, User, UserPreference
router = APIRouter(prefix="/v1/notifications", tags=["notifications"])
@@ -157,9 +158,23 @@ async def create_notification(
conversation is a single line saying when it last spoke rather than forty
saying that it spoke.
- Returns None when the person muted this group: the point of muting is that
- nothing is created, not that something is created and hidden.
+ Returns None when the person muted this group, or turned every notification
+ off: the point of muting is that nothing is created, not that something is
+ created and hidden — and nothing created is nothing pushed to a phone.
+
+ The account-wide switch used to be read by the interface alone, which hid
+ the list while rows went on accumulating; with a phone that is told about
+ each row, a switch only the interface honours is a switch that does nothing.
"""
+ disabled = await db.execute(
+ select(UserPreference.value).where(
+ UserPreference.user_id == user_id,
+ UserPreference.key == "notifications_disabled",
+ )
+ )
+ if disabled.scalar() == "true":
+ return None
+
if group_id is not None:
muted = await db.execute(
select(GroupMember.muted).where(
@@ -185,6 +200,7 @@ async def create_notification(
existing.read = False
existing.created_at = datetime.now(UTC)
await db.flush()
+ await push_notification(db, existing)
return existing
notif = Notification(
@@ -193,4 +209,5 @@ async def create_notification(
)
db.add(notif)
await db.flush()
+ await push_notification(db, notif)
return notif
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/push.py b/packages/meshbay-hub/src/meshbay_hub/api/push.py
new file mode 100644
index 0000000..0c4a411
--- /dev/null
+++ b/packages/meshbay-hub/src/meshbay_hub/api/push.py
@@ -0,0 +1,277 @@
+"""
+Notifications on a phone — /v1/push/*: pushed when it can be, fetched when not.
+
+A phone registers once and gets a row here. **With a UnifiedPush distributor**
+it gives an endpoint and a P-256 key, and every notification
+`create_notification` lets through is sent there, encrypted to the phone
+(`webpush.py`). **Without one** — nothing to install is the default — the row has
+no endpoint, and the phone fetches what is new with `POST /v1/push/poll` every
+quarter of an hour or so. Both are the same rows and the same payload, so a phone
+can move between them (a distributor installed, removed, refusing) without the
+hub caring which.
+
+**Nothing reaches a phone that was not created**: a muted group and an account
+with every notification turned off stop at `create_notification`, before this
+module is reached, so the two switches the person sees are the only two there are.
+
+The poll is authenticated by a secret issued with the row, not by a session. It
+reads notification lines and nothing else, so a phone running in the background
+holds no token that could do anything more — and a sign-out, which deletes the
+row, ends it.
+
+Who pays (§13.5b): a member's chat costs every other member's phones a push.
+That fan-out is already bounded where it starts — `chat_notify` is budgeted per
+node — and here a conversation reaches each phone at most once per
+`CHAT_COALESCE` seconds: the phone shows one line per group, so the pushes in
+between would only have replaced it. An account holds `MAX_SUBSCRIPTIONS` rows
+at most, because each is one outbound request per notification; a row is polled
+at most once per `POLL_MIN_INTERVAL`.
+"""
+
+import asyncio
+import hashlib
+import hmac
+import logging
+import math
+import secrets
+import time
+from datetime import UTC, datetime
+
+from fastapi import APIRouter, Depends, HTTPException
+from pydantic import BaseModel, Field
+from sqlalchemy import func, select, update
+from sqlalchemy.ext.asyncio import AsyncSession
+
+from meshbay_hub import webpush
+from meshbay_hub.api.deps import require_user_scope
+from meshbay_hub.db.engine import get_db
+from meshbay_hub.db.models import Notification, PushSubscription, User
+
+log = logging.getLogger(__name__)
+
+router = APIRouter(prefix="/v1/push", tags=["push"])
+
+MAX_SUBSCRIPTIONS = 10
+CHAT_COALESCE = 30.0
+_COALESCE_ENTRIES = 10_000
+POLL_MIN_INTERVAL = 60.0
+POLL_LIMIT = 20
+
+# Strong references: asyncio holds a task weakly, and a collected one is a push
+# that silently never went (CLAUDE.md, "a background task nobody holds").
+_tasks: set[asyncio.Task] = set()
+_last_chat: dict[tuple[str, str], float] = {}
+_last_poll: dict[str, float] = {}
+# Replaced by the tests; the real one never raises.
+_send = webpush.send
+
+
+class SubscriptionIn(BaseModel):
+ # The row this phone already has, to update rather than add one: a phone
+ # with no endpoint has nothing else to be recognised by.
+ id: str | None = Field(default=None, max_length=36)
+ endpoint: str | None = Field(default=None, max_length=webpush.MAX_ENDPOINT)
+ # Lengths bounded before decoding: base64 decoding skips characters outside
+ # its alphabet, so an unbounded string could still decode to 65 bytes.
+ p256dh: str | None = Field(default=None, max_length=128)
+ auth: str | None = Field(default=None, max_length=32)
+
+
+def _hash(secret: str) -> str:
+ return hashlib.sha256(secret.encode()).hexdigest()
+
+
+def _payload(notif: Notification) -> dict:
+ """What a phone is told, pushed or fetched: the hub's own line, never a message."""
+ return {
+ "v": 1,
+ "id": notif.id,
+ "kind": notif.kind,
+ "title": notif.title,
+ "link": notif.link,
+ "group_id": notif.group_id,
+ "created_at": _iso(notif.created_at),
+ }
+
+
+def _iso(at: datetime) -> str:
+ # SQLite hands back naive datetimes; every one stored here is UTC.
+ return (at if at.tzinfo else at.replace(tzinfo=UTC)).isoformat()
+
+
+@router.post("/subscriptions")
+async def subscribe(
+ body: SubscriptionIn,
+ current_user: User = Depends(require_user_scope),
+ db: AsyncSession = Depends(get_db),
+):
+ """
+ Register this phone, or update its row: with an endpoint and keys when it
+ has a push distributor, without them when it will fetch instead.
+
+ Answers the row's id, a fresh secret for `POST /v1/push/poll` (the previous
+ one stops working) and the hub's time, from which the phone counts what is
+ new — what was there before it registered is not news.
+ """
+ if body.endpoint is not None:
+ if body.p256dh is None or body.auth is None:
+ raise HTTPException(status_code=422, detail="an endpoint needs its keys")
+ try:
+ webpush.check_endpoint(body.endpoint)
+ webpush.check_keys(body.p256dh, body.auth)
+ except ValueError as e:
+ raise HTTPException(status_code=422, detail=str(e)) from e
+
+ sub = None
+ if body.id is not None:
+ sub = await db.get(PushSubscription, body.id)
+ if sub is not None and sub.user_id != current_user.id:
+ sub = None
+ if body.endpoint is not None:
+ same = (await db.execute(
+ select(PushSubscription).where(
+ PushSubscription.user_id == current_user.id,
+ PushSubscription.endpoint == body.endpoint))).scalar_one_or_none()
+ if sub is None:
+ sub = same
+ elif same is not None and same.id != sub.id:
+ # The endpoint moved to this row; the old one would only repeat it.
+ await db.delete(same)
+ await db.flush()
+ if sub is None:
+ held = (await db.execute(
+ select(func.count()).select_from(PushSubscription)
+ .where(PushSubscription.user_id == current_user.id))).scalar() or 0
+ if held >= MAX_SUBSCRIPTIONS:
+ raise HTTPException(status_code=429, detail="Too many push subscriptions")
+ sub = PushSubscription(user_id=current_user.id)
+ db.add(sub)
+ secret = secrets.token_urlsafe(32)
+ sub.endpoint, sub.p256dh, sub.auth = body.endpoint, body.p256dh, body.auth
+ sub.poll_hash = _hash(secret)
+ await db.commit()
+ return {"id": sub.id, "poll_secret": secret, "now": datetime.now(UTC).isoformat()}
+
+
+class PollIn(BaseModel):
+ id: str = Field(max_length=36)
+ secret: str = Field(max_length=64)
+ since: datetime
+
+
+@router.post("/poll")
+async def poll(body: PollIn, db: AsyncSession = Depends(get_db)):
+ """
+ What is new for this phone since `since`: the same payloads a push carries,
+ oldest first, at most twenty.
+
+ Authenticated by the row's secret rather than a session, so what a phone
+ keeps for running in the background reads notification lines and nothing
+ else. A wrong secret and an unknown row answer the same 404.
+ """
+ sub = await db.get(PushSubscription, body.id)
+ if (sub is None or sub.poll_hash is None
+ or not hmac.compare_digest(sub.poll_hash, _hash(body.secret))):
+ raise HTTPException(status_code=404, detail="Subscription not found")
+ now = time.monotonic()
+ last = _last_poll.get(sub.id)
+ if last is not None and now - last < POLL_MIN_INTERVAL:
+ raise HTTPException(status_code=429, detail="Polled too often",
+ headers={"Retry-After": str(int(POLL_MIN_INTERVAL - (now - last)) + 1)})
+ if len(_last_poll) >= _COALESCE_ENTRIES:
+ for k in [k for k, at in _last_poll.items() if now - at >= POLL_MIN_INTERVAL]:
+ del _last_poll[k]
+ _last_poll[sub.id] = now
+
+ since = body.since if body.since.tzinfo else body.since.replace(tzinfo=UTC)
+ rows = (await db.execute(
+ select(Notification).where(
+ Notification.user_id == sub.user_id,
+ Notification.created_at > since.astimezone(UTC),
+ ).order_by(Notification.created_at.desc()).limit(POLL_LIMIT)
+ )).scalars().all()
+ return {"notifications": [_payload(n) for n in reversed(rows)]}
+
+
+@router.delete("/subscriptions/{subscription_id}")
+async def unsubscribe(
+ subscription_id: str,
+ current_user: User = Depends(require_user_scope),
+ db: AsyncSession = Depends(get_db),
+):
+ """Stop telling one phone anything: turned off there, or signed out of."""
+ sub = await db.get(PushSubscription, subscription_id)
+ if sub is None or sub.user_id != current_user.id:
+ raise HTTPException(status_code=404, detail="Subscription not found")
+ await db.delete(sub)
+ await db.commit()
+ return {"status": "ok"}
+
+
+def _coalesced(sub_id: str, group_id: str, now: float) -> bool:
+ key = (sub_id, group_id)
+ if now - _last_chat.get(key, -math.inf) < CHAT_COALESCE:
+ return True
+ if len(_last_chat) >= _COALESCE_ENTRIES:
+ for k in [k for k, at in _last_chat.items() if now - at >= CHAT_COALESCE]:
+ del _last_chat[k]
+ _last_chat[key] = now
+ return False
+
+
+async def push_notification(db: AsyncSession, notif: Notification) -> None:
+ """
+ Send `notif` to the person's phones, off the caller's path.
+
+ Called by `create_notification` once the row exists, inside the caller's
+ transaction: the subscriptions are read there, the requests leave in a task
+ of their own, so a slow push server delays nobody's request.
+ """
+ subs = (await db.execute(
+ select(PushSubscription).where(PushSubscription.user_id == notif.user_id,
+ PushSubscription.endpoint.is_not(None))
+ )).scalars().all()
+ if not subs:
+ return
+ payload = _payload(notif)
+ now = time.monotonic()
+ targets = [
+ (s.id, webpush.Target(s.endpoint, s.p256dh, s.auth)) for s in subs
+ if not (notif.kind == "chat_message" and notif.group_id
+ and _coalesced(s.id, notif.group_id, now))
+ ]
+ if not targets:
+ return
+ ttl = webpush.TTL_CHAT if notif.kind == "chat_message" else webpush.TTL_OTHER
+ task = asyncio.get_running_loop().create_task(_deliver(targets, payload, ttl))
+ _tasks.add(task)
+ task.add_done_callback(_tasks.discard)
+
+
+async def _deliver(targets: list[tuple[str, webpush.Target]], payload: dict,
+ ttl: int) -> None:
+ results = await asyncio.gather(*(_send(t, payload, ttl=ttl) for _, t in targets),
+ return_exceptions=True)
+ gone = [sid for (sid, _), r in zip(targets, results, strict=True) if r == webpush.GONE]
+ if not gone:
+ return
+ # The distributor dropped the registration: pushing there would cost a
+ # request per notification for ever and reach nothing. The row stays, without
+ # its endpoint — the phone still fetches, and registers again when it can.
+ try:
+ from meshbay_hub.db.engine import get_session_factory
+ async with get_session_factory()() as db:
+ await db.execute(
+ update(PushSubscription).where(PushSubscription.id.in_(gone))
+ .values(endpoint=None, p256dh=None, auth=None))
+ await db.commit()
+ except Exception as e:
+ log.warning("Could not drop %d gone push subscription(s): %s", len(gone), e)
+
+
+async def drain() -> None:
+ """Wait for every push in flight — for the tests, and for a clean shutdown."""
+ # Done tasks leave the set from a callback the loop has not run yet, and
+ # awaiting a finished gather never yields to it: wait on the unfinished only.
+ while pending := [t for t in _tasks if not t.done()]:
+ await asyncio.gather(*pending, return_exceptions=True)
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/users.py b/packages/meshbay-hub/src/meshbay_hub/api/users.py
index 795a902..130240e 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/users.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/users.py
@@ -46,6 +46,7 @@ from meshbay_hub.db.models import (
KnownBrowser,
Node,
Notification,
+ PushSubscription,
RefreshToken,
User,
UserDevice,
@@ -1361,6 +1362,7 @@ async def password_reset(
.values(revoked=True))
await db.execute(delete(UserDevice).where(UserDevice.user_id == user.id))
await db.execute(delete(KnownBrowser).where(KnownBrowser.user_id == user.id))
+ await db.execute(delete(PushSubscription).where(PushSubscription.user_id == user.id))
# A code sent to the address on file is a stronger proof than a passphrase,
# and it is the way out of a lockout somebody else caused.
await login_throttle.clear(db, user.username)
@@ -1579,6 +1581,7 @@ async def erase_account(db: AsyncSession, user: User, owned_groups: str = "refus
await db.execute(delete(Node).where(Node.user_id == user.id))
await db.execute(delete(UserDevice).where(UserDevice.user_id == user.id))
await db.execute(delete(KnownBrowser).where(KnownBrowser.user_id == user.id))
+ await db.execute(delete(PushSubscription).where(PushSubscription.user_id == user.id))
await db.execute(delete(EmailVerification).where(EmailVerification.user_id == user.id))
# Links this account issued for a group it no longer owns; the ones for its
# own groups went with them above. A used link keeps pointing at the