import webPush from "web-push"; import { and, eq, inArray } from "drizzle-orm"; import { db } from "@/lib/db"; import logger from "@/lib/logger"; import { pushSubscriptions } from "./schema"; function ensureVapidConfigured() { const subject = process.env["VAPID_SUBJECT"]; const publicKey = process.env["VAPID_PUBLIC_KEY"]; const privateKey = process.env["VAPID_PRIVATE_KEY"]; if (!subject || !publicKey || !privateKey) { throw new Error("VAPID_SUBJECT, VAPID_PUBLIC_KEY, and VAPID_PRIVATE_KEY must be set"); } webPush.setVapidDetails(subject, publicKey, privateKey); } export async function sendPushToEndpoint( endpoint: string, payload: { title: string; body: string; url?: string }, ) { ensureVapidConfigured(); const [sub] = await db .select() .from(pushSubscriptions) .where(eq(pushSubscriptions.endpoint, endpoint)) .limit(1); if (!sub) return; try { await webPush.sendNotification( { endpoint: sub.endpoint, keys: { p256dh: sub.p256dh, auth: sub.auth } }, JSON.stringify({ title: payload.title, body: payload.body, url: payload.url ?? "/" }), ); } catch (err) { const status = (err as { statusCode?: number }).statusCode; if (status === 404 || status === 410) { await db.delete(pushSubscriptions).where(eq(pushSubscriptions.endpoint, endpoint)); } else { logger.error({ err }, "push delivery failed"); } } } export async function sendPush( userId: string, payload: { title: string; body: string; url?: string }, ) { ensureVapidConfigured(); const subs = await db .select() .from(pushSubscriptions) .where(eq(pushSubscriptions.userId, userId)); if (subs.length === 0) return; const staleIds: string[] = []; await Promise.allSettled( subs.map(async (sub) => { try { await webPush.sendNotification( { endpoint: sub.endpoint, keys: { p256dh: sub.p256dh, auth: sub.auth } }, JSON.stringify({ title: payload.title, body: payload.body, url: payload.url ?? "/" }), ); } catch (err) { const status = (err as { statusCode?: number }).statusCode; if (status === 404 || status === 410) { staleIds.push(sub.id); } else { logger.error({ err }, "push delivery failed"); } } }), ); if (staleIds.length > 0) { await db .delete(pushSubscriptions) .where(and(eq(pushSubscriptions.userId, userId), inArray(pushSubscriptions.id, staleIds))); } }