import 'server-only'; import { and, desc, eq, isNull, or, sql } from 'drizzle-orm'; import { alias } from 'drizzle-orm/pg-core'; import { getDb } from '@/db'; import { notifications, comments, commentThreads, guides } from '@/db/schema'; import { notifyNotificationChange } from './events'; type Recipient = { userId: string; adminOnly?: boolean; url: string }; type Message = { eventKey: string; kind: string; title: string; body: string; commentId?: string; messageId?: string }; export async function createNotifications(recipients: Recipient[], message: Message) { for (let offset = 0; offset < recipients.length; offset += 100) { const inserted = await getDb().insert(notifications).values(recipients.slice(offset, offset + 100).map(recipient => ({ ...recipient, ...message }))) .onConflictDoNothing({ target: [notifications.userId, notifications.eventKey] }).returning({ userId: notifications.userId }); for (let start = 0; start < inserted.length; start += 10) await Promise.all(inserted.slice(start, start + 10).map(row => notifyNotificationChange(row.userId))); } } const root = alias(comments, 'notification_root'); export async function listNotifications(user: { id: string; role: string | null; emailVerified: boolean }, before?: { time: string; id: string }) { const admin = user.role === 'admin' && user.emailVerified; const visible = and(eq(notifications.userId, user.id), admin ? undefined : and(eq(notifications.adminOnly, false), or(isNull(notifications.commentId), and( eq(comments.hidden, false), sql`coalesce(${root.hidden}, false) = false`, or(isNull(commentThreads.guideId), and(eq(guides.isPublic, true), isNull(guides.trashedAt))), )))); const query = () => getDb().select({ notification: notifications, cursorTime: sql`to_char(${notifications.createdAt} at time zone 'UTC', 'YYYY-MM-DD"T"HH24:MI:SS.US"Z"')` }).from(notifications) .leftJoin(comments, eq(comments.id, notifications.commentId)).leftJoin(root, eq(root.id, comments.rootId)) .leftJoin(commentThreads, eq(commentThreads.id, comments.threadId)).leftJoin(guides, eq(guides.id, commentThreads.guideId)); const [rows, count] = await Promise.all([ query().where(and(visible, before ? sql`(${notifications.createdAt}, ${notifications.id}) < (${before.time}::timestamptz, ${before.id}::uuid)` : undefined)).orderBy(desc(notifications.createdAt), desc(notifications.id)).limit(21), getDb().select({ count: sql`count(*)::int` }).from(notifications) .leftJoin(comments, eq(comments.id, notifications.commentId)).leftJoin(root, eq(root.id, comments.rootId)) .leftJoin(commentThreads, eq(commentThreads.id, comments.threadId)).leftJoin(guides, eq(guides.id, commentThreads.guideId)) .where(and(visible, isNull(notifications.readAt))), ]); const items = rows.slice(0, 20).map(row => row.notification); const last = rows[19]; return { items, unread: count[0]?.count ?? 0, nextCursor: rows.length > 20 ? `${last.cursorTime}|${last.notification.id}` : null }; } export async function readNotifications(userId: string, id?: string) { await getDb().update(notifications).set({ readAt: sql`now()` }).where(and(eq(notifications.userId, userId), isNull(notifications.readAt), id ? eq(notifications.id, id) : undefined)); await notifyNotificationChange(userId); }