feat : ask for noti
This commit is contained in:
+42
-16
@@ -1,33 +1,59 @@
|
||||
import "server-only";
|
||||
|
||||
import { and, eq } from "drizzle-orm";
|
||||
import { and, eq, inArray, ne, sql } from "drizzle-orm";
|
||||
import webPush from "web-push";
|
||||
import { getDb } from "@/db";
|
||||
import { comments, commentPushSubscriptions, commentRevisions, users } from "@/db/schema";
|
||||
import { pushPublicKey, validPushEndpoint } from "@/lib/commission/push";
|
||||
import { createNotifications } from "@/lib/notifications/repository";
|
||||
import { authorizeComment } from "./repository";
|
||||
|
||||
export async function sendCommentReplyPush(id: string) {
|
||||
if (!pushPublicKey()) return;
|
||||
export async function sendCommentPush(id: string) {
|
||||
const [reply] = await getDb().select().from(comments).where(eq(comments.id, id)).limit(1);
|
||||
if (!reply?.replyToId || reply.hidden || reply.deletedAt) return;
|
||||
const [parent] = await getDb().select().from(comments).where(eq(comments.id, reply.replyToId)).limit(1);
|
||||
if (!parent || parent.authorId === reply.authorId || parent.hidden || parent.deletedAt) return;
|
||||
const [recipient] = await getDb().select().from(users).where(eq(users.id, parent.authorId)).limit(1);
|
||||
if (!recipient || recipient.banned) return;
|
||||
let destination;
|
||||
try {
|
||||
({ destination } = await authorizeComment(id, { id: recipient.id, admin: recipient.role === "admin" && recipient.emailVerified }));
|
||||
} catch { return; }
|
||||
if (!reply || reply.hidden || reply.deletedAt) return;
|
||||
let recipients: { id: string; admin: boolean }[];
|
||||
let url: string;
|
||||
let title: string;
|
||||
if (reply.replyToId) {
|
||||
const [parent] = await getDb().select().from(comments).where(eq(comments.id, reply.replyToId)).limit(1);
|
||||
if (!parent || parent.authorId === reply.authorId || parent.hidden || parent.deletedAt) return;
|
||||
const [recipient] = await getDb().select().from(users).where(eq(users.id, parent.authorId)).limit(1);
|
||||
if (!recipient || recipient.banned) return;
|
||||
let destination;
|
||||
try {
|
||||
({ destination } = await authorizeComment(id, { id: recipient.id, admin: recipient.role === "admin" && recipient.emailVerified }));
|
||||
} catch { return; }
|
||||
recipients = [{ id: recipient.id, admin: false }];
|
||||
url = `${destination.href}${destination.href.includes("?") ? "&" : "?"}reply=${id}#comment-${reply.rootId}`;
|
||||
title = "ตอบกลับความคิดเห็นของคุณ";
|
||||
} else {
|
||||
recipients = (await getDb().select({ id: users.id }).from(users)
|
||||
.where(and(eq(users.role, "admin"), eq(users.emailVerified, true), ne(users.id, reply.authorId), sql`${users.banned} is not true`)))
|
||||
.map((user) => ({ ...user, admin: true }));
|
||||
if (!recipients.length) return;
|
||||
let destination;
|
||||
try {
|
||||
({ destination } = await authorizeComment(id, { id: recipients[0].id, admin: true }));
|
||||
} catch { return; }
|
||||
if (!destination.writable) return;
|
||||
url = `/admin/comments?target=${encodeURIComponent(destination.target)}&comment=${id}`;
|
||||
title = `แสดงความคิดเห็นใหม่ใน ${destination.name}`;
|
||||
}
|
||||
const [author] = await getDb().select({ name: users.name }).from(users).where(eq(users.id, reply.authorId)).limit(1);
|
||||
const [revision] = await getDb().select({ text: commentRevisions.text }).from(commentRevisions)
|
||||
.where(and(eq(commentRevisions.commentId, id), eq(commentRevisions.version, reply.version))).limit(1);
|
||||
const subscriptions = await getDb().select().from(commentPushSubscriptions).where(eq(commentPushSubscriptions.userId, recipient.id));
|
||||
const notificationTitle = `${author?.name ?? "ผู้ใช้"} ${title}`;
|
||||
const body = revision?.text.slice(0, 300) || "ส่งรูปภาพ";
|
||||
await createNotifications(recipients.map(recipient => ({ userId: recipient.id, adminOnly: recipient.admin, url })), {
|
||||
eventKey: `comment:${id}`, kind: reply.replyToId ? "comment_reply" : "comment_new", title: notificationTitle, body, commentId: id,
|
||||
});
|
||||
if (!pushPublicKey()) return;
|
||||
const subscriptions = await getDb().select().from(commentPushSubscriptions).where(inArray(commentPushSubscriptions.userId, recipients.map(recipient => recipient.id)));
|
||||
if (!subscriptions.length) return;
|
||||
webPush.setVapidDetails(process.env.WEB_PUSH_SUBJECT!, process.env.WEB_PUSH_PUBLIC_KEY!, process.env.WEB_PUSH_PRIVATE_KEY!);
|
||||
const payload = JSON.stringify({ id: `comment-${id}`, title: `${author?.name ?? "ผู้ใช้"} ตอบกลับความคิดเห็นของคุณ`,
|
||||
body: revision?.text.slice(0, 300) || "ส่งรูปภาพ", icon: "/icon/nav/Comment.webp",
|
||||
url: `${destination.href}${destination.href.includes("?") ? "&" : "?"}reply=${id}#comment-${reply.rootId}` });
|
||||
const payload = JSON.stringify({ id: `comment-${id}`, title: notificationTitle,
|
||||
body, icon: "/icon/nav/Comment.webp",
|
||||
url });
|
||||
for (let offset = 0; offset < subscriptions.length; offset += 10) {
|
||||
await Promise.all(subscriptions.slice(offset, offset + 10).map(async (subscription) => {
|
||||
if (!validPushEndpoint(subscription.endpoint)) return;
|
||||
|
||||
@@ -5,7 +5,8 @@ import { notifyCommentChange } from "./events";
|
||||
import { and, asc, desc, eq, inArray, isNull, lt, or, sql, type SQL } from "drizzle-orm";
|
||||
import { alias } from "drizzle-orm/pg-core";
|
||||
import { getDb, type Database } from "@/db";
|
||||
import { commentThreads, comments, commentRevisions, commentAttachments, commentRevisionAttachments, commentReactions, catalogCharacters, guides, stygianSchedules, users } from "@/db/schema";
|
||||
import { commentThreads, comments, commentRevisions, commentAttachments, commentRevisionAttachments, commentReactions, catalogCharacters, guides, stygianSchedules, users, notifications } from "@/db/schema";
|
||||
import { notifyNotificationChange } from "@/lib/notifications/events";
|
||||
import { getMediaStorage, publicMediaUrl } from "@/lib/media/storage";
|
||||
import { getCustomerSession } from "@/lib/auth/server";
|
||||
import { isAuthorizedAdmin } from "@/lib/auth/authorization";
|
||||
@@ -14,7 +15,7 @@ import { HttpError } from "@/lib/security/http";
|
||||
import { decodeCursor, encodeCursor, parseTarget, uuidSchema } from "./validation";
|
||||
import type { CommentHistoryItem, CommentImage, CommentItem, CommentPage, CommentViewer } from "./types";
|
||||
|
||||
type Reader = Pick<Database, "select" | "insert" | "update" | "delete" | "execute">;
|
||||
type Reader = Pick<Database, "select" | "selectDistinct" | "insert" | "update" | "delete" | "execute">;
|
||||
export async function getCommentViewer(): Promise<CommentViewer | null> {
|
||||
const session = await getCustomerSession();
|
||||
if (!session) return null;
|
||||
@@ -199,11 +200,14 @@ export async function withCommentLock<T>(id: string, viewer: CommentViewer, task
|
||||
export async function mutateComment(id: string, viewer: CommentViewer, action: "delete" | "reaction" | "moderation" | "heart", value?: number | boolean) {
|
||||
let target = "";
|
||||
let deletedImages: { objectKey: string }[] = [];
|
||||
let notificationUsers: { userId: string }[] = [];
|
||||
const result = await withCommentLock(id, viewer, async (tx, context) => {
|
||||
target = context.destination.target;
|
||||
const c = context.comment;
|
||||
if (action === "moderation") {
|
||||
requireCommentAdmin(viewer);
|
||||
notificationUsers = await tx.selectDistinct({ userId: notifications.userId }).from(notifications)
|
||||
.innerJoin(comments, eq(comments.id, notifications.commentId)).where(or(eq(comments.id, id), eq(comments.rootId, id)));
|
||||
await tx.update(comments).set({ hidden: value as boolean }).where(eq(comments.id, id));
|
||||
const [actor] = await tx.select({ id: users.id, name: users.name }).from(users).where(eq(users.id, viewer.id));
|
||||
await writeAuditLog(tx, auditActor(actor), { action: value ? "comment.hidden" : "comment.restored", targetType: "comment", targetId: id,
|
||||
@@ -221,6 +225,8 @@ export async function mutateComment(id: string, viewer: CommentViewer, action: "
|
||||
) select id from descendants`;
|
||||
deletedImages = await tx.select({ objectKey: commentAttachments.objectKey }).from(commentAttachments)
|
||||
.where(sql`${commentAttachments.commentId} in (${descendants})`);
|
||||
notificationUsers = await tx.selectDistinct({ userId: notifications.userId }).from(notifications)
|
||||
.where(sql`${notifications.commentId} in (${descendants})`);
|
||||
// Delete the entire subtree in one statement so self-referencing foreign keys remain valid.
|
||||
// Revisions, attachment metadata, revision links, and reactions cascade automatically.
|
||||
await tx.delete(comments).where(sql`${comments.id} in (${descendants})`);
|
||||
@@ -234,6 +240,8 @@ export async function mutateComment(id: string, viewer: CommentViewer, action: "
|
||||
.onConflictDoUpdate({ target: [commentReactions.commentId, commentReactions.userId], set: { value: value as number } });
|
||||
}
|
||||
});
|
||||
for (let offset = 0; offset < notificationUsers.length; offset += 10)
|
||||
await Promise.all(notificationUsers.slice(offset, offset + 10).map(user => notifyNotificationChange(user.userId)));
|
||||
await notifyCommentChange(target);
|
||||
if (deletedImages.length) {
|
||||
try {
|
||||
|
||||
@@ -1,7 +1,8 @@
|
||||
import "server-only";
|
||||
|
||||
import { and, eq, ne, or } from "drizzle-orm";
|
||||
import { and, eq, ne, or, sql } from "drizzle-orm";
|
||||
import webPush from "web-push";
|
||||
import { createNotifications } from "@/lib/notifications/repository";
|
||||
import { getDb } from "@/db";
|
||||
import { commissionPushSubscriptions, users } from "@/db/schema";
|
||||
|
||||
@@ -25,6 +26,12 @@ export async function sendCommissionMessagePush(message: {
|
||||
id: string; ticketId: string; ticketTitle: string; customerId: string; authorId: string;
|
||||
authorName: string; text: string | null;
|
||||
}) {
|
||||
const recipients = await getDb().select({ id: users.id, role: users.role, verified: users.emailVerified }).from(users)
|
||||
.where(and(ne(users.id, message.authorId), sql`${users.banned} is not true`, or(eq(users.id, message.customerId), and(eq(users.role, "admin"), eq(users.emailVerified, true)))));
|
||||
await createNotifications(recipients.map(user => {
|
||||
const admin = user.role === "admin" && user.verified && user.id !== message.customerId;
|
||||
return { userId: user.id, adminOnly: admin, url: `/${admin ? "admin/commission" : "commission/tickets"}/${message.ticketId}` };
|
||||
}), { eventKey: `commission:${message.id}`, kind: "commission_message", title: `${message.authorName} - ${message.ticketTitle}`, body: message.text?.slice(0, 300) || "ส่งรูปภาพ", messageId: message.id });
|
||||
if (!pushPublicKey()) return;
|
||||
const subscriptions = await getDb().select({ subscription: commissionPushSubscriptions })
|
||||
.from(commissionPushSubscriptions)
|
||||
|
||||
@@ -0,0 +1,6 @@
|
||||
import 'server-only';
|
||||
import { getRedisClient, redisEventChannel } from '@/lib/redis/client';
|
||||
export async function notifyNotificationChange(userId: string) {
|
||||
try { await (await getRedisClient()).publish(redisEventChannel(`notifications:user:${userId}`), 'changed'); }
|
||||
catch { /* Saved notifications remain available after reconnect. */ }
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
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);
|
||||
}
|
||||
Reference in New Issue
Block a user