diff --git a/app/stygian/page.tsx b/app/stygian/page.tsx index fb47d5a..6751e58 100644 --- a/app/stygian/page.tsx +++ b/app/stygian/page.tsx @@ -1,6 +1,5 @@ import Link from "next/link"; import { Suspense } from "react"; -import Image from "@/components/ui/resilient-image"; import { connection } from "next/server"; import { SiteFooter } from "@/components/public/site-footer"; @@ -37,12 +36,6 @@ export default async function StygianPage({ searchParams }: { searchParams: Prom {schedules.length ?

เลือกเวอร์ชัน

: null} - {schedule && } {schedule && catalog ? : ( ยังไม่มีไกด์ Stygianกำลังเตรียมข้อมูล กลับมาอ่านได้ภายหลัง หรือเลือกไกด์ตัวละครระหว่างรอเลือกไกด์ตัวละคร )} diff --git a/lib/comments/heart-notification.test.ts b/lib/comments/heart-notification.test.ts new file mode 100644 index 0000000..bff5bf1 --- /dev/null +++ b/lib/comments/heart-notification.test.ts @@ -0,0 +1,95 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const state = vi.hoisted(() => ({ + comment: {} as Record, author: {} as Record | null, + publicGuide: true, notifications: [] as Record[], + notify: vi.fn(), changed: vi.fn(), +})); +vi.mock("server-only", () => ({})); +vi.mock("@/lib/notifications/events", () => ({ notifyNotificationChange: state.notify })); +vi.mock("@/lib/comments/events", () => ({ notifyCommentChange: state.changed })); +vi.mock("@/lib/auth/server", () => ({ getCustomerSession: vi.fn() })); +vi.mock("@/lib/audit-log", () => ({ auditActor: vi.fn(), writeAuditLog: vi.fn() })); +vi.mock("@/lib/media/storage", () => ({ getMediaStorage: vi.fn(), publicMediaUrl: vi.fn() })); +vi.mock("@/db", () => ({ getDb: () => database })); + +import { comments, guides, users } from "@/db/schema"; +import { mutateComment } from "./repository"; + +const id = "72df08ab-50dd-4cbd-9a69-70949d34cf9f"; +const guideId = "72df08ab-50dd-4cbd-9a69-70949d34cf9e"; +const admin = { id: "admin", admin: true }; +const database = { + transaction: async (task: (tx: unknown) => unknown) => task(database), + select: (fields: Record) => { + let table: unknown; + const rows = () => { + if (table === comments) return fields.comment + ? [{ comment: state.comment, thread: { id: "thread", guideId } }] + : [{ hidden: false }]; + if (table === guides) return [{ id: guideId, name: "Amber", slug: "amber", public: state.publicGuide, trashedAt: null }]; + if (table === users) return state.author ? [state.author] : []; + return []; + }; + const query = { + from(value: unknown) { table = value; return query; }, + innerJoin() { return query; }, where() { return query; }, + limit: async () => rows(), for: async () => rows(), + }; + return query; + }, + update: () => { + const query = { set: () => query, where: async () => undefined }; + return query; + }, + insert: () => ({ values: (value: Record) => ({ onConflictDoNothing: () => ({ returning: async () => { + if (state.notifications.some(old => old.eventKey === value.eventKey)) return []; + state.notifications.push(value); + return [{ userId: value.userId }]; + } }) }) }), +}; + +beforeEach(() => { + state.comment = { id, authorId: "reader", rootId: null, hidden: false, deletedAt: null, heartedById: null }; + state.author = { id: "reader", email: "reader@example.test", role: "user", emailVerified: true, banned: false }; + state.publicGuide = true; state.notifications.length = 0; state.notify.mockReset(); state.changed.mockReset(); +}); + +describe("admin heart notifications", () => { + it("notifies the comment author and refreshes their bell after saving", async () => { + await mutateComment(id, admin, "heart", true); + expect(state.notifications).toEqual([expect.objectContaining({ userId: "reader", kind: "comment_heart", + title: "คุณได้รับหัวใจจากแอดมิน", url: `/amber/comment#comment-${id}`, commentId: id })]); + expect(state.notify).toHaveBeenCalledWith("reader"); + }); + + it("links a reply to its exact thread", async () => { + state.comment.rootId = "root-id"; + await mutateComment(id, admin, "heart", true); + expect(state.notifications[0].url).toBe(`/amber/comment?reply=${id}#comment-root-id`); + }); + + it.each(["removed", "already hearted", "self", "banned", "missing", "private guide"])("skips notifications for %s", async (scenario) => { + if (scenario === "already hearted") state.comment.heartedById = "other-admin"; + if (scenario === "self") state.comment.authorId = admin.id; + if (scenario === "banned") state.author!.banned = true; + if (scenario === "missing") state.author = null; + if (scenario === "private guide") state.publicGuide = false; + await mutateComment(id, admin, "heart", scenario !== "removed"); + expect(state.notifications).toHaveLength(0); + expect(state.notify).not.toHaveBeenCalled(); + }); + + it("does not duplicate a notification when the heart is added again", async () => { + await mutateComment(id, admin, "heart", true); + await mutateComment(id, admin, "heart", false); + await mutateComment(id, admin, "heart", true); + expect(state.notifications).toHaveLength(1); + expect(state.notify).toHaveBeenCalledTimes(1); + }); + + it("rejects hearts from regular users", async () => { + await expect(mutateComment(id, { id: "reader", admin: false }, "heart", true)).rejects.toMatchObject({ status: 403 }); + expect(state.notifications).toHaveLength(0); + }); +}); diff --git a/lib/comments/repository.ts b/lib/comments/repository.ts index 3cf7cb7..ddaf8c8 100644 --- a/lib/comments/repository.ts +++ b/lib/comments/repository.ts @@ -241,6 +241,24 @@ export async function mutateComment(id: string, viewer: CommentViewer, action: " } else if (action === "heart") { requireCommentAdmin(viewer); await tx.update(comments).set({ heartedById: value ? viewer.id : null }).where(eq(comments.id, id)); + if (value && !c.heartedById && c.authorId !== viewer.id) { + const [author] = await tx.select({ id: users.id, email: users.email, role: users.role, emailVerified: users.emailVerified, banned: users.banned }) + .from(users).where(eq(users.id, c.authorId)).limit(1); + if (!author || author.banned) return; + let destination; + try { + ({ destination } = await authorizeComment(id, { id: author.id, admin: isAuthorizedAdmin(author) }, tx)); + } catch (cause) { + if (cause instanceof HttpError && cause.status === 404) return; + throw cause; + } + const url = `${destination.href}${c.rootId ? `${destination.href.includes("?") ? "&" : "?"}reply=${id}` : ""}#comment-${c.rootId ?? id}`; + notificationUsers = await tx.insert(notifications).values({ + userId: author.id, eventKey: `comment:${id}:heart`, kind: "comment_heart", + title: "คุณได้รับหัวใจจากแอดมิน", body: `แอดมินให้หัวใจความคิดเห็นของคุณใน ${destination.name}`, + url, commentId: id, + }).onConflictDoNothing({ target: [notifications.userId, notifications.eventKey] }).returning({ userId: notifications.userId }); + } } else if (value === 0) { await tx.delete(commentReactions).where(and(eq(commentReactions.commentId, id), eq(commentReactions.userId, viewer.id))); } else { diff --git a/lib/guides/mutations.ts b/lib/guides/mutations.ts index e18fc12..d484ca8 100644 --- a/lib/guides/mutations.ts +++ b/lib/guides/mutations.ts @@ -599,6 +599,11 @@ export async function setGuideVisibility( context: GuideMutationContext, ) { return getDb().transaction(async (tx) => { + const [previous] = await tx + .select({ isPublic: guides.isPublic, publishedAt: guides.publishedAt }) + .from(guides) + .where(eq(guides.id, guideId)) + .for("update"); const [guide] = await tx .update(guides) .set({ @@ -624,6 +629,14 @@ export async function setGuideVisibility( throw new GuideVersionConflictError(current.version); } await emitGuideEvents(tx, guide.id, guide.version, true); + if (input.isPublic && previous && !previous.isPublic && !previous.publishedAt) { + await tx.insert(outboxEvents).values({ + topic: "notifications:guide", + aggregateId: guide.id, + eventType: "guide.published", + payload: { publishedAt: guide.publishedAt!.toISOString() }, + }); + } await auditGuide( tx, context, diff --git a/lib/guides/notifications.test.ts b/lib/guides/notifications.test.ts new file mode 100644 index 0000000..564045a --- /dev/null +++ b/lib/guides/notifications.test.ts @@ -0,0 +1,64 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { PgDialect } from "drizzle-orm/pg-core"; + +const mocks = vi.hoisted(() => ({ rows: [] as unknown[][], filters: [] as unknown[], create: vi.fn() })); +vi.mock("server-only", () => ({})); +vi.mock("@/lib/notifications/repository", () => ({ createNotifications: mocks.create })); +vi.mock("@/db", () => ({ getDb: () => ({ select: () => { + const query = { + from() { return query; }, + where(filter: unknown) { mocks.filters.push(filter); return query; }, + orderBy() { return query; }, + limit() { return Promise.resolve(mocks.rows.shift() ?? []); }, + }; + return query; +} }) })); + +import { sendGuidePublicationNotifications } from "./notifications"; + +const event = { + aggregateId: "72df08ab-50dd-4cbd-9a69-70949d34cf9f", + eventType: "guide.published", topic: "notifications:guide", + payload: { publishedAt: "2026-10-08T10:00:00.000Z" }, +}; + +beforeEach(() => { mocks.rows.length = 0; mocks.filters.length = 0; mocks.create.mockReset(); }); + +describe("guide publication notifications", () => { + it("notifies all existing active users in batches with a retry-safe event key and guide link", async () => { + mocks.rows.push([{ name: "Amber", slug: "amber" }], + Array.from({ length: 100 }, (_, index) => ({ id: `user-${index.toString().padStart(3, "0")}` })), + [{ id: "user-100" }], []); + await sendGuidePublicationNotifications(event); + expect(mocks.create).toHaveBeenCalledTimes(2); + expect(mocks.create.mock.calls[0][0]).toHaveLength(100); + expect(mocks.create.mock.calls[1]).toEqual([ + [{ userId: "user-100", url: "/amber" }], + { eventKey: `guide:${event.aggregateId}:published`, kind: "guide_published", + title: "เผยแพร่ Guide ใหม่", body: "อ่าน Guide Amber ได้แล้ว" }, + ]); + const query = new PgDialect().sqlToQuery(mocks.filters[2] as Parameters[0]); + expect(query.sql).toContain('"banned" is not true'); + expect(query.sql).toContain('"created_at" <='); + expect(query.sql).toContain('"id" >'); + expect(query.params).toContain("user-099"); + }); + + it("skips a guide that is private, trashed, or deleted before delivery", async () => { + mocks.rows.push([]); + await sendGuidePublicationNotifications(event); + expect(mocks.create).not.toHaveBeenCalled(); + }); + + it("propagates delivery failures so the outbox retries", async () => { + mocks.rows.push([{ name: "Amber", slug: "amber" }], [{ id: "user" }]); + mocks.create.mockRejectedValueOnce(new Error("database unavailable")); + await expect(sendGuidePublicationNotifications(event)).rejects.toThrow("database unavailable"); + }); + + it("rejects malformed events before reading users", async () => { + await expect(sendGuidePublicationNotifications({ ...event, payload: {} })).rejects.toThrow("Invalid guide publication"); + await expect(sendGuidePublicationNotifications({ ...event, topic: "admin" })).rejects.toThrow("Invalid guide publication"); + expect(mocks.filters).toHaveLength(0); + }); +}); diff --git a/lib/guides/notifications.ts b/lib/guides/notifications.ts new file mode 100644 index 0000000..c592925 --- /dev/null +++ b/lib/guides/notifications.ts @@ -0,0 +1,36 @@ +import "server-only"; + +import { and, asc, eq, gt, isNull, lte, sql } from "drizzle-orm"; +import { getDb } from "@/db"; +import { guides, users } from "@/db/schema"; +import type { OutboxEventLike } from "@/lib/events/invalidation"; +import { createNotifications } from "@/lib/notifications/repository"; + +export async function sendGuidePublicationNotifications(event: OutboxEventLike) { + const publishedAt = typeof event.payload.publishedAt === "string" + ? new Date(event.payload.publishedAt) : null; + if (event.eventType !== "guide.published" || event.topic !== "notifications:guide" + || !publishedAt || Number.isNaN(publishedAt.getTime())) { + throw new Error("Invalid guide publication notification event"); + } + const db = getDb(); + const [guide] = await db.select({ name: guides.name, slug: guides.slug }).from(guides) + .where(and(eq(guides.id, event.aggregateId), eq(guides.isPublic, true), isNull(guides.trashedAt))).limit(1); + if (!guide) return; + + let cursor: string | undefined; + while (true) { + const recipients: { id: string }[] = await db.select({ id: users.id }).from(users) + .where(and(sql`${users.banned} is not true`, lte(users.createdAt, publishedAt), + cursor === undefined ? undefined : gt(users.id, cursor))) + .orderBy(asc(users.id)).limit(100); + if (!recipients.length) return; + await createNotifications(recipients.map(user => ({ userId: user.id, url: `/${guide.slug}` })), { + eventKey: `guide:${event.aggregateId}:published`, + kind: "guide_published", + title: "เผยแพร่ Guide ใหม่", + body: `อ่าน Guide ${guide.name} ได้แล้ว`, + }); + cursor = recipients[recipients.length - 1].id; + } +} diff --git a/lib/guides/visibility.test.ts b/lib/guides/visibility.test.ts new file mode 100644 index 0000000..728aeee --- /dev/null +++ b/lib/guides/visibility.test.ts @@ -0,0 +1,52 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { outboxEvents } from "@/db/schema"; + +const mocks = vi.hoisted(() => ({ previous: { isPublic: false, publishedAt: null as Date | null }, + updated: [] as unknown[], current: [] as unknown[], inserts: [] as { table: unknown; values: unknown }[] })); +vi.mock("server-only", () => ({})); +vi.mock("@/lib/audit-log", () => ({ writeAuditLog: vi.fn() })); +vi.mock("@/db", () => ({ getDb: () => ({ transaction: async (callback: (tx: unknown) => unknown) => { + const tx = { + select: () => { + const query = { from: () => query, where: () => query, + for: async () => [mocks.previous], limit: async () => mocks.current }; + return query; + }, + update: () => { const query = { set: () => query, where: () => query, returning: async () => mocks.updated }; return query; }, + insert: (table: unknown) => ({ values: async (values: unknown) => { mocks.inserts.push({ table, values }); } }), + }; + return callback(tx); +} }) })); + +import { GuideVersionConflictError, setGuideVisibility } from "./mutations"; + +const guide = { id: "guide-id", version: 2, name: "Amber", characterKey: "amber", publishedAt: new Date("2026-10-08T10:00:00Z") }; +const context = { actor: { id: "admin", label: "Admin" } }; +beforeEach(() => { + mocks.previous = { isPublic: false, publishedAt: null }; + mocks.updated = [guide]; mocks.current = []; mocks.inserts.length = 0; +}); +const publicationEvents = () => mocks.inserts.filter(row => row.table === outboxEvents && !Array.isArray(row.values)); + +describe("guide publication trigger", () => { + it("queues delivery in the publishing transaction", async () => { + await setGuideVisibility(guide.id, { isPublic: true, expectedVersion: 1 }, context); + expect(publicationEvents().map(row => row.values)).toEqual([{ + topic: "notifications:guide", aggregateId: guide.id, eventType: "guide.published", + payload: { publishedAt: guide.publishedAt.toISOString() }, + }]); + }); + + it.each(["unpublish", "republish", "already public"])("does not notify on %s", async (scenario) => { + if (scenario === "republish") mocks.previous.publishedAt = guide.publishedAt; + if (scenario === "already public") mocks.previous.isPublic = true; + await setGuideVisibility(guide.id, { isPublic: scenario !== "unpublish", expectedVersion: 1 }, context); + expect(publicationEvents()).toHaveLength(0); + }); + + it("preserves version conflicts without queuing notifications", async () => { + mocks.updated = []; mocks.current = [{ version: 5 }]; + await expect(setGuideVisibility(guide.id, { isPublic: true, expectedVersion: 1 }, context)).rejects.toBeInstanceOf(GuideVersionConflictError); + expect(mocks.inserts).toHaveLength(0); + }); +}); diff --git a/scripts/outbox-worker.ts b/scripts/outbox-worker.ts index 2f3cf06..6d8ec41 100644 --- a/scripts/outbox-worker.ts +++ b/scripts/outbox-worker.ts @@ -10,6 +10,7 @@ import { } from "@/lib/outbox/processor"; import { closeRedisClient, getRedisClient } from "@/lib/redis/client"; import { removeExpiredUnpaidCheckouts } from "@/lib/commission/cleanup"; +import { sendGuidePublicationNotifications } from "@/lib/guides/notifications"; let stopping = false; const stop = () => { @@ -58,7 +59,11 @@ async function run(): Promise { await Promise.all( events.map(async (event) => { try { - await processOutboxEvent(event, transport); + if (event.eventType === "guide.published") { + await sendGuidePublicationNotifications(event); + } else { + await processOutboxEvent(event, transport); + } await markOutboxProcessed(event.id); } catch (cause) { await releaseOutboxEvent(event.id, event.attempts, cause);