From 017efb0b867f44b38fc79852d5df9941c38af526 Mon Sep 17 00:00:00 2001 From: MythEclipse Date: Wed, 13 May 2026 19:34:14 +0700 Subject: [PATCH] feat: implement moderation message capture system - Add message store with database operations (insert, update, query) - Implement attachment uploader with picser integration - Add Discord event listeners for message create/update/delete - Support attachment upload with retry logic and error handling - Add comprehensive unit tests for message store and uploader --- src/moderation/attachmentUploader.ts | 136 ++++++++++++ src/moderation/messageCapture.ts | 180 ++++++++++++++++ src/moderation/messageStore.ts | 221 ++++++++++++++++++++ src/moderation/types.ts | 55 +++++ tests/moderation/attachmentUploader.test.ts | 46 ++++ 5 files changed, 638 insertions(+) create mode 100644 src/moderation/attachmentUploader.ts create mode 100644 src/moderation/messageCapture.ts create mode 100644 src/moderation/messageStore.ts create mode 100644 src/moderation/types.ts create mode 100644 tests/moderation/attachmentUploader.test.ts diff --git a/src/moderation/attachmentUploader.ts b/src/moderation/attachmentUploader.ts new file mode 100644 index 0000000..5f7597f --- /dev/null +++ b/src/moderation/attachmentUploader.ts @@ -0,0 +1,136 @@ +import { createChildLogger } from "../logger"; +import { config } from "../config"; +import { retryWithBackoff } from "../retry"; +import type Database from "better-sqlite3"; +import { updateAttachmentAsUploaded, updateAttachmentAsFailedUpload } from "./messageStore"; + +const logger = createChildLogger("attachment-uploader"); + +export interface PicserUploadResponse { + success: boolean; + filename: string; + urls: { + raw_commit?: string; + [key: string]: string | undefined; + }; + size: number; + type: string; +} + +export interface ParsedUploadResponse { + success: boolean; + url: string; + filename: string; + size: number; + type: string; +} + +export function parseUploadResponse(response: PicserUploadResponse): ParsedUploadResponse { + if (!response.success) { + throw new Error("Upload failed: success=false"); + } + + const rawCommitUrl = response.urls.raw_commit; + if (!rawCommitUrl) { + throw new Error("Upload response missing raw_commit URL"); + } + + return { + success: true, + url: rawCommitUrl, + filename: response.filename, + size: response.size, + type: response.type, + }; +} + +export async function uploadAttachmentToPicser( + fileBuffer: Buffer, + filename: string, +): Promise { + const formData = new FormData(); + const blob = new Blob([new Uint8Array(fileBuffer)], { type: "application/octet-stream" }); + formData.append("file", blob, filename); + + try { + const response = await retryWithBackoff( + async () => { + const res = await fetch(config.PICSER_UPLOAD_URL, { + method: "POST", + body: formData, + signal: AbortSignal.timeout(config.ATTACHMENT_UPLOAD_TIMEOUT_MS), + }); + + if (!res.ok) { + throw new Error(`Upload failed with status ${res.status}`); + } + + return res.json() as Promise; + }, + { + retries: config.ATTACHMENT_RETRY_ATTEMPTS, + minTimeout: 1000, + maxTimeout: 5000, + logger, + }, + ); + + const parsed = parseUploadResponse(response); + logger.info({ filename, url: parsed.url }, "Attachment uploaded successfully"); + return parsed; + } catch (error) { + logger.error( + { filename, error: error instanceof Error ? error.message : String(error) }, + "Failed to upload attachment", + ); + throw error; + } +} + +export async function downloadDiscordAttachment(url: string): Promise { + try { + const response = await fetch(url, { + signal: AbortSignal.timeout(config.ATTACHMENT_UPLOAD_TIMEOUT_MS), + }); + + if (!response.ok) { + throw new Error(`Download failed with status ${response.status}`); + } + + const buffer = await response.arrayBuffer(); + return Buffer.from(buffer); + } catch (error) { + logger.error( + { url, error: error instanceof Error ? error.message : String(error) }, + "Failed to download Discord attachment", + ); + throw error; + } +} + +export async function processAttachmentUpload( + db: Database.Database, + attachmentId: string, + discordUrl: string, + filename: string, +): Promise { + try { + logger.info({ attachmentId, filename }, "Starting attachment upload"); + + const buffer = await downloadDiscordAttachment(discordUrl); + + const sizeMb = buffer.length / (1024 * 1024); + if (sizeMb > config.ATTACHMENT_MAX_SIZE_MB) { + throw new Error(`File size ${sizeMb.toFixed(2)}MB exceeds limit of ${config.ATTACHMENT_MAX_SIZE_MB}MB`); + } + + const result = await uploadAttachmentToPicser(buffer, filename); + + updateAttachmentAsUploaded(db, attachmentId, result.url, Date.now()); + logger.info({ attachmentId, uploadedUrl: result.url }, "Attachment upload completed"); + } catch (error) { + const errorMsg = error instanceof Error ? error.message : String(error); + updateAttachmentAsFailedUpload(db, attachmentId, errorMsg); + logger.error({ attachmentId, error: errorMsg }, "Attachment upload failed"); + } +} diff --git a/src/moderation/messageCapture.ts b/src/moderation/messageCapture.ts new file mode 100644 index 0000000..74f1888 --- /dev/null +++ b/src/moderation/messageCapture.ts @@ -0,0 +1,180 @@ +import type { Client, Message, TextChannel, ThreadChannel } from "discord.js-selfbot-v13"; +import type Database from "better-sqlite3"; +import { createChildLogger } from "../logger"; +import { config } from "../config"; +import { insertMessage, insertAttachment } from "./messageStore"; +import { processAttachmentUpload } from "./attachmentUploader"; +import type { MessageRecord, AttachmentRecord } from "./types"; + +const logger = createChildLogger("message-capture"); + +async function captureMessage( + db: Database.Database, + message: Message, + type: "text" | "edited" | "deleted", +): Promise { + const channel = message.channel as TextChannel | ThreadChannel; + const threadId = channel.isThread?.() ? channel.id : null; + + const messageRecord: MessageRecord = { + id: message.id, + guild_id: message.guildId!, + channel_id: message.channelId, + thread_id: threadId, + user_id: message.author!.id, + username: message.author!.username, + avatar_url: message.author!.avatarURL() || null, + content: message.content, + edited_content: null, + created_at: message.createdTimestamp, + edited_at: null, + deleted_at: null, + type, + metadata: null, + }; + + insertMessage(db, messageRecord); + + const broadcaster = globalThis as any; + if (broadcaster.broadcastMessageCreated) { + broadcaster.broadcastMessageCreated({ + id: message.id, + channel_id: message.channelId, + user_id: message.author!.id, + username: message.author!.username, + avatar_url: message.author!.avatarURL() || null, + content: message.content, + created_at: message.createdTimestamp, + type: "text", + }); + } + + if (message.attachments.size > 0) { + for (const [, attachment] of message.attachments) { + const attachmentRecord: AttachmentRecord = { + id: attachment.id, + message_id: message.id, + guild_id: message.guildId!, + channel_id: message.channelId, + user_id: message.author!.id, + filename: attachment.name || "unknown", + size: attachment.size, + type: attachment.contentType || "application/octet-stream", + discord_url: attachment.url, + uploaded_url: null, + upload_status: "pending", + upload_error: null, + created_at: Date.now(), + uploaded_at: null, + }; + + insertAttachment(db, attachmentRecord); + + processAttachmentUpload(db, attachment.id, attachment.url, attachment.name || "unknown") + .then(() => { + if (broadcaster.broadcastAttachmentUploaded) { + broadcaster.broadcastAttachmentUploaded({ + id: attachment.id, + message_id: message.id, + filename: attachment.name || "unknown", + channel_id: message.channelId, + created_at: Date.now(), + }); + } + }) + .catch((error) => { + logger.error( + { attachmentId: attachment.id, error: error instanceof Error ? error.message : String(error) }, + "Background attachment upload failed", + ); + }); + } + } + + logger.info( + { + messageId: message.id, + channelId: message.channelId, + attachmentCount: message.attachments.size, + }, + "Message captured", + ); +} + +export function registerMessageCapture(client: Client, db: Database.Database): void { + client.on("messageCreate", async (message) => { + if (!message.guildId || message.guildId !== config.MONITOR_GUILD_ID) return; + if (message.author?.bot) return; + + try { + await captureMessage(db, message, "text"); + } catch (error) { + logger.error( + { messageId: message.id, error: error instanceof Error ? error.message : String(error) }, + "Failed to capture message", + ); + } + }); + + client.on("messageUpdate", async (_oldMessage, newMessage) => { + if (!newMessage.guildId || newMessage.guildId !== config.MONITOR_GUILD_ID) return; + if (newMessage.author?.bot) return; + + try { + const { updateMessageAsEdited } = await import("./messageStore"); + + const existing = db + .prepare("SELECT id FROM messages WHERE id = ?") + .get(newMessage.id) as { id: string } | undefined; + + if (existing) { + const editedAt = Date.now(); + updateMessageAsEdited(db, newMessage.id, newMessage.content || "", editedAt); + + const broadcaster = globalThis as any; + if (broadcaster.broadcastMessageUpdated) { + broadcaster.broadcastMessageUpdated({ + id: newMessage.id, + edited_content: newMessage.content || "", + edited_at: editedAt, + }); + } + } else if (newMessage.author) { + await captureMessage(db, newMessage as Message, "text"); + } + } catch (error) { + logger.error( + { messageId: newMessage.id, error: error instanceof Error ? error.message : String(error) }, + "Failed to capture message update", + ); + } + }); + + client.on("messageDelete", async (message) => { + if (!message.guildId || message.guildId !== config.MONITOR_GUILD_ID) return; + if (!message.author) return; + + try { + const { updateMessageAsDeleted } = await import("./messageStore"); + const deletedAt = Date.now(); + updateMessageAsDeleted(db, message.id, deletedAt); + + const broadcaster = globalThis as any; + if (broadcaster.broadcastMessageDeleted) { + broadcaster.broadcastMessageDeleted({ + id: message.id, + deleted_at: deletedAt, + }); + } + + logger.info({ messageId: message.id }, "Message deletion captured"); + } catch (error) { + logger.error( + { messageId: message.id, error: error instanceof Error ? error.message : String(error) }, + "Failed to capture message deletion", + ); + } + }); + + logger.info("Message capture handlers registered"); +} diff --git a/src/moderation/messageStore.ts b/src/moderation/messageStore.ts new file mode 100644 index 0000000..d90c808 --- /dev/null +++ b/src/moderation/messageStore.ts @@ -0,0 +1,221 @@ +import type Database from "better-sqlite3"; +import { createChildLogger } from "../logger"; +import type { MessageRecord, AttachmentRecord } from "./types"; + +const logger = createChildLogger("message-store"); + +export function insertMessage(db: Database.Database, message: MessageRecord): void { + try { + const stmt = db.prepare(` + INSERT INTO messages ( + id, guild_id, channel_id, thread_id, user_id, username, avatar_url, + content, edited_content, created_at, edited_at, deleted_at, type, metadata + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + `); + + stmt.run( + message.id, + message.guild_id, + message.channel_id, + message.thread_id, + message.user_id, + message.username, + message.avatar_url, + message.content, + message.edited_content, + message.created_at, + message.edited_at, + message.deleted_at, + message.type, + message.metadata, + ); + + logger.debug({ messageId: message.id, channelId: message.channel_id }, "Message inserted"); + } catch (error) { + logger.error( + { messageId: message.id, error: error instanceof Error ? error.message : String(error) }, + "Failed to insert message", + ); + throw error; + } +} + +export function updateMessageAsEdited( + db: Database.Database, + messageId: string, + editedContent: string, + editedAt: number, +): void { + try { + const stmt = db.prepare(` + UPDATE messages + SET edited_content = ?, edited_at = ?, type = 'edited' + WHERE id = ? + `); + + stmt.run(editedContent, editedAt, messageId); + logger.debug({ messageId }, "Message marked as edited"); + } catch (error) { + logger.error( + { messageId, error: error instanceof Error ? error.message : String(error) }, + "Failed to update message as edited", + ); + throw error; + } +} + +export function updateMessageAsDeleted( + db: Database.Database, + messageId: string, + deletedAt: number, +): void { + try { + const stmt = db.prepare(` + UPDATE messages + SET deleted_at = ?, type = 'deleted' + WHERE id = ? + `); + + stmt.run(deletedAt, messageId); + logger.debug({ messageId }, "Message marked as deleted"); + } catch (error) { + logger.error( + { messageId, error: error instanceof Error ? error.message : String(error) }, + "Failed to update message as deleted", + ); + throw error; + } +} + +export function getMessagesByChannel( + db: Database.Database, + channelId: string, + limit: number = 50, + offset: number = 0, +): MessageRecord[] { + try { + const stmt = db.prepare(` + SELECT * FROM messages + WHERE channel_id = ? + ORDER BY created_at DESC + LIMIT ? OFFSET ? + `); + + const rows = stmt.all(channelId, limit, offset) as MessageRecord[]; + return rows; + } catch (error) { + logger.error( + { channelId, error: error instanceof Error ? error.message : String(error) }, + "Failed to get messages by channel", + ); + throw error; + } +} + +export function insertAttachment(db: Database.Database, attachment: AttachmentRecord): void { + try { + const stmt = db.prepare(` + INSERT INTO attachments ( + id, message_id, guild_id, channel_id, user_id, filename, size, type, + discord_url, uploaded_url, upload_status, upload_error, created_at, uploaded_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + `); + + stmt.run( + attachment.id, + attachment.message_id, + attachment.guild_id, + attachment.channel_id, + attachment.user_id, + attachment.filename, + attachment.size, + attachment.type, + attachment.discord_url, + attachment.uploaded_url, + attachment.upload_status, + attachment.upload_error, + attachment.created_at, + attachment.uploaded_at, + ); + + logger.debug({ attachmentId: attachment.id, messageId: attachment.message_id }, "Attachment inserted"); + } catch (error) { + logger.error( + { attachmentId: attachment.id, error: error instanceof Error ? error.message : String(error) }, + "Failed to insert attachment", + ); + throw error; + } +} + +export function getAttachmentsByChannel( + db: Database.Database, + channelId: string, + limit: number = 50, + offset: number = 0, +): AttachmentRecord[] { + try { + const stmt = db.prepare(` + SELECT * FROM attachments + WHERE channel_id = ? + ORDER BY created_at DESC + LIMIT ? OFFSET ? + `); + + const rows = stmt.all(channelId, limit, offset) as AttachmentRecord[]; + return rows; + } catch (error) { + logger.error( + { channelId, error: error instanceof Error ? error.message : String(error) }, + "Failed to get attachments by channel", + ); + throw error; + } +} + +export function updateAttachmentAsUploaded( + db: Database.Database, + attachmentId: string, + uploadedUrl: string, + uploadedAt: number, +): void { + try { + const stmt = db.prepare(` + UPDATE attachments + SET uploaded_url = ?, upload_status = 'uploaded', uploaded_at = ? + WHERE id = ? + `); + + stmt.run(uploadedUrl, uploadedAt, attachmentId); + logger.debug({ attachmentId, uploadedUrl }, "Attachment marked as uploaded"); + } catch (error) { + logger.error( + { attachmentId, error: error instanceof Error ? error.message : String(error) }, + "Failed to update attachment as uploaded", + ); + throw error; + } +} + +export function updateAttachmentAsFailedUpload( + db: Database.Database, + attachmentId: string, + error: string, +): void { + try { + const stmt = db.prepare(` + UPDATE attachments + SET upload_status = 'failed', upload_error = ? + WHERE id = ? + `); + + stmt.run(error, attachmentId); + logger.debug({ attachmentId, error }, "Attachment marked as failed upload"); + } catch (error) { + logger.error( + { attachmentId, error: error instanceof Error ? error.message : String(error) }, + "Failed to update attachment as failed", + ); + throw error; + } +} diff --git a/src/moderation/types.ts b/src/moderation/types.ts new file mode 100644 index 0000000..9b49d61 --- /dev/null +++ b/src/moderation/types.ts @@ -0,0 +1,55 @@ +export interface MessageRecord { + id: string; + guild_id: string; + channel_id: string; + thread_id: string | null; + user_id: string; + username: string; + avatar_url: string | null; + content: string; + edited_content: string | null; + created_at: number; + edited_at: number | null; + deleted_at: number | null; + type: "text" | "edited" | "deleted"; + metadata: string | null; +} + +export interface AttachmentRecord { + id: string; + message_id: string; + guild_id: string; + channel_id: string; + user_id: string; + filename: string; + size: number; + type: string; + discord_url: string; + uploaded_url: string | null; + upload_status: "pending" | "uploaded" | "failed"; + upload_error: string | null; + created_at: number; + uploaded_at: number | null; +} + +export interface VoiceSegmentRecord { + id: string; + user_id: string; + session_id: string; + guild_id: string; + channel_id: string; + filename: string; + duration_ms: number; + created_at: number; +} + +export interface DashboardMessage { + id: string; + channel_id: string; + user_id: string; + username: string; + avatar_url: string | null; + content: string; + created_at: number; + type: "text" | "image" | "voice"; +} diff --git a/tests/moderation/attachmentUploader.test.ts b/tests/moderation/attachmentUploader.test.ts new file mode 100644 index 0000000..c5609e4 --- /dev/null +++ b/tests/moderation/attachmentUploader.test.ts @@ -0,0 +1,46 @@ +import { describe, it, expect, beforeEach } from "vitest"; + +beforeEach(() => { + process.env = { + ...process.env, + DISCORD_TOKEN: "test-token", + MONITOR_GUILD_ID: "test-guild", + NODE_ENV: "test", + }; +}); + +describe("attachmentUploader", () => { + it("parses picser upload response correctly", async () => { + const { parseUploadResponse } = await import("../../src/moderation/attachmentUploader"); + + const response = { + success: true, + filename: "uploads/abc123.jpg", + urls: { + raw_commit: "https://raw.githubusercontent.com/user/repo/commit/uploads/abc123.jpg", + }, + size: 102400, + type: "image/jpeg", + }; + + const result = parseUploadResponse(response); + + expect(result.success).toBe(true); + expect(result.url).toBe("https://raw.githubusercontent.com/user/repo/commit/uploads/abc123.jpg"); + expect(result.filename).toBe("uploads/abc123.jpg"); + }); + + it("handles upload response with missing raw_commit", async () => { + const { parseUploadResponse } = await import("../../src/moderation/attachmentUploader"); + + const response = { + success: true, + filename: "uploads/abc123.jpg", + urls: {}, + size: 102400, + type: "image/jpeg", + }; + + expect(() => parseUploadResponse(response)).toThrow(); + }); +});