import "server-only"; import { notifyCommentChange } from "./events"; import { and, eq, inArray } from "drizzle-orm"; import { getDb } from "@/db"; import { commentAttachments, commentRevisions, commentRevisionAttachments, comments, commentThreads } from "@/db/schema"; import { inspectImage } from "@/lib/media/inspect"; import { getMediaStorage } from "@/lib/media/storage"; import { boundedBody, HttpError, withUploadSlot } from "@/lib/security/http"; import { authorizeComment, commentId, ensureCommentThread, getCommentTarget, withCommentLock } from "./repository"; import { MAX_COMMENT_BODY_BYTES, parseCommentForm } from "./validation"; import type { CommentViewer } from "./types"; import { withCommentSpamProtection } from "./spam"; type ParsedForm = ReturnType; type Upload = { id: string; objectKey: string; mimeType: string; byteSize: number }; type Writer = Parameters["transaction"]>[0]>[0]; async function uploadCommentImage(storage: Awaited>, image: Upload, bytes: Uint8Array) { for (let attempt = 0; ; attempt++) { try { // Retry the same object key, never the database mutation. await storage.write(image.objectKey, bytes, { type: image.mimeType, acl: "public-read", retry: 0 }); return; } catch (cause) { const error = cause as { code?: string; name?: string; statusCode?: number } | null; const transient = error && ( ["ECONNRESET", "ECONNREFUSED", "EPIPE", "ETIMEDOUT", "ConnectionClosed", "SocketClosed", "InternalError", "SlowDown", "ServiceUnavailable", "RequestTimeout"].includes(error.code ?? "") || error.name === "TimeoutError" || error.statusCode === 429 || (error.statusCode !== undefined && error.statusCode >= 500) ); if (!transient || attempt === 2) throw new HttpError(503, "comment-image-upload-failed"); await new Promise((resolve) => setTimeout(resolve, 200 * (attempt + 1))); } } } async function saveRevision(tx: Pick, id: string, version: number, form: ParsedForm, uploaded: Upload[]) { if (form.keep.length) { const own = await tx.select({ id: commentAttachments.id }).from(commentAttachments) .where(and(eq(commentAttachments.commentId, id), inArray(commentAttachments.id, form.keep))); if (own.length !== form.keep.length) throw new HttpError(400, "invalid-retained-image"); } if (uploaded.length) await tx.insert(commentAttachments).values(uploaded.map((image) => ({ ...image, commentId: id }))); const [revision] = await tx.insert(commentRevisions).values({ commentId: id, version, text: form.text }).returning({ id: commentRevisions.id }); const images = [...form.keep, ...uploaded.map((image) => image.id)]; if (images.length) await tx.insert(commentRevisionAttachments).values(images.map((attachmentId, position) => ({ revisionId: revision.id, attachmentId, position }))); } export async function publishComment(request: Request, viewer: CommentViewer, options: { target?: string; id?: string }) { if (!request.headers.get("content-type")?.startsWith("multipart/form-data;")) throw new HttpError(415, "expected-multipart"); const initial = options.id ? await authorizeComment(commentId(options.id), viewer) : null; if (initial && (initial.comment.authorId !== viewer.id || initial.comment.deletedAt || initial.comment.hidden || initial.rootHidden)) throw new HttpError(403, "comment-not-editable"); const destination = initial?.destination ?? await getCommentTarget(options.target!, viewer); const result = await withUploadSlot(async () => { const form = parseCommentForm(await boundedBody(request, MAX_COMMENT_BODY_BYTES).formData(), Boolean(options.id)); return withCommentSpamProtection(viewer.id, form.text, Boolean(options.id), async () => { const uploaded: Upload[] = []; try { const storage = form.files.length ? await getMediaStorage() : null; for (const file of form.files) { const bytes = new Uint8Array(await file.arrayBuffer()); await inspectImage(bytes, file.type as "image/png" | "image/jpeg" | "image/webp"); const image = { id: crypto.randomUUID(), objectKey: `comments/${crypto.randomUUID()}`, mimeType: file.type, byteSize: file.size }; uploaded.push(image); // Browser images load directly from the configured S3 public CDN. await uploadCommentImage(storage!, image, bytes); } if (options.id) { const version = await withCommentLock(options.id, viewer, async (tx, context) => { const c = context.comment; if (c.authorId !== viewer.id) throw new HttpError(403, "not-comment-author"); if (c.deletedAt || c.hidden || context.rootHidden) throw new HttpError(409, "comment-unavailable"); if (c.version !== form.version) throw new HttpError(409, "comment-edited-reload"); await saveRevision(tx, c.id, c.version + 1, form, uploaded); await tx.update(comments).set({ version: c.version + 1 }).where(eq(comments.id, c.id)); return c.version + 1; }); return { id: options.id, version }; } const thread = await ensureCommentThread(options.target!, viewer); return await getDb().transaction(async (tx) => { await tx.select({ id: commentThreads.id }).from(commentThreads).where(eq(commentThreads.id, thread.id)).for("update"); const target = await getCommentTarget(options.target!, viewer, tx); if (!target.writable) throw new HttpError(409, "guide-trashed"); let rootId: string | null = null; if (form.replyToId) { const parent = await authorizeComment(form.replyToId, viewer, tx); if (parent.comment.threadId !== thread.id) throw new HttpError(400, "cross-thread-reply"); if (parent.comment.deletedAt || parent.comment.hidden || parent.rootHidden) throw new HttpError(409, "comment-unavailable"); rootId = parent.comment.rootId ?? parent.comment.id; } const [comment] = await tx.insert(comments).values({ threadId: thread.id, authorId: viewer.id, rootId, replyToId: form.replyToId }).returning({ id: comments.id }); await saveRevision(tx, comment.id, 1, form, uploaded); return { id: comment.id, version: 1 }; }); } catch (cause) { if (uploaded.length) { const storage = await getMediaStorage(); for (const image of uploaded) await storage.delete(image.objectKey).catch(() => undefined); } throw cause; } }); }); await notifyCommentChange(destination.target); return result; }