42 lines
3.3 KiB
TypeScript
42 lines
3.3 KiB
TypeScript
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<string>`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<number>`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);
|
|
}
|