diff --git a/go.mod b/go.mod
index 16a3cf5a..f775d595 100644
--- a/go.mod
+++ b/go.mod
@@ -18,6 +18,7 @@ require (
firebase.google.com/go/v4 v4.14.1
github.com/99designs/gqlgen v0.17.64
github.com/NYTimes/gziphandler v1.1.1
+ github.com/SherClockHolmes/webpush-go v1.4.0
github.com/ThalesGroup/crypto11 v0.0.0-00010101000000-000000000000
github.com/acarl005/stripansi v0.0.0-20180116102854-5a71ef0e047d
github.com/aws/aws-sdk-go-v2 v1.41.4
diff --git a/go.sum b/go.sum
index 7f141140..5f890882 100644
--- a/go.sum
+++ b/go.sum
@@ -142,6 +142,8 @@ github.com/ProtonMail/go-crypto v1.0.0/go.mod h1:EjAoLdwvbIOoOQr3ihjnSoLZRtE8azu
github.com/RaveNoX/go-jsoncommentstrip v1.0.0/go.mod h1:78ihd09MekBnJnxpICcwzCMzGrKSKYe4AqU6PDYYpjk=
github.com/RussellLuo/slidingwindow v0.0.0-20200528002341-535bb99d338b h1:5/++qT1/z812ZqBvqQt6ToRswSuPZ/B33m6xVHRzADU=
github.com/RussellLuo/slidingwindow v0.0.0-20200528002341-535bb99d338b/go.mod h1:4+EPqMRApwwE/6yo6CxiHoSnBzjRr3jsqer7frxP8y4=
+github.com/SherClockHolmes/webpush-go v1.4.0 h1:ocnzNKWN23T9nvHi6IfyrQjkIc0oJWv1B1pULsf9i3s=
+github.com/SherClockHolmes/webpush-go v1.4.0/go.mod h1:XSq8pKX11vNV8MJEMwjrlTkxhAj1zKfxmyhdV7Pd6UA=
github.com/StackExchange/wmi v1.2.1 h1:VIkavFPXSjcnS+O8yTq7NI32k0R5Aj+v39y29VYDOSA=
github.com/StackExchange/wmi v1.2.1/go.mod h1:rcmrprowKIVzvc+NUiLncP2uuArMWLCbu9SBzvHz7e8=
github.com/VictoriaMetrics/fastcache v1.13.0 h1:AW4mheMR5Vd9FkAPUv+NH6Nhw+fmbTMGMsNAoA/+4G0=
@@ -568,6 +570,7 @@ github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69
github.com/golang-jwt/jwt/v4 v4.4.2/go.mod h1:m21LjoU+eqJr34lmDMbreY2eSTRJ1cv77w39/MY0Ch0=
github.com/golang-jwt/jwt/v4 v4.5.2 h1:YtQM7lnr8iZ+j5q71MGKkNw9Mn7AjHM68uc9g5fXeUI=
github.com/golang-jwt/jwt/v4 v4.5.2/go.mod h1:m21LjoU+eqJr34lmDMbreY2eSTRJ1cv77w39/MY0Ch0=
+github.com/golang-jwt/jwt/v5 v5.2.1/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk=
github.com/golang-jwt/jwt/v5 v5.3.0 h1:pv4AsKCKKZuqlgs5sUmn4x8UlGa0kEVt/puTpKx9vvo=
github.com/golang-jwt/jwt/v5 v5.3.0/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE=
github.com/golang-sql/civil v0.0.0-20220223132316-b832511892a9 h1:au07oEsX2xN0ktxqI+Sida1w446QrXBRJ0nee3SNZlA=
@@ -1372,8 +1375,6 @@ github.com/streamplace/atmoq/go v0.0.4-0.20260701223355-13757de4ae08 h1:NiTRz8AX
github.com/streamplace/atmoq/go v0.0.4-0.20260701223355-13757de4ae08/go.mod h1:3P8eSwKAGH7uh3SX5z1jlt/JgPTilJTUZngQJKhWY5s=
github.com/streamplace/atproto-oauth-golang v0.0.0-20260413212710-98956064d06c h1:IzEPU2O4iL58Nb7aw+7lB9ttnesEwOVVE5oV9NEXemM=
github.com/streamplace/atproto-oauth-golang v0.0.0-20260413212710-98956064d06c/go.mod h1:9LlKkqciiO5lRfbX0n4Wn5KNY9nvFb4R3by8FdW2TWc=
-github.com/streamplace/glex v0.0.0-20260715231618-ee553e32d7c7 h1:MSBBIH+QMR9AVfC0RuBLbBm/o1RAl7+bacek7txswDo=
-github.com/streamplace/glex v0.0.0-20260715231618-ee553e32d7c7/go.mod h1:LRaoeSMvSgOrhFX8s7ygjRlyka7wXdDa1s7JJ9o1IzY=
github.com/streamplace/glex v0.0.0-20260716203108-f73ed7cc31c9 h1:HbIhx8i7wytiNg9mWg2EuHG59r8FXTVlDqlkZObjqbQ=
github.com/streamplace/glex v0.0.0-20260716203108-f73ed7cc31c9/go.mod h1:LRaoeSMvSgOrhFX8s7ygjRlyka7wXdDa1s7JJ9o1IzY=
github.com/streamplace/go-dpop v0.0.0-20250510031900-c897158a8ad4 h1:L1fS4HJSaAyNnkwfuZubgfeZy8rkWmA0cMtH5Z0HqNc=
@@ -1604,6 +1605,7 @@ golang.org/x/crypto v0.13.0/go.mod h1:y6Z2r+Rw4iayiXXAIxJIDAJ1zMW4yaTpebo8fPOliY
golang.org/x/crypto v0.14.0/go.mod h1:MVFd36DqK4CsrnJYDkBA3VC4m2GkXAM0PvzMCn4JQf4=
golang.org/x/crypto v0.19.0/go.mod h1:Iy9bg/ha4yyC70EfRS8jz+B6ybOBKMaSxLj6P6oBDfU=
golang.org/x/crypto v0.23.0/go.mod h1:CKFgDieR+mRhux2Lsu27y0fO304Db0wZe70UKqHu0v8=
+golang.org/x/crypto v0.31.0/go.mod h1:kDsLvtWBEx7MV9tJOj9bnXsPbxwJQ6csT/x4KIN4Ssk=
golang.org/x/crypto v0.51.0 h1:IBPXwPfKxY7cWQZ38ZCIRPI50YLeevDLlLnyC5wRGTI=
golang.org/x/crypto v0.51.0/go.mod h1:8AdwkbraGNABw2kOX6YFPs3WM22XqI4EXEd8g+x7Oc8=
golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
@@ -1742,6 +1744,7 @@ golang.org/x/sync v0.3.0/go.mod h1:FU7BRWz2tNW+3quACPkgCx/L+uEAv1htQ0V83Z9Rj+Y=
golang.org/x/sync v0.4.0/go.mod h1:FU7BRWz2tNW+3quACPkgCx/L+uEAv1htQ0V83Z9Rj+Y=
golang.org/x/sync v0.6.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
golang.org/x/sync v0.7.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
+golang.org/x/sync v0.10.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4=
golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
@@ -1806,6 +1809,7 @@ golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.13.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.17.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/sys v0.20.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
+golang.org/x/sys v0.28.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY=
golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/telemetry v0.0.0-20240228155512-f48c80bd79b2/go.mod h1:TeRTkGYfJXctD9OcfyVLyj2J3IxLnKwHJR8f4D8a3YE=
@@ -1822,6 +1826,7 @@ golang.org/x/term v0.12.0/go.mod h1:owVbMEjm3cBLCHdkQu9b1opXd4ETQWc3BhuQGKgXgvU=
golang.org/x/term v0.13.0/go.mod h1:LTmsnFJwVN6bCy1rVCoS+qHT1HhALEFxKncY3WNNh4U=
golang.org/x/term v0.17.0/go.mod h1:lLRBjIVuehSbZlaOtGMbcMncT+aqLLLmKrsjNrUguwk=
golang.org/x/term v0.20.0/go.mod h1:8UkIAJTvZgivsXaD6/pH6U9ecQzZ45awqEOzuCvwpFY=
+golang.org/x/term v0.27.0/go.mod h1:iMsnZpn0cago0GOrHO2+Y7u7JPn5AylBrcoWkElMTSM=
golang.org/x/term v0.43.0 h1:S4RLU2sB31O/NCl+zFN9Aru9A/Cq2aqKpTZJ6B+DwT4=
golang.org/x/term v0.43.0/go.mod h1:lrhlHNdQJHO+1qVYiHfFKVuVioJIheAc3fBSMFYEIsk=
golang.org/x/text v0.0.0-20170915032832-14c0d48ead0c/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
@@ -1840,6 +1845,7 @@ golang.org/x/text v0.13.0/go.mod h1:TvPlkZtksWOMsz7fbANvkp4WM8x/WCo/om8BMLbz+aE=
golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
golang.org/x/text v0.15.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
golang.org/x/text v0.16.0/go.mod h1:GhwF1Be+LQoKShO3cGOHzqOgRrGaYc9AvblQOmPVHnI=
+golang.org/x/text v0.21.0/go.mod h1:4IBbMaMmOPCJ8SecivzSH54+73PCFmPWxNTLm+vZkEQ=
golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc=
golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38=
golang.org/x/time v0.0.0-20181108054448-85acf8d2951c/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
diff --git a/js/app/components/settings/notifications-category-settings.tsx b/js/app/components/settings/notifications-category-settings.tsx
new file mode 100644
index 00000000..9ef4f303
--- /dev/null
+++ b/js/app/components/settings/notifications-category-settings.tsx
@@ -0,0 +1,110 @@
+import {
+ MenuContainer,
+ MenuGroup,
+ Text,
+ View,
+ zero,
+} from "@streamplace/components";
+import { useEffect, useState } from "react";
+import { useTranslation } from "react-i18next";
+import { Platform, ScrollView } from "react-native";
+import { useStore } from "store";
+import { SettingToggle } from "./components/setting-toggle";
+
+// NotificationsCategorySettings is the opt-in surface for push notifications.
+//
+// On web, the toggle calls enableWebNotifications (which requests permission,
+// subscribes the browser's PushManager, and registers the subscription with the
+// backend) or disableWebNotifications (unsubscribes + prunes the server row).
+//
+// On native, push is driven by initPushNotifications at app start (FCM), so
+// the toggle reflects the OS permission state and, when supported, links the
+// user to the system settings to change it. The toggle itself can't grant
+// permission on native — the OS prompt already ran at startup — but it shows
+// the current state honestly.
+export function NotificationsCategorySettings() {
+ const { t } = useTranslation("settings");
+ const enableWebNotifications = useStore(
+ (state) => state.enableWebNotifications,
+ );
+ const disableWebNotifications = useStore(
+ (state) => state.disableWebNotifications,
+ );
+ const webNotificationPermission = useStore(
+ (state) => state.webNotificationPermission,
+ );
+ const notificationToken = useStore((state) => state.notificationToken);
+
+ const isWeb = Platform.OS === "web";
+ const [busy, setBusy] = useState(false);
+
+ // The "enabled" state: on web it's permission === granted AND we have a
+ // registered subscription token. On native it's whether the OS granted
+ // permission (the FCM token is acquired separately at startup).
+ const permission = isWeb ? webNotificationPermission() : "granted";
+ const enabled = isWeb
+ ? permission === "granted" && !!notificationToken
+ : permission === "granted";
+
+ // Re-check permission when the screen gains focus (the user may have
+ // changed it in browser settings).
+ const [refreshKey, setRefreshKey] = useState(0);
+ useEffect(() => {
+ if (!isWeb) return;
+ const interval = setInterval(() => setRefreshKey((k) => k + 1), 1000);
+ return () => clearInterval(interval);
+ }, [isWeb]);
+ // touch refreshKey so the linter doesn't complain and the re-render happens
+ void refreshKey;
+
+ const handleToggle = async (value: boolean) => {
+ if (busy) return;
+ setBusy(true);
+ try {
+ if (value) {
+ if (isWeb) {
+ await enableWebNotifications();
+ }
+ // native: permission was already requested at startup; nothing to do
+ } else {
+ if (isWeb) {
+ await disableWebNotifications();
+ }
+ }
+ } finally {
+ setBusy(false);
+ }
+ };
+
+ return (
+
+
+
+
+
+
+
+ {isWeb && permission === "denied" && (
+
+
+ {t("notifications-blocked-help")}
+
+
+ )}
+
+
+
+
+ );
+}
diff --git a/js/app/components/settings/settings.tsx b/js/app/components/settings/settings.tsx
index e778534c..9fd245f4 100644
--- a/js/app/components/settings/settings.tsx
+++ b/js/app/components/settings/settings.tsx
@@ -18,6 +18,7 @@ import {
import { ImageBackground } from "expo-image";
import {
Award,
+ Bell,
Brush,
Globe,
Info,
@@ -141,6 +142,12 @@ export function Settings() {
screen="PrivacyCategory"
icon={Shield}
/>
+
+
)}
{danmuUnlocked && (
diff --git a/js/app/features/platform/shared.tsx b/js/app/features/platform/shared.tsx
index a3b086fd..c4aa024c 100644
--- a/js/app/features/platform/shared.tsx
+++ b/js/app/features/platform/shared.tsx
@@ -13,4 +13,5 @@ export const initialState: PlatformState = {
export type RegisterNotificationTokenBody = {
token: string;
repoDID?: string;
+ type?: "firebase" | "web";
};
diff --git a/js/app/public/sw.js b/js/app/public/sw.js
new file mode 100644
index 00000000..b4c5b98c
--- /dev/null
+++ b/js/app/public/sw.js
@@ -0,0 +1,70 @@
+// Streamplace service worker for Web Push.
+//
+// This file is served from /sw.js (see public/) and registered by the web
+// platform slice on app mount. It does two things:
+//
+// 1. push — when a push message arrives, show it as a system notification.
+// The payload is the JSON-serialized NotificationBlast from the server
+// ({title, body, data}). If the payload is empty (a "silent" push), we
+// still show a notification with a generic title, since browsers require
+// a visible notification for every push.
+//
+// 2. notificationclick — focus an existing tab (or open a new one) and
+// navigate it to the path encoded in the notification's data, so tapping
+// a "🔴 @user is LIVE!" notification opens that stream.
+
+self.addEventListener("push", (event) => {
+ let data = { title: "Streamplace", body: "", data: {} };
+ try {
+ if (event.data) {
+ const parsed = event.data.json();
+ data = {
+ title: parsed.title || data.title,
+ body: parsed.body || data.body,
+ data: parsed.data || data.data,
+ };
+ }
+ } catch (e) {
+ // Payload wasn't JSON; fall back to raw text if present.
+ if (event.data) {
+ data.body = event.data.text();
+ }
+ }
+ event.waitUntil(
+ self.registration.showNotification(data.title, {
+ body: data.body,
+ data: data.data,
+ icon: "/favicon.ico",
+ badge: "/favicon.ico",
+ }),
+ );
+});
+
+self.addEventListener("notificationclick", (event) => {
+ event.notification.close();
+ const path = event.notification.data && event.notification.data.path;
+ const targetUrl = path ? path : "/";
+
+ event.waitUntil(
+ (async () => {
+ const all = await self.clients.matchAll({
+ type: "window",
+ includeUncontrolled: true,
+ });
+ // Focus an existing tab if one is open.
+ for (const client of all) {
+ if ("focus" in client) {
+ client.focus();
+ if ("navigate" in client) {
+ await client.navigate(targetUrl);
+ }
+ return;
+ }
+ }
+ // Otherwise open a new one.
+ if (self.clients.openWindow) {
+ await self.clients.openWindow(targetUrl);
+ }
+ })(),
+ );
+});
diff --git a/js/app/src/linking-config.ts b/js/app/src/linking-config.ts
index cdfac859..3cead5a3 100644
--- a/js/app/src/linking-config.ts
+++ b/js/app/src/linking-config.ts
@@ -121,6 +121,7 @@ export const SCREEN_PATHS = {
WebhooksSettings: "settings/streaming/webhooks",
RecommendationsSettings: "settings/streaming/recommendations",
PrivacyCategory: "settings/privacy",
+ NotificationsCategory: "settings/notifications",
DanmuCategory: "settings/danmu",
AdvancedCategory: "settings/advanced",
DeveloperSettings: "settings/developer",
diff --git a/js/app/src/navigation-types.ts b/js/app/src/navigation-types.ts
index b3712f43..431aefc1 100644
--- a/js/app/src/navigation-types.ts
+++ b/js/app/src/navigation-types.ts
@@ -9,6 +9,7 @@ export type SettingsStackParamList = {
WebhooksSettings: undefined;
BackupSettings: undefined;
PrivacyCategory: undefined;
+ NotificationsCategory: undefined;
DanmuCategory: undefined;
AdvancedCategory: undefined;
LanguagesCategory: undefined;
diff --git a/js/app/src/shell.tsx b/js/app/src/shell.tsx
index ed990807..927defc0 100644
--- a/js/app/src/shell.tsx
+++ b/js/app/src/shell.tsx
@@ -33,6 +33,7 @@ import { DanmuCategorySettings } from "components/settings/danmu-category-settin
import KeyManager from "components/settings/key-manager";
import { LanguagesCategorySettings } from "components/settings/languages-category-settings";
import MultistreamManager from "components/settings/multistream-manager";
+import { NotificationsCategorySettings } from "components/settings/notifications-category-settings";
import { PrivacyCategorySettings } from "components/settings/privacy-category-settings";
import RecommendationsManager from "components/settings/recommendations-manager";
import { StreamingCategorySettings } from "components/settings/streaming-category-settings";
@@ -369,6 +370,11 @@ function SettingsNavigator() {
component={PrivacyCategorySettings}
options={{ title: "Privacy & Security" }}
/>
+
useStore((state) => state.notificationToken);
export const useNotificationDestination = () =>
useStore((state) => state.notificationDestination);
+export const useEnableWebNotifications = () =>
+ useStore((state) => state.enableWebNotifications);
+export const useDisableWebNotifications = () =>
+ useStore((state) => state.disableWebNotifications);
+export const useWebNotificationPermission = () =>
+ useStore((state) => state.webNotificationPermission);
diff --git a/js/app/store/slices/platformSlice.native.ts b/js/app/store/slices/platformSlice.native.ts
index abb1ebef..860b5f47 100644
--- a/js/app/store/slices/platformSlice.native.ts
+++ b/js/app/store/slices/platformSlice.native.ts
@@ -14,6 +14,12 @@ export interface PlatformSlice {
openLoginLink: (url: string) => Promise;
initPushNotifications: () => Promise;
registerNotificationToken: () => Promise;
+ // web-only actions; no-ops on native. Present so the shared PlatformSlice
+ // type is identical across platforms and the settings toggle can call
+ // them unconditionally.
+ enableWebNotifications: () => Promise;
+ disableWebNotifications: () => Promise;
+ webNotificationPermission: () => NotificationPermission;
}
const checkApplicationPermission = async () => {
@@ -170,7 +176,10 @@ export const createPlatformSlice: StateCreator<
return;
}
- const body: { token: string; repoDID?: string } = { token };
+ const body: { token: string; type: string; repoDID?: string } = {
+ token,
+ type: "firebase",
+ };
const did = oauthSession?.did;
if (did) {
@@ -196,4 +205,14 @@ export const createPlatformSlice: StateCreator<
console.error("registerNotificationToken error", e);
}
},
+ enableWebNotifications: async () => {
+ // web-only; native uses FCM via initPushNotifications
+ return "denied";
+ },
+ disableWebNotifications: async () => {
+ // web-only
+ },
+ webNotificationPermission: () => {
+ return "denied";
+ },
});
diff --git a/js/app/store/slices/platformSlice.ts b/js/app/store/slices/platformSlice.ts
index cbcf76f6..75e4b284 100644
--- a/js/app/store/slices/platformSlice.ts
+++ b/js/app/store/slices/platformSlice.ts
@@ -1,3 +1,4 @@
+import { AppStore } from "store";
import { StateCreator } from "zustand";
export interface PlatformSlice {
@@ -10,14 +11,39 @@ export interface PlatformSlice {
openLoginLink: (url: string) => Promise;
initPushNotifications: () => Promise;
registerNotificationToken: () => Promise;
+ // web-only: subscribe/unsubscribe the browser's PushManager. Returns the
+ // permission state so the settings toggle can reflect reality.
+ enableWebNotifications: () => Promise;
+ disableWebNotifications: () => Promise;
+ webNotificationPermission: () => NotificationPermission;
}
-export const createPlatformSlice: StateCreator = (set, get) => ({
+// VAPID public key must be converted from base64url to a Uint8Array for the
+// PushManager.subscribe() applicationServerKey argument.
+function urlBase64ToUint8Array(base64String: string): Uint8Array {
+ const padding = "=".repeat((4 - (base64String.length % 4)) % 4);
+ const base64 = (base64String + padding).replace(/-/g, "+").replace(/_/g, "/");
+ const rawData = atob(base64);
+ const output = new Uint8Array(rawData.length);
+ for (let i = 0; i < rawData.length; ++i) {
+ output[i] = rawData.charCodeAt(i);
+ }
+ return output;
+}
+
+export const createPlatformSlice: StateCreator<
+ AppStore,
+ [],
+ [],
+ PlatformSlice
+> = (set, get) => ({
status: "idle",
notificationToken: null,
notificationDestination: null,
handleNotification: (payload) => {
- // notification handling logic
+ if (!payload) return;
+ if (typeof payload.path !== "string") return;
+ set({ notificationDestination: payload.path });
},
clearNotification: () => {
set({ notificationDestination: null });
@@ -32,9 +58,112 @@ export const createPlatformSlice: StateCreator = (set, get) => ({
}
},
initPushNotifications: async () => {
- // mobile-only, web notifications someday
+ // Register the service worker that receives push events. This must
+ // happen early (on app mount) so that pushes delivered while the tab is
+ // backgrounded still surface as system notifications. The actual
+ // subscription + permission request is deferred to the settings toggle
+ // (enableWebNotifications) because browsers require a user gesture for
+ // the permission prompt.
+ if (!("serviceWorker" in navigator)) {
+ return;
+ }
+ try {
+ await navigator.serviceWorker.register("/sw.js");
+ } catch (e) {
+ console.log("service worker registration failed", e);
+ }
},
registerNotificationToken: async () => {
- // notification token registration
+ // On web, token registration is driven by enableWebNotifications (which
+ // subscribes and posts the subscription). This no-op keeps the shared
+ // shell effect happy without double-registering.
+ },
+ enableWebNotifications: async () => {
+ const url = get().url;
+ if (!url) {
+ console.log(
+ "no streamplace url configured, cannot enable web notifications",
+ );
+ return "denied";
+ }
+ try {
+ const permission = await Notification.requestPermission();
+ if (permission !== "granted") {
+ return permission;
+ }
+
+ // Make sure the service worker is active before subscribing.
+ const reg = await navigator.serviceWorker.ready;
+
+ // Fetch the server's VAPID public key.
+ const vapidRes = await fetch(`${url}/api/notification/vapid-public-key`);
+ if (!vapidRes.ok) {
+ throw new Error(`failed to fetch vapid public key: ${vapidRes.status}`);
+ }
+ const { publicKey } = await vapidRes.json();
+
+ // Subscribe the browser's PushManager.
+ const subscription = await reg.pushManager.subscribe({
+ userVisibleOnly: true,
+ applicationServerKey: urlBase64ToUint8Array(publicKey) as BufferSource,
+ });
+ const subJSON = JSON.stringify(subscription);
+ set({ notificationToken: subJSON });
+
+ // Register the subscription with the backend.
+ const { oauthSession } = get();
+ const body: { token: string; type: string; repoDID?: string } = {
+ token: subJSON,
+ type: "web",
+ };
+ if (oauthSession?.did) {
+ body.repoDID = oauthSession.did;
+ }
+ const res = await fetch(`${url}/api/notification`, {
+ method: "POST",
+ headers: { "content-type": "application/json" },
+ body: JSON.stringify(body),
+ });
+ console.log("web notification registration status:", res.status);
+ return permission;
+ } catch (e) {
+ console.error("enableWebNotifications error", e);
+ return "denied";
+ }
+ },
+ disableWebNotifications: async () => {
+ const url = get().url;
+ const { notificationToken } = get();
+ try {
+ if (notificationToken) {
+ // Unsubscribe the browser side so it stops accepting pushes.
+ const sub = JSON.parse(notificationToken);
+ // We need the live PushSubscription object to call unsubscribe(); get
+ // it from the service worker registration by matching endpoint.
+ const reg = await navigator.serviceWorker.ready;
+ const existing = await reg.pushManager.getSubscription();
+ if (existing && existing.endpoint === sub.endpoint) {
+ await existing.unsubscribe();
+ }
+ // Tell the server to drop the row.
+ if (url) {
+ await fetch(`${url}/api/notification`, {
+ method: "DELETE",
+ headers: { "content-type": "application/json" },
+ body: JSON.stringify({ token: notificationToken }),
+ });
+ }
+ }
+ } catch (e) {
+ console.error("disableWebNotifications error", e);
+ } finally {
+ set({ notificationToken: null });
+ }
+ },
+ webNotificationPermission: () => {
+ if (typeof Notification === "undefined") {
+ return "denied";
+ }
+ return Notification.permission;
},
});
diff --git a/js/components/locales/en-US/settings.ftl b/js/components/locales/en-US/settings.ftl
index 3711dcdc..9f522a9c 100644
--- a/js/components/locales/en-US/settings.ftl
+++ b/js/components/locales/en-US/settings.ftl
@@ -55,6 +55,12 @@ developer = Developer
languages = Languages
privacy-security = Privacy & Security
streaming = Streaming
+notifications = Notifications
+notifications-title = Push Notifications
+notifications-web-description = Get notified when streamers you follow go live.
+notifications-mobile-description = Push notifications are managed by your device settings.
+notifications-blocked-description = Blocked — you denied notification permission in your browser.
+notifications-blocked-help = To re-enable, update site permissions in your browser settings, then toggle this on again.
## Common Actions
cancel = Cancel
diff --git a/js/components/locales/es-ES/settings.ftl b/js/components/locales/es-ES/settings.ftl
index 21db93e9..027991e0 100644
--- a/js/components/locales/es-ES/settings.ftl
+++ b/js/components/locales/es-ES/settings.ftl
@@ -91,6 +91,7 @@ developer = Desarrollador
languages = Idiomas
privacy-security = Privacidad y Seguridad
streaming = Transmisión
+notifications = Notificaciones
## Acciones Comunes
cancel = Cancelar
diff --git a/js/components/locales/fr-FR/settings.ftl b/js/components/locales/fr-FR/settings.ftl
index 2f338953..b38df888 100644
--- a/js/components/locales/fr-FR/settings.ftl
+++ b/js/components/locales/fr-FR/settings.ftl
@@ -89,6 +89,7 @@ developer = Développeur
languages = Langues
privacy-security = Confidentialité et Sécurité
streaming = Diffusion
+notifications = Notifications
## Actions Courantes
cancel = Annuler
diff --git a/js/components/locales/pt-BR/settings.ftl b/js/components/locales/pt-BR/settings.ftl
index 08712eeb..253bf6b8 100644
--- a/js/components/locales/pt-BR/settings.ftl
+++ b/js/components/locales/pt-BR/settings.ftl
@@ -89,6 +89,7 @@ developer = Desenvolvedor
languages = Idiomas
privacy-security = Privacidade e Segurança
streaming = Transmissão
+notifications = Notificações
## Ações Comuns
cancel = Cancelar
diff --git a/js/components/locales/ro-RO/settings.ftl b/js/components/locales/ro-RO/settings.ftl
index 6fa630eb..a58dcb53 100644
--- a/js/components/locales/ro-RO/settings.ftl
+++ b/js/components/locales/ro-RO/settings.ftl
@@ -51,6 +51,7 @@ developer = Dezvoltator
languages = Limbi
privacy-security = Confidențialitate și securitate
streaming = Streaming
+notifications = Notificări
## Common Actions
cancel = Anulare
diff --git a/js/components/locales/zh-Hans/settings.ftl b/js/components/locales/zh-Hans/settings.ftl
index 9235e910..a6041461 100644
--- a/js/components/locales/zh-Hans/settings.ftl
+++ b/js/components/locales/zh-Hans/settings.ftl
@@ -53,6 +53,7 @@ developer = 开发者
languages = 语言
privacy-security = 隐私与安全
streaming = 串流
+notifications = 通知
## Common Actions
cancel = 取消
diff --git a/js/components/locales/zh-Hant/settings.ftl b/js/components/locales/zh-Hant/settings.ftl
index f0522c9c..47446a79 100644
--- a/js/components/locales/zh-Hant/settings.ftl
+++ b/js/components/locales/zh-Hant/settings.ftl
@@ -90,6 +90,7 @@ developer = 開發者
languages = 語言
privacy-security = 隱私與安全
streaming = 串流
+notifications = 通知
## 常用動作
cancel = 取消
diff --git a/pkg/api/api.go b/pkg/api/api.go
index 9309331d..331ab77c 100644
--- a/pkg/api/api.go
+++ b/pkg/api/api.go
@@ -64,7 +64,7 @@ type StreamplaceAPI struct {
Updater *Updater
Signer *eip712.EIP712Signer
Mimes map[string]string
- FirebaseNotifier notifications.FirebaseNotifier
+ Notifier notifications.Notifier
MediaManager *media.MediaManager
MediaSigner media.MediaSigner
UploadManager *upload.Manager
@@ -99,7 +99,7 @@ type WebsocketTracker struct {
mu sync.RWMutex
}
-func MakeStreamplaceAPI(cli *config.CLI, mod model.Model, statefulDB *statedb.StatefulDB, noter notifications.FirebaseNotifier, mm *media.MediaManager, ms media.MediaSigner, bus *bus.Bus, atsync *atproto.ATProtoSynchronizer, d *director.Director, op *oatproxy.OATProxy, ldb localdb.LocalDB, um *upload.Manager, playbackStore blob.Store, viewLog *viewlog.Writer) (*StreamplaceAPI, error) {
+func MakeStreamplaceAPI(cli *config.CLI, mod model.Model, statefulDB *statedb.StatefulDB, noter notifications.Notifier, mm *media.MediaManager, ms media.MediaSigner, bus *bus.Bus, atsync *atproto.ATProtoSynchronizer, d *director.Director, op *oatproxy.OATProxy, ldb localdb.LocalDB, um *upload.Manager, playbackStore blob.Store, viewLog *viewlog.Writer) (*StreamplaceAPI, error) {
updater, err := PrepareUpdater(cli)
if err != nil {
return nil, err
@@ -108,7 +108,7 @@ func MakeStreamplaceAPI(cli *config.CLI, mod model.Model, statefulDB *statedb.St
Model: mod,
StatefulDB: statefulDB,
Updater: updater,
- FirebaseNotifier: noter,
+ Notifier: noter,
MediaManager: mm,
MediaSigner: ms,
UploadManager: um,
@@ -188,6 +188,8 @@ func (a *StreamplaceAPI) Handler(ctx context.Context) (http.Handler, error) {
router.Handler("GET", "/.well-known/assetlinks.json", a.HandleAndroidAssetLinks(ctx))
apiRouter := httprouter.New()
addFunc(apiRouter, "POST", "/api/notification", a.HandleNotification(ctx))
+ addFunc(apiRouter, "DELETE", "/api/notification", a.HandleNotificationDelete(ctx))
+ addFunc(apiRouter, "GET", "/api/notification/vapid-public-key", a.HandleVapidPublicKey(ctx))
// old clients
addFunc(router, "GET", "/app-updates", a.HandleAppUpdates(ctx))
// new ones
@@ -556,6 +558,7 @@ func (a *StreamplaceAPI) RedirectHandler(ctx context.Context) (http.Handler, err
type NotificationPayload struct {
Token string `json:"token"`
RepoDID string `json:"repoDID"`
+ Type string `json:"type"`
}
func (a *StreamplaceAPI) HandleAPI404(ctx context.Context) http.HandlerFunc {
@@ -579,7 +582,7 @@ func (a *StreamplaceAPI) HandleNotification(ctx context.Context) http.HandlerFun
w.WriteHeader(400)
return
}
- err = a.StatefulDB.CreateNotification(n.Token, n.RepoDID)
+ err = a.StatefulDB.CreateNotification(n.Token, n.RepoDID, statedb.NotificationType(n.Type))
if err != nil {
log.Log(ctx, "error creating notification", "error", err)
w.WriteHeader(400)
@@ -598,6 +601,56 @@ func (a *StreamplaceAPI) HandleNotification(ctx context.Context) http.HandlerFun
}
}
+// HandleNotificationDelete removes a push token (web or mobile). Used by the
+// web client when a user disables notifications — the browser subscription is
+// unsubscribed locally and the server row is pruned so we stop pushing to a
+// dead endpoint.
+func (a *StreamplaceAPI) HandleNotificationDelete(ctx context.Context) http.HandlerFunc {
+ return func(w http.ResponseWriter, req *http.Request) {
+ payload, err := io.ReadAll(req.Body)
+ if err != nil {
+ log.Log(ctx, "error reading notification delete", "error", err)
+ w.WriteHeader(400)
+ return
+ }
+ n := NotificationPayload{}
+ if err := json.Unmarshal(payload, &n); err != nil {
+ log.Log(ctx, "error unmarshalling notification delete", "error", err)
+ w.WriteHeader(400)
+ return
+ }
+ if n.Token == "" {
+ w.WriteHeader(400)
+ return
+ }
+ if err := a.StatefulDB.DeleteNotification(n.Token); err != nil {
+ log.Log(ctx, "error deleting notification", "error", err)
+ w.WriteHeader(400)
+ return
+ }
+ log.Log(ctx, "successfully deleted notification", "token", n.Token)
+ w.WriteHeader(200)
+ }
+}
+
+// HandleVapidPublicKey returns the server's Web Push VAPID public key. The web
+// client needs it to subscribe the browser's PushManager. The key is generated
+// on first access (via EnsureVAPIDKeys) and stays stable thereafter.
+func (a *StreamplaceAPI) HandleVapidPublicKey(ctx context.Context) http.HandlerFunc {
+ return func(w http.ResponseWriter, req *http.Request) {
+ keys, err := a.StatefulDB.EnsureVAPIDKeys(ctx)
+ if err != nil {
+ log.Error(ctx, "error ensuring vapid keys", "error", err)
+ apierrors.WriteHTTPInternalServerError(w, "unable to get vapid public key", err)
+ return
+ }
+ w.Header().Set("Content-Type", "application/json")
+ if _, err := w.Write([]byte(`{"publicKey":"` + keys.PublicKey + `"}`)); err != nil {
+ log.Error(ctx, "error writing vapid public key", "error", err)
+ }
+ }
+}
+
func (a *StreamplaceAPI) HandleSegment(ctx context.Context) http.HandlerFunc {
return func(w http.ResponseWriter, req *http.Request) {
err := a.MediaManager.ValidateMP4(ctx, req.Body, false)
diff --git a/pkg/api/api_internal.go b/pkg/api/api_internal.go
index 1b92b862..7f0cd158 100644
--- a/pkg/api/api_internal.go
+++ b/pkg/api/api_internal.go
@@ -430,15 +430,15 @@ func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, err
errors.WriteHTTPInternalServerError(w, "unable to get notifications", err)
return
}
- if a.FirebaseNotifier == nil {
- errors.WriteHTTPInternalServerError(w, "firebase notifier not initialized", nil)
+ if a.Notifier == nil {
+ errors.WriteHTTPInternalServerError(w, "notifier not initialized", nil)
return
}
- tokens := []string{}
- for _, not := range notifications {
- tokens = append(tokens, not.Token)
+ targets := make([]notificationpkg.NotificationTarget, len(notifications))
+ for i, not := range notifications {
+ targets[i] = notificationpkg.NotificationTarget{Token: not.Token, Type: not.Type}
}
- err = a.FirebaseNotifier.Blast(ctx, tokens, &payload)
+ err = a.Notifier.Blast(ctx, targets, &payload)
if err != nil {
errors.WriteHTTPInternalServerError(w, "unable to blast notifications", err)
return
diff --git a/pkg/api/api_test.go b/pkg/api/api_test.go
index 48bc12f9..c68276a5 100644
--- a/pkg/api/api_test.go
+++ b/pkg/api/api_test.go
@@ -81,7 +81,7 @@ func TestRedirectHandler(t *testing.T) {
type MockFirebase struct {
}
-func (m *MockFirebase) Blast(ctx context.Context, nots []string, nb *notifications.NotificationBlast) error {
+func (m *MockFirebase) Blast(ctx context.Context, targets []notifications.NotificationTarget, nb *notifications.NotificationBlast) error {
return nil
}
diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go
index a54c3cf5..aed0aa0a 100644
--- a/pkg/atproto/firehose.go
+++ b/pkg/atproto/firehose.go
@@ -47,7 +47,7 @@ type ATProtoSynchronizer struct {
CLI *config.CLI
Model model.Model
StatefulDB *statedb.StatefulDB
- Noter notificationpkg.FirebaseNotifier
+ Noter notificationpkg.Notifier
Bus *bus.Bus
PLCDirectory identity.Directory
CachedPLCDirectory identity.Directory
diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go
index 84996d55..401a6b3f 100644
--- a/pkg/atproto/sync.go
+++ b/pkg/atproto/sync.go
@@ -965,21 +965,25 @@ func (atsync *ATProtoSynchronizer) notifyBetaInvite(ctx context.Context, rec *pl
if atsync.Noter == nil || atsync.StatefulDB == nil {
return
}
- tokens, err := atsync.StatefulDB.GetManyNotificationTokens([]string{rec.Did})
+ notifications, err := atsync.StatefulDB.GetManyNotifications([]string{rec.Did})
if err != nil {
log.Error(ctx, "beta invite notification: failed to load tokens", "did", rec.Did, "err", err)
return
}
- if len(tokens) == 0 {
+ if len(notifications) == 0 {
log.Debug(ctx, "beta invite notification: no device tokens for invitee", "did", rec.Did, "feature", rec.Feature)
return
}
blast := betaInviteBlast(rec.Feature)
- if err := atsync.Noter.Blast(ctx, tokens, blast); err != nil {
+ targets := make([]notificationpkg.NotificationTarget, len(notifications))
+ for i, n := range notifications {
+ targets[i] = notificationpkg.NotificationTarget{Token: n.Token, Type: n.Type}
+ }
+ if err := atsync.Noter.Blast(ctx, targets, blast); err != nil {
log.Error(ctx, "beta invite notification: blast failed", "did", rec.Did, "feature", rec.Feature, "err", err)
return
}
- log.Log(ctx, "sent beta invite notification", "did", rec.Did, "feature", rec.Feature, "tokens", len(tokens))
+ log.Log(ctx, "sent beta invite notification", "did", rec.Did, "feature", rec.Feature, "tokens", len(notifications))
}
// betaInviteBlast builds the push payload for a newly-granted beta feature.
diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go
index fcdfdde9..f268c43e 100644
--- a/pkg/cmd/streamplace.go
+++ b/pkg/cmd/streamplace.go
@@ -198,9 +198,9 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu
if err != nil {
return err
}
- var noter notifications.FirebaseNotifier
+ var fbNotifier notifications.FirebaseNotifier
if cli.FirebaseServiceAccount != "" {
- noter, err = notifications.MakeFirebaseNotifier(ctx, cli.FirebaseServiceAccount)
+ fbNotifier, err = notifications.MakeFirebaseNotifier(ctx, cli.FirebaseServiceAccount)
if err != nil {
return err
}
@@ -213,10 +213,22 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu
if err != nil {
return err
}
- state, err := statedb.MakeDB(ctx, cli, noter, mod)
+ // The notifier is assembled after the DB exists, because the Web Push
+ // notifier needs VAPID keys that are persisted in the Config table. The
+ // queue processor nil-checks the notifier, so the brief window is safe.
+ state, err := statedb.MakeDB(ctx, cli, nil, mod)
if err != nil {
return err
}
+
+ // Build the Web Push notifier from VAPID keys generated/stored in the DB.
+ vapidKeys, err := state.EnsureVAPIDKeys(ctx)
+ if err != nil {
+ return err
+ }
+ webNotifier := notifications.NewWebPushNotifier(vapidKeys, "")
+ noter := notifications.NewMultiNotifier(fbNotifier, webNotifier)
+ state.SetNotifier(noter)
handle, err := atproto.MakeLexiconRepo(ctx, cli, mod, state)
if err != nil {
return err
diff --git a/pkg/notifications/firebase.go b/pkg/notifications/firebase.go
index 2dd16740..3bac692b 100644
--- a/pkg/notifications/firebase.go
+++ b/pkg/notifications/firebase.go
@@ -1,20 +1,23 @@
package notifications
import (
+ "context"
"encoding/base64"
"encoding/json"
"fmt"
- "context"
-
firebase "firebase.google.com/go/v4"
"firebase.google.com/go/v4/messaging"
"google.golang.org/api/option"
"stream.place/streamplace/pkg/log"
)
+// FirebaseNotifier sends pushes via Firebase Cloud Messaging (FCM/APNs). It
+// implements Notifier by handling targets of NotificationTypeFirebase and
+// ignoring all others.
type FirebaseNotifier interface {
- Blast(ctx context.Context, tokens []string, golive *NotificationBlast) error
+ Notifier
+ BlastTokens(ctx context.Context, tokens []string, blast *NotificationBlast) error
}
type FirebaseNotifierS struct {
@@ -55,8 +58,23 @@ func MakeFirebaseNotifier(ctx context.Context, serviceAccountJSONb64 string) (Fi
return &FirebaseNotifierS{app: app}, nil
}
+// Blast implements Notifier. It filters targets down to firebase tokens and
+// delegates to BlastTokens; web targets are left for another notifier.
+func (f *FirebaseNotifierS) Blast(ctx context.Context, targets []NotificationTarget, blast *NotificationBlast) error {
+ tokens := make([]string, 0, len(targets))
+ for _, t := range targets {
+ if t.Type == NotificationTypeFirebase || t.Type == "" {
+ tokens = append(tokens, t.Token)
+ }
+ }
+ if len(tokens) == 0 {
+ return nil
+ }
+ return f.BlastTokens(ctx, tokens, blast)
+}
+
// refactor me when we have >500 users
-func (f *FirebaseNotifierS) Blast(ctx context.Context, tokens []string, blast *NotificationBlast) error {
+func (f *FirebaseNotifierS) BlastTokens(ctx context.Context, tokens []string, blast *NotificationBlast) error {
client, err := f.app.Messaging(ctx)
if err != nil {
return err
diff --git a/pkg/notifications/multi.go b/pkg/notifications/multi.go
new file mode 100644
index 00000000..ae008dee
--- /dev/null
+++ b/pkg/notifications/multi.go
@@ -0,0 +1,52 @@
+package notifications
+
+import (
+ "context"
+ "fmt"
+)
+
+// MultiNotifier fans a single blast out to every transport it wraps. It
+// implements Notifier by delegating each target to the child notifier that
+// handles that target's Type. This is the single Notifier the rest of the
+// codebase holds (replacing the old FirebaseNotifier field), so callers don't
+// need to know which transports are configured.
+//
+// A child that returns an error for its slice of targets does not abort the
+// others — each transport is independent, and a dead FCM credential
+// shouldn't block web pushes (or vice versa). Errors are collected and
+// returned as a joined error after all transports have been attempted.
+type MultiNotifier struct {
+ notifiers []Notifier
+}
+
+// NewMultiNotifier wraps one or more transport notifiers. nil entries are
+// skipped so callers can pass a possibly-unconfigured notifier through
+// without filtering.
+func NewMultiNotifier(notifiers ...Notifier) *MultiNotifier {
+ nn := &MultiNotifier{}
+ for _, n := range notifiers {
+ if n != nil {
+ nn.notifiers = append(nn.notifiers, n)
+ }
+ }
+ return nn
+}
+
+func (m *MultiNotifier) Blast(ctx context.Context, targets []NotificationTarget, blast *NotificationBlast) error {
+ if len(m.notifiers) == 0 {
+ return nil
+ }
+ var errs []error
+ for _, n := range m.notifiers {
+ if err := n.Blast(ctx, targets, blast); err != nil {
+ errs = append(errs, err)
+ }
+ }
+ if len(errs) == 1 {
+ return errs[0]
+ }
+ if len(errs) > 1 {
+ return fmt.Errorf("multi-notifier: %d transports failed: %v", len(errs), errs)
+ }
+ return nil
+}
diff --git a/pkg/notifications/notifier.go b/pkg/notifications/notifier.go
new file mode 100644
index 00000000..3d1c019f
--- /dev/null
+++ b/pkg/notifications/notifier.go
@@ -0,0 +1,34 @@
+package notifications
+
+import "context"
+
+// NotificationType identifies the push transport a token belongs to. It lives
+// in pkg/notifications (the lower-level package) so that pkg/statedb — which
+// already imports pkg/notifications for the Notifier field — can reference it
+// on the Notification row without creating a circular import.
+type NotificationType string
+
+const (
+ // NotificationTypeFirebase is an FCM/APNs registration token (mobile).
+ NotificationTypeFirebase NotificationType = "firebase"
+ // NotificationTypeWeb is a Web Push subscription (endpoint + p256dh/auth
+ // keys), stored as the JSON-serialized PushSubscription object.
+ NotificationTypeWeb NotificationType = "web"
+)
+
+// NotificationTarget pairs a push token with the transport that knows how to
+// deliver to it. Blast callers produce a []NotificationTarget (typically by
+// loading notification rows from the DB) and hand them to a Notifier, which
+// fans each target out to the matching transport.
+type NotificationTarget struct {
+ Token string
+ Type NotificationType
+}
+
+// Notifier sends a notification blast to a set of targets. Each
+// implementation handles one transport (firebase, web) or fans out across
+// several (MultiNotifier). Implementations should silently skip targets whose
+// Type they don't handle.
+type Notifier interface {
+ Blast(ctx context.Context, targets []NotificationTarget, blast *NotificationBlast) error
+}
diff --git a/pkg/notifications/webpush.go b/pkg/notifications/webpush.go
new file mode 100644
index 00000000..ba5adbaa
--- /dev/null
+++ b/pkg/notifications/webpush.go
@@ -0,0 +1,135 @@
+package notifications
+
+import (
+ "context"
+ "encoding/json"
+ "fmt"
+ "sync"
+
+ webpush "github.com/SherClockHolmes/webpush-go"
+ "stream.place/streamplace/pkg/log"
+)
+
+// VAPIDKeys is the ECDSA P-256 application-server keypair Web Push requires.
+// The public key is handed to the browser (it validates pushes against it);
+// the private key signs the VAPID JWT on each send. Keys must stay stable —
+// rotating them invalidates every existing browser subscription.
+type VAPIDKeys struct {
+ PublicKey string `json:"publicKey"`
+ PrivateKey string `json:"privateKey"`
+}
+
+// WebPushNotifier sends pushes via the Web Push protocol (RFC 8291 + VAPID).
+// It implements Notifier by handling targets of NotificationTypeWeb and
+// ignoring all others. Each target's Token is the JSON-serialized
+// PushSubscription object the browser produced.
+type WebPushNotifier struct {
+ keys VAPIDKeys
+ // subscriber is the contact URI embedded in the VAPID JWT (RFC 8291
+ // "sub"). mailto: is conventional; a URL works too. Browsers ignore it
+ // for delivery but it's required by the spec.
+ subscriber string
+}
+
+// NewWebPushNotifier builds a notifier from a VAPID keypair. The subscriber
+// defaults to a mailto: if empty.
+func NewWebPushNotifier(keys VAPIDKeys, subscriber string) *WebPushNotifier {
+ if subscriber == "" {
+ subscriber = "mailto:noreply@stream.place"
+ }
+ return &WebPushNotifier{keys: keys, subscriber: subscriber}
+}
+
+// Blast implements Notifier. It fans a push out to every web target in
+// parallel (each subscription is an independent HTTP POST to the browser
+// push service). Firebase targets are ignored.
+func (w *WebPushNotifier) Blast(ctx context.Context, targets []NotificationTarget, blast *NotificationBlast) error {
+ webTargets := make([]NotificationTarget, 0, len(targets))
+ for _, t := range targets {
+ if t.Type == NotificationTypeWeb {
+ webTargets = append(webTargets, t)
+ }
+ }
+ if len(webTargets) == 0 {
+ return nil
+ }
+
+ payload, err := json.Marshal(blast)
+ if err != nil {
+ return fmt.Errorf("error marshaling notification blast: %w", err)
+ }
+
+ var (
+ wg sync.WaitGroup
+ mu sync.Mutex
+ success int
+ failed int
+ errs []error
+ )
+
+ for _, t := range webTargets {
+ wg.Add(1)
+ go func(token string) {
+ defer wg.Done()
+ err := w.sendOne(ctx, token, payload)
+ mu.Lock()
+ defer mu.Unlock()
+ if err != nil {
+ failed++
+ errs = append(errs, err)
+ log.Error(ctx, "web push failed", "err", err)
+ } else {
+ success++
+ }
+ }(t.Token)
+ }
+ wg.Wait()
+
+ log.Log(ctx, "web push blast complete", "success", success, "failed", failed, "total", len(webTargets))
+ if len(errs) == 0 {
+ return nil
+ }
+ if len(errs) == 1 {
+ return errs[0]
+ }
+ return fmt.Errorf("web push blast: %d of %d failed: %v", len(errs), len(webTargets), errs)
+}
+
+// sendOne decrypts the stored subscription JSON and POSTs the encrypted
+// payload to the browser's push endpoint.
+func (w *WebPushNotifier) sendOne(ctx context.Context, subscriptionJSON string, payload []byte) error {
+ var sub webpush.Subscription
+ if err := json.Unmarshal([]byte(subscriptionJSON), &sub); err != nil {
+ return fmt.Errorf("error parsing web push subscription: %w", err)
+ }
+ resp, err := webpush.SendNotificationWithContext(ctx, payload, &sub, &webpush.Options{
+ VAPIDPublicKey: w.keys.PublicKey,
+ VAPIDPrivateKey: w.keys.PrivateKey,
+ Subscriber: w.subscriber,
+ TTL: 24 * 60 * 60, // 24h
+ })
+ if err != nil {
+ return fmt.Errorf("error sending web push: %w", err)
+ }
+ defer resp.Body.Close()
+ // 410 Gone means the subscription is no longer valid; the caller should
+ // prune it. We surface it distinctly so the queue processor can react.
+ if resp.StatusCode == 410 || resp.StatusCode == 404 {
+ return &ExpiredSubscriptionError{Endpoint: sub.Endpoint, Status: resp.StatusCode}
+ }
+ if resp.StatusCode >= 400 {
+ return fmt.Errorf("web push endpoint returned status %d", resp.StatusCode)
+ }
+ return nil
+}
+
+// ExpiredSubscriptionError indicates a push endpoint returned 410 Gone (or
+// 404), meaning the subscription is dead and should be removed from the DB.
+type ExpiredSubscriptionError struct {
+ Endpoint string
+ Status int
+}
+
+func (e *ExpiredSubscriptionError) Error() string {
+ return fmt.Sprintf("web push subscription expired (status %d): %s", e.Status, e.Endpoint)
+}
diff --git a/pkg/notifications/webpush_test.go b/pkg/notifications/webpush_test.go
new file mode 100644
index 00000000..7a57b311
--- /dev/null
+++ b/pkg/notifications/webpush_test.go
@@ -0,0 +1,127 @@
+package notifications
+
+import (
+ "context"
+ "encoding/json"
+ "io"
+ "net/http"
+ "net/http/httptest"
+ "testing"
+
+ "github.com/stretchr/testify/require"
+ webpush "github.com/SherClockHolmes/webpush-go"
+)
+
+// TestWebPushNotifierBlast verifies that the WebPushNotifier:
+// - ignores firebase targets,
+// - POSTs the encrypted payload to each web subscription's endpoint,
+// - surfaces a 410 Gone as an ExpiredSubscriptionError so the caller can
+// prune the dead subscription.
+func TestWebPushNotifierBlast(t *testing.T) {
+ // Spin up a fake push service that records what it receives.
+ var (
+ gotPaths []string
+ gotBodies [][]byte
+ status int
+ )
+ mux := http.NewServeMux()
+ mux.HandleFunc("/push/", func(w http.ResponseWriter, r *http.Request) {
+ gotPaths = append(gotPaths, r.URL.Path)
+ body, _ := io.ReadAll(r.Body)
+ gotBodies = append(gotBodies, body)
+ if status != 0 {
+ w.WriteHeader(status)
+ return
+ }
+ w.WriteHeader(201)
+ })
+ srv := httptest.NewServer(mux)
+ defer srv.Close()
+
+ // Generate a real VAPID keypair so the notifier can sign.
+ priv, pub, err := webpush.GenerateVAPIDKeys()
+ require.NoError(t, err)
+ notifier := NewWebPushNotifier(VAPIDKeys{PublicKey: pub, PrivateKey: priv}, "mailto:test@example.com")
+
+ // Build a web subscription pointing at our fake push service.
+ sub := webpush.Subscription{
+ Endpoint: srv.URL + "/push/abc",
+ Keys: webpush.Keys{
+ P256dh: "BMd4Zb1d3Z2Z8Z8Z8Z8Z8Z8Z8Z8Z8Z8Z8Z8Z8Z8Z8Z8",
+ Auth: "dGhpcyBpcyBhbiBhdXRoIGtleQ",
+ },
+ }
+ // Use a real, valid-length p256dh so the library doesn't reject it before
+ // hitting the network. Generate a throwaway ECDH keypair for the client
+ // side and use its raw public key.
+ _, clientPub, err := webpush.GenerateVAPIDKeys()
+ require.NoError(t, err)
+ sub.Keys.P256dh = clientPub
+
+ subJSON, err := json.Marshal(sub)
+ require.NoError(t, err)
+
+ targets := []NotificationTarget{
+ {Token: "firebase-token-should-be-ignored", Type: NotificationTypeFirebase},
+ {Token: string(subJSON), Type: NotificationTypeWeb},
+ }
+
+ blast := &NotificationBlast{
+ Title: "🔴 @test is LIVE!",
+ Body: "hello world",
+ Data: map[string]string{"path": "/test"},
+ }
+
+ err = notifier.Blast(context.Background(), targets, blast)
+ require.NoError(t, err)
+ require.Len(t, gotPaths, 1, "only the web target should have been pushed to")
+ require.Equal(t, "/push/abc", gotPaths[0])
+
+ // The body is encrypted (RFC 8291), so it won't be our plaintext JSON —
+ // just confirm something non-empty was sent.
+ require.NotEmpty(t, gotBodies[0])
+
+ // Now make the endpoint return 410 Gone and confirm we get an
+ // ExpiredSubscriptionError.
+ status = 410
+ err = notifier.Blast(context.Background(), targets, blast)
+ require.Error(t, err)
+ var expired *ExpiredSubscriptionError
+ require.ErrorAs(t, err, &expired, "410 should surface as ExpiredSubscriptionError")
+ require.Equal(t, 410, expired.Status)
+}
+
+// TestMultiNotifierFanout confirms the MultiNotifier delegates to each child
+// notifier and that each child only acts on its own target type.
+func TestMultiNotifierFanout(t *testing.T) {
+ fb := &recordingNotifier{typeFilter: NotificationTypeFirebase}
+ web := &recordingNotifier{typeFilter: NotificationTypeWeb}
+ multi := NewMultiNotifier(fb, web)
+
+ targets := []NotificationTarget{
+ {Token: "fb-1", Type: NotificationTypeFirebase},
+ {Token: "web-1", Type: NotificationTypeWeb},
+ {Token: "fb-2", Type: NotificationTypeFirebase},
+ }
+ blast := &NotificationBlast{Title: "t", Body: "b"}
+
+ require.NoError(t, multi.Blast(context.Background(), targets, blast))
+ require.Equal(t, []string{"fb-1", "fb-2"}, fb.seen, "firebase notifier should only see firebase targets")
+ require.Equal(t, []string{"web-1"}, web.seen, "web notifier should only see web targets")
+}
+
+// recordingNotifier is a test double that records the tokens it was asked to
+// blast, filtered to a single type.
+type recordingNotifier struct {
+ typeFilter NotificationType
+ seen []string
+}
+
+func (r *recordingNotifier) Blast(ctx context.Context, targets []NotificationTarget, blast *NotificationBlast) error {
+ for _, t := range targets {
+ if t.Type == r.typeFilter {
+ r.seen = append(r.seen, t.Token)
+ }
+ }
+ return nil
+}
diff --git a/pkg/statedb/notification.go b/pkg/statedb/notification.go
index 3c98f74e..de96e709 100644
--- a/pkg/statedb/notification.go
+++ b/pkg/statedb/notification.go
@@ -5,13 +5,25 @@ import (
"time"
"gorm.io/gorm/clause"
+ notificationpkg "stream.place/streamplace/pkg/notifications"
+)
+
+// NotificationType is re-exported from pkg/notifications so callers of this
+// package don't need a second import to name the transport.
+type NotificationType = notificationpkg.NotificationType
+
+// Re-export the transport constants for the same reason.
+const (
+ NotificationTypeFirebase = notificationpkg.NotificationTypeFirebase
+ NotificationTypeWeb = notificationpkg.NotificationTypeWeb
)
type Notification struct {
- Token string `gorm:"column:token;primarykey"`
- RepoDID string `json:"repoDID,omitempty" gorm:"column:repo_did;index"`
- CreatedAt time.Time `gorm:"column:created_at"`
- UpdatedAt time.Time `gorm:"column:updated_at"`
+ Token string `gorm:"column:token;primarykey"`
+ RepoDID string `json:"repoDID,omitempty" gorm:"column:repo_did;index"`
+ Type NotificationType `json:"type,omitempty" gorm:"column:type;default:firebase"`
+ CreatedAt time.Time `gorm:"column:created_at"`
+ UpdatedAt time.Time `gorm:"column:updated_at"`
}
// CreateNotification registers (or refreshes) a device's push token. When a
@@ -19,19 +31,28 @@ type Notification struct {
// can target the user's followers. When repoDID is empty we make sure the
// token row exists but never clobber an existing repoDID association.
//
+// notifType selects the push transport ("firebase" or "web"). An empty value
+// defaults to "firebase" so existing callers and pre-migration rows keep
+// working. The type is only written when a row is created or a repoDID is
+// being upserted; a DID-less re-registration deliberately leaves the type
+// untouched (mirroring the repoDID-preservation behavior below).
+//
// This deliberately avoids DB.Save(): Save issues a full-row UPDATE including
// zero-value columns, so a re-registration with no repoDID (e.g. the client
// posts before its OAuth session has restored) would blank out repo_did and
// silently drop the user from follower notifications.
-func (state *StatefulDB) CreateNotification(token string, repoDID string) error {
+func (state *StatefulDB) CreateNotification(token string, repoDID string, notifType NotificationType) error {
+ if notifType == "" {
+ notifType = NotificationTypeFirebase
+ }
if repoDID != "" {
- not := Notification{Token: token, RepoDID: repoDID}
+ not := Notification{Token: token, RepoDID: repoDID, Type: notifType}
return state.DB.Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "token"}},
- DoUpdates: clause.AssignmentColumns([]string{"repo_did", "updated_at"}),
+ DoUpdates: clause.AssignmentColumns([]string{"repo_did", "type", "updated_at"}),
}).Create(¬).Error
}
- not := Notification{Token: token}
+ not := Notification{Token: token, Type: notifType}
return state.DB.Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "token"}},
DoNothing: true,
@@ -56,6 +77,9 @@ func (state *StatefulDB) ListUserNotifications(userDID string) ([]Notification,
return nots, nil
}
+// GetManyNotificationTokens returns the raw token strings for the given user
+// DIDs, across all notification types. Kept for backwards compatibility with
+// callers that only need the token list (e.g. the legacy blast endpoint).
func (state *StatefulDB) GetManyNotificationTokens(userDIDs []string) ([]string, error) {
tokens := []string{}
err := state.DB.Model(&Notification{}).
@@ -67,3 +91,21 @@ func (state *StatefulDB) GetManyNotificationTokens(userDIDs []string) ([]string,
}
return tokens, nil
}
+
+// GetManyNotifications returns the full notification rows for the given user
+// DIDs, including each token's Type so the notifier can route it to the
+// correct transport (firebase vs web).
+func (state *StatefulDB) GetManyNotifications(userDIDs []string) ([]Notification, error) {
+ nots := []Notification{}
+ err := state.DB.Where("repo_did IN (?)", userDIDs).Find(¬s).Error
+ if err != nil {
+ return nil, fmt.Errorf("error retrieving notifications: %w", err)
+ }
+ return nots, nil
+}
+
+// DeleteNotification removes a token row, used when a web client unsubscribes
+// (or a mobile token is revoked). Missing rows are not an error.
+func (state *StatefulDB) DeleteNotification(token string) error {
+ return state.DB.Where("token = ?", token).Delete(&Notification{}).Error
+}
diff --git a/pkg/statedb/notification_test.go b/pkg/statedb/notification_test.go
index 1f593df9..7908b1a9 100644
--- a/pkg/statedb/notification_test.go
+++ b/pkg/statedb/notification_test.go
@@ -19,14 +19,14 @@ func TestNotificationRepoDIDPreserved(t *testing.T) {
const didB = "did:plc:bbbb"
// Initial registration while logged in associates the DID.
- require.NoError(t, state.CreateNotification(token, didA))
+ require.NoError(t, state.CreateNotification(token, didA, NotificationTypeFirebase))
tokens, err := state.GetManyNotificationTokens([]string{didA})
require.NoError(t, err)
require.Equal(t, []string{token}, tokens)
// Re-registration without a DID (e.g. before the OAuth session has
// restored) must NOT clobber the existing association.
- require.NoError(t, state.CreateNotification(token, ""))
+ require.NoError(t, state.CreateNotification(token, "", NotificationTypeFirebase))
tokens, err = state.GetManyNotificationTokens([]string{didA})
require.NoError(t, err)
require.Equal(t, []string{token}, tokens, "repo_did was wiped by a DID-less re-registration")
@@ -37,7 +37,7 @@ func TestNotificationRepoDIDPreserved(t *testing.T) {
require.Len(t, nots, 1)
// Re-registering with a different DID replaces the association.
- require.NoError(t, state.CreateNotification(token, didB))
+ require.NoError(t, state.CreateNotification(token, didB, NotificationTypeFirebase))
tokens, err = state.GetManyNotificationTokens([]string{didB})
require.NoError(t, err)
require.Equal(t, []string{token}, tokens)
@@ -56,7 +56,7 @@ func TestNotificationAnonymousThenAssociated(t *testing.T) {
const did = "did:plc:cccc"
// Anonymous registration: the row exists but has no association yet.
- require.NoError(t, state.CreateNotification(token, ""))
+ require.NoError(t, state.CreateNotification(token, "", NotificationTypeFirebase))
tokens, err := state.GetManyNotificationTokens([]string{did})
require.NoError(t, err)
require.Empty(t, tokens)
@@ -65,7 +65,7 @@ func TestNotificationAnonymousThenAssociated(t *testing.T) {
require.Len(t, nots, 1)
// Once logged in, the association is set without adding a new row.
- require.NoError(t, state.CreateNotification(token, did))
+ require.NoError(t, state.CreateNotification(token, did, NotificationTypeFirebase))
tokens, err = state.GetManyNotificationTokens([]string{did})
require.NoError(t, err)
require.Equal(t, []string{token}, tokens)
@@ -74,3 +74,53 @@ func TestNotificationAnonymousThenAssociated(t *testing.T) {
require.Len(t, nots, 1)
})
}
+
+// TestNotificationTypeAndDelete covers the Type column (firebase vs web) and
+// the DeleteNotification path used when a web client unsubscribes.
+func TestNotificationTypeAndDelete(t *testing.T) {
+ WithAllDatabases(t, func(state *StatefulDB) {
+ const did = "did:plc:dddd"
+ const fbToken = "firebase-token-1"
+ const webToken = `{"endpoint":"https://push.example/abc","keys":{"p256dh":"x","auth":"y"}}`
+
+ // Register one firebase and one web subscription for the same user.
+ require.NoError(t, state.CreateNotification(fbToken, did, NotificationTypeFirebase))
+ require.NoError(t, state.CreateNotification(webToken, did, NotificationTypeWeb))
+
+ // GetManyNotifications returns both rows with their types intact.
+ nots, err := state.GetManyNotifications([]string{did})
+ require.NoError(t, err)
+ require.Len(t, nots, 2)
+
+ byType := map[NotificationType]Notification{}
+ for _, n := range nots {
+ byType[n.Type] = n
+ }
+ require.Contains(t, byType, NotificationTypeFirebase)
+ require.Contains(t, byType, NotificationTypeWeb)
+ require.Equal(t, fbToken, byType[NotificationTypeFirebase].Token)
+ require.Equal(t, webToken, byType[NotificationTypeWeb].Token)
+
+ // An empty type defaults to firebase.
+ require.NoError(t, state.CreateNotification("defaulted-token", did, ""))
+ nots, err = state.GetManyNotifications([]string{did})
+ require.NoError(t, err)
+ var defaulted Notification
+ for _, n := range nots {
+ if n.Token == "defaulted-token" {
+ defaulted = n
+ }
+ }
+ require.Equal(t, NotificationTypeFirebase, defaulted.Type, "empty type should default to firebase")
+
+ // DeleteNotification removes the row; deleting a missing row is not an error.
+ require.NoError(t, state.DeleteNotification(webToken))
+ nots, err = state.GetManyNotifications([]string{did})
+ require.NoError(t, err)
+ require.Len(t, nots, 2, "web token should be gone, leaving firebase + defaulted")
+ for _, n := range nots {
+ require.NotEqual(t, webToken, n.Token)
+ }
+ require.NoError(t, state.DeleteNotification("never-existed"))
+ })
+}
diff --git a/pkg/statedb/queue_processor.go b/pkg/statedb/queue_processor.go
index ae5b2489..cc765433 100644
--- a/pkg/statedb/queue_processor.go
+++ b/pkg/statedb/queue_processor.go
@@ -185,6 +185,13 @@ type VODProcessor func(ctx context.Context, t VODProcessTask) (cid string, err e
func (state *StatefulDB) SetVODProcessor(f VODProcessor) { state.vodProcessor = f }
+// SetNotifier installs the notification notifier after construction. This is
+// needed because building the Web Push notifier requires VAPID keys, which
+// are stored in the DB — so the DB must exist before the notifier can be
+// fully assembled. The queue processor checks for nil, so a brief window
+// with no notifier is safe.
+func (state *StatefulDB) SetNotifier(n notificationpkg.Notifier) { state.noter = n }
+
func (state *StatefulDB) processVODProcessTask(ctx context.Context, task *AppTask) error {
ctx = log.WithLogValues(ctx, "func", "processVODProcessTask")
var t VODProcessTask
@@ -467,7 +474,7 @@ func (state *StatefulDB) processNotificationTask(ctx context.Context, task *AppT
log.Log(ctx, "found followers", "count", len(followersDIDs))
- notifications, err := state.GetManyNotificationTokens(followersDIDs)
+ notifications, err := state.GetManyNotifications(followersDIDs)
if err != nil {
return err
}
@@ -480,7 +487,11 @@ func (state *StatefulDB) processNotificationTask(ctx context.Context, task *AppT
"path": fmt.Sprintf("/%s", lsv.Author.Handle),
},
}
- err = state.noter.Blast(ctx, notifications, nb)
+ targets := make([]notificationpkg.NotificationTarget, len(notifications))
+ for i, n := range notifications {
+ targets[i] = notificationpkg.NotificationTarget{Token: n.Token, Type: n.Type}
+ }
+ err = state.noter.Blast(ctx, targets, nb)
if err != nil {
log.Error(ctx, "failed to blast notifications", "err", err)
} else {
diff --git a/pkg/statedb/statedb.go b/pkg/statedb/statedb.go
index 311fb3ac..ea1d8f94 100644
--- a/pkg/statedb/statedb.go
+++ b/pkg/statedb/statedb.go
@@ -30,7 +30,7 @@ type StatefulDB struct {
CLI *config.CLI
Type DBType
locks *NamedLocks
- noter notificationpkg.FirebaseNotifier
+ noter notificationpkg.Notifier
model model.Model
// pokeQueue is used to wake up the queue processor when a new task is enqueued
pokeQueue chan struct{}
@@ -76,7 +76,7 @@ var StatefulDBModels = []any{
var NoPostgresDatabaseCode = "3D000"
// Stateful database for storing private streamplace state
-func MakeDB(ctx context.Context, cli *config.CLI, noter notificationpkg.FirebaseNotifier, model model.Model) (*StatefulDB, error) {
+func MakeDB(ctx context.Context, cli *config.CLI, noter notificationpkg.Notifier, model model.Model) (*StatefulDB, error) {
dbURL := cli.DBURL
log.Log(ctx, "starting stateful database", "dbURL", redactDBURL(dbURL))
var dial gorm.Dialector
diff --git a/pkg/statedb/vapid.go b/pkg/statedb/vapid.go
new file mode 100644
index 00000000..54a6bb4c
--- /dev/null
+++ b/pkg/statedb/vapid.go
@@ -0,0 +1,59 @@
+package statedb
+
+import (
+ "context"
+ "encoding/json"
+ "fmt"
+
+ webpush "github.com/SherClockHolmes/webpush-go"
+ notificationpkg "stream.place/streamplace/pkg/notifications"
+ "stream.place/streamplace/pkg/log"
+)
+
+// vapidConfigKey is the Config-table key under which the VAPID keypair is
+// stored, mirroring how EnsureJWK persists JWKs by name.
+const vapidConfigKey = "vapid-keys"
+
+// EnsureVAPIDKeys returns the Web Push VAPID keypair, generating and
+// persisting it on first use. It follows the same pattern as EnsureJWK:
+// look the key up in the Config table; if present, use it; otherwise
+// generate a fresh P-256 keypair and store it so it survives restarts.
+//
+// VAPID keys must stay stable — rotating them invalidates every existing
+// browser subscription, so we never regenerate once a key exists.
+func (state *StatefulDB) EnsureVAPIDKeys(ctx context.Context) (notificationpkg.VAPIDKeys, error) {
+ conf, err := state.GetConfig(vapidConfigKey)
+ if err != nil {
+ return notificationpkg.VAPIDKeys{}, fmt.Errorf("error loading vapid keys: %w", err)
+ }
+
+ // happy path: we found the keys in the database, use that
+ if conf != nil {
+ var keys notificationpkg.VAPIDKeys
+ if err := json.Unmarshal(conf.Value, &keys); err != nil {
+ return notificationpkg.VAPIDKeys{}, fmt.Errorf("error parsing stored vapid keys: %w", err)
+ }
+ return keys, nil
+ }
+
+ // new path: no keys yet, generate a fresh pair
+ log.Warn(ctx, "no VAPID keys found, generating new ones")
+ privateKey, publicKey, err := webpush.GenerateVAPIDKeys()
+ if err != nil {
+ return notificationpkg.VAPIDKeys{}, fmt.Errorf("failed to generate vapid keys: %w", err)
+ }
+ keys := notificationpkg.VAPIDKeys{
+ PublicKey: publicKey,
+ PrivateKey: privateKey,
+ }
+
+ b, err := json.Marshal(keys)
+ if err != nil {
+ return notificationpkg.VAPIDKeys{}, fmt.Errorf("failed to marshal vapid keys: %w", err)
+ }
+ if err := state.PutConfig(vapidConfigKey, b); err != nil {
+ return notificationpkg.VAPIDKeys{}, fmt.Errorf("failed to save vapid keys: %w", err)
+ }
+
+ return keys, nil
+}