diff --git a/packages/shared/src/database/schema.ts b/packages/shared/src/database/schema.ts index 86890e6..e5e3d8a 100644 --- a/packages/shared/src/database/schema.ts +++ b/packages/shared/src/database/schema.ts @@ -83,6 +83,12 @@ export const pgMessagesTable = pgTable( table.created_at, table.id, ), + guildAiStatusAnalyzedIdx: pgIndex("idx_messages_guild_ai_status_analyzed").on( + table.guild_id, + table.ai_status, + table.ai_analyzed_at, + table.id, + ), guildCreatedDeletedIdx: pgIndex("idx_messages_guild_created_deleted").on( table.guild_id, table.created_at, diff --git a/packages/shared/src/utils/index.ts b/packages/shared/src/utils/index.ts index 1204ede..4bf5b86 100644 --- a/packages/shared/src/utils/index.ts +++ b/packages/shared/src/utils/index.ts @@ -6,6 +6,40 @@ export function delay(ms: number): Promise { export * from "./pagination.js"; +// --------------------------------------------------------------------------- +// Centralized AbortController with guaranteed cleanup +// --------------------------------------------------------------------------- + +/** + * Creates an AbortController with a timeout that is ALWAYS cleaned up, + * even if the caller throws or returns early without calling clear(). + * + * Returns both the controller and a cleanup handle. + * + * Usage: + * const { controller, clear } = createAbortControllerWithTimeout(8000); + * try { + * const res = await fetch(url, { signal: controller.signal }); + * // ... work ... + * } finally { + * clear(); // guaranteed to clear the timeout + * } + */ +export function createAbortControllerWithTimeout( + timeoutMs: number, +): { controller: AbortController; clear: () => void } { + const controller = new AbortController(); + const timeoutId = setTimeout(() => controller.abort(), timeoutMs); + // Unref so the timeout doesn't keep the process alive + timeoutId?.unref?.(); + return { + controller, + clear: () => { + clearTimeout(timeoutId); + }, + }; +} + // --------------------------------------------------------------------------- // Retry with exponential backoff // --------------------------------------------------------------------------- diff --git a/services/backend/src/http/app.ts b/services/backend/src/http/app.ts index f83f9d7..32f10b6 100644 --- a/services/backend/src/http/app.ts +++ b/services/backend/src/http/app.ts @@ -19,6 +19,7 @@ import { createUiStateRouter } from "../modules/ui-state/ui-state.routes.js"; import { createGuildsRouter } from "../modules/voice/guilds.routes.js"; import { createVoiceRouter } from "../modules/voice/voice.routes.js"; import { + adminAuth, errorHandler, } from "../shared/middlewares/index.js"; import { config } from "../shared/config/index.js"; @@ -73,17 +74,18 @@ export function createHttpApp(): Express { app.use("/api", createConfigRouter()); app.use("/api", createDashboardRouter()); - // Protected routes — all routes are now public - app.use("/api", createMessagesRouter()); - app.use("/api", createAnalysisRouter()); - app.use("/api", createMascotChatRouter()); - app.use("/api", createMediaRouter()); - app.use("/api", createVoiceRouter()); - app.use("/api", createRecordingsRouter()); - app.use("/api", createUiStateRouter()); + // Protected routes — require admin authentication (X-Admin-Password header) + const adminAuthMiddleware = adminAuth(ADMIN_PASSWORD); + app.use("/api", adminAuthMiddleware, createMessagesRouter()); + app.use("/api", adminAuthMiddleware, createAnalysisRouter()); + app.use("/api", adminAuthMiddleware, createMascotChatRouter()); + app.use("/api", adminAuthMiddleware, createMediaRouter()); + app.use("/api", adminAuthMiddleware, createVoiceRouter()); + app.use("/api", adminAuthMiddleware, createRecordingsRouter()); + app.use("/api", adminAuthMiddleware, createUiStateRouter()); // Guilds routes - app.use("/api/guilds", createGuildsRouter()); + app.use("/api/guilds", adminAuthMiddleware, createGuildsRouter()); // 404 handler app.use((_req: Request, res: Response) => { diff --git a/services/backend/src/modules/auth/auth.routes.ts b/services/backend/src/modules/auth/auth.routes.ts index 10ec58f..94e5fcb 100644 --- a/services/backend/src/modules/auth/auth.routes.ts +++ b/services/backend/src/modules/auth/auth.routes.ts @@ -3,18 +3,22 @@ import { createChildLogger } from "@bete/shared/logger"; import type { Request, Response, Router } from "express"; import express from "express"; import { config } from "../../shared/config/index.js"; -import { asyncHandler } from "../../shared/middlewares/index.js"; +import { asyncHandler, rateLimit } from "../../shared/middlewares/index.js"; const logger = createChildLogger("auth.routes"); const adminPassword = config.ADMIN_PASSWORD || "admin"; +// Rate limit: max 10 login attempts per IP per 15 minutes +const loginRateLimit = rateLimit({ windowMs: 15 * 60 * 1000, max: 10 }); + export function createAuthRouter(): Router { const router = express.Router(); // POST /api/auth/login router.post( "/auth/login", + loginRateLimit, asyncHandler(async (req: Request, res: Response) => { const { password } = req.body as { password?: string }; diff --git a/services/backend/src/modules/dashboard/dashboard.repository.ts b/services/backend/src/modules/dashboard/dashboard.repository.ts index b7aef81..1759075 100644 --- a/services/backend/src/modules/dashboard/dashboard.repository.ts +++ b/services/backend/src/modules/dashboard/dashboard.repository.ts @@ -43,9 +43,10 @@ export class DashboardRepository { // Top channels by message count const topChannels = await pool.query(` SELECT channel_id, - (metadata::jsonb -> 'channel' ->> 'channelName') AS channel_name, + COALESCE(NULLIF((metadata::jsonb -> 'channel' ->> 'channelName'), ''), channel_id) AS channel_name, COUNT(*)::int AS message_count FROM messages + WHERE metadata IS NOT NULL AND metadata != '' GROUP BY channel_id, (metadata::jsonb -> 'channel' ->> 'channelName') ORDER BY COUNT(*) DESC LIMIT 10 diff --git a/services/backend/src/shared/middlewares/index.ts b/services/backend/src/shared/middlewares/index.ts index c0ac9a0..492bfb9 100644 --- a/services/backend/src/shared/middlewares/index.ts +++ b/services/backend/src/shared/middlewares/index.ts @@ -50,6 +50,48 @@ export function asyncHandler( }; } +/** + * Simple in-memory rate limiter (no external dependency). + * Tracks request counts per IP within a rolling window. + * Use for auth endpoints to prevent brute-force attacks. + */ +export function rateLimit(opts: { windowMs: number; max: number }) { + const { windowMs, max } = opts; + const hits = new Map(); + + // Periodic cleanup of stale entries to prevent unbounded memory growth + const cleanupInterval = setInterval(() => { + const now = Date.now(); + for (const [key, value] of hits) { + if (now >= value.resetAt) hits.delete(key); + } + }, windowMs * 2); + cleanupInterval.unref(); + + return (req: Request, res: Response, next: NextFunction) => { + const ip = req.ip ?? req.socket.remoteAddress ?? "unknown"; + const now = Date.now(); + const entry = hits.get(ip); + + if (!entry || now >= entry.resetAt) { + hits.set(ip, { count: 1, resetAt: now + windowMs }); + next(); + return; + } + + entry.count++; + if (entry.count > max) { + res.status(429).json({ + error: "TOO_MANY_REQUESTS", + message: `Rate limit exceeded. Try again in ${Math.ceil((entry.resetAt - now) / 1000)}s.`, + }); + return; + } + + next(); + }; +} + /** * Validate that a value is a non-empty string, or throw a descriptive error. * Use for both route params and query string values. diff --git a/services/discord-gateway/src/modules/ai-moderation/mediaAnalysisClient.ts b/services/discord-gateway/src/modules/ai-moderation/mediaAnalysisClient.ts index 460f820..80d3f01 100644 --- a/services/discord-gateway/src/modules/ai-moderation/mediaAnalysisClient.ts +++ b/services/discord-gateway/src/modules/ai-moderation/mediaAnalysisClient.ts @@ -11,7 +11,7 @@ import { readFile, writeFile, unlink, rm, mkdtemp } from "node:fs/promises"; import { tmpdir } from "node:os"; import path from "node:path"; import { promisify } from "node:util"; -import { delay } from "@bete/shared/utils"; +import { createAbortControllerWithTimeout, delay } from "@bete/shared/utils"; import { LRUCache } from "lru-cache"; import { config } from "../../shared/config/config.js"; import { resizeImageForVision } from "../attachment-upload/imageResizer.js"; @@ -297,8 +297,7 @@ async function downloadSingleAttachment( const urlToUse = att.uploaded_url ?? att.discord_url ?? null; if (!urlToUse) return; - const controller = new AbortController(); - const timeoutId = setTimeout(() => controller.abort(), 15000); + const { controller, clear } = createAbortControllerWithTimeout(15000); try { const res = await fetch(urlToUse, { signal: controller.signal }); if (!res.ok || !res.body) return; @@ -322,7 +321,35 @@ async function downloadSingleAttachment( await extractVideoFrames(att, imageBytes, targetId, maxDimension, imageMap); return; } - if (!sniffedMime) return; + + // Fallback: try attachment type metadata, then filename extension + let resolvedMime = sniffedMime; + if (!resolvedMime) { + if (att.type.startsWith("image/")) { + resolvedMime = att.type; + log.warn({ attachmentId: att.id, filename: att.filename, type: att.type }, + "Image MIME sniff failed — using attachment metadata type as fallback"); + } else { + // Last resort: check file extension + const ext = att.filename?.toLowerCase().split(".").pop(); + if (ext && ["jpg", "jpeg", "png", "gif", "webp", "bmp"].includes(ext)) { + const mimeMap: Record = { + jpg: "image/jpeg", jpeg: "image/jpeg", png: "image/png", + gif: "image/gif", webp: "image/webp", bmp: "image/bmp", + }; + resolvedMime = mimeMap[ext]; + log.warn({ attachmentId: att.id, filename: att.filename, ext }, + "Image MIME sniff failed — using file extension fallback"); + } + } + } + + // If all fallbacks fail, still try with generic image/jpeg (better than silent skip) + if (!resolvedMime) { + resolvedMime = "image/jpeg"; + log.warn({ attachmentId: att.id, filename: att.filename }, + "All MIME detection failed — forcing image/jpeg as last resort"); + } const { data: resizedBuffer, mimeType: resizedMime } = await resizeImageForVision(imageBytes, maxDimension); const dataUrl = `data:${resizedMime};base64,${resizedBuffer.toString("base64")}`; @@ -334,7 +361,7 @@ async function downloadSingleAttachment( } catch (err) { log.warn({ attachmentId: att.id, error: err instanceof Error ? err.message : String(err) }, "Download failed"); } finally { - clearTimeout(timeoutId); + clear(); } } @@ -404,6 +431,8 @@ async function downloadMediaCandidate( const existing = mediaAnalysisMap.get(targetId) ?? []; existing.push(`[Media analysis for message ${candidate.messageId}] ${candidate.label}: ${cached}`); mediaAnalysisMap.set(targetId, existing); + // Warm the LRU cache so subsequent calls in the same process skip DB query + visionLruCache.set(vck, cached); return; } } diff --git a/services/discord-gateway/src/modules/ai-moderation/searxngSearch.ts b/services/discord-gateway/src/modules/ai-moderation/searxngSearch.ts index b9e6b3b..0b05897 100644 --- a/services/discord-gateway/src/modules/ai-moderation/searxngSearch.ts +++ b/services/discord-gateway/src/modules/ai-moderation/searxngSearch.ts @@ -1,5 +1,6 @@ import Redis from "ioredis"; import { createChildLogger } from "@bete/shared/logger"; +import { createAbortControllerWithTimeout } from "@bete/shared/utils"; const log = createChildLogger("searxng-search"); @@ -68,43 +69,45 @@ export async function searchSearxng( // Cache miss — hit SearXNG API try { const url = `${SEARXNG_BASE_URL}/search?q=${encodeURIComponent(query)}&format=json&language=id&categories=${category}`; - const controller = new AbortController(); - const timeoutId = setTimeout(() => controller.abort(), TIMEOUT_MS); + const { controller, clear } = createAbortControllerWithTimeout(TIMEOUT_MS); - const response = await fetch(url, { - signal: controller.signal, - headers: { - Accept: "application/json", - "User-Agent": - "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36", - }, - }); - clearTimeout(timeoutId); - - if (!response.ok) { - log.warn({ status: response.status, query }, "SearXNG search failed"); - return []; - } - - const data = (await response.json()) as { - results?: Array<{ title?: string; url?: string; content?: string }>; - }; - const results = data.results ?? []; - const mapped = results.slice(0, MAX_RESULTS).map((r) => ({ - title: r.title ?? "", - url: r.url ?? "", - snippet: (r.content ?? "").slice(0, 500), - })); - - // Store in cache (fire and forget — don't block on write) - if (redis) { - redis.setex(cacheKey, CACHE_TTL, JSON.stringify(mapped)).catch(() => { - // Cache write failed silently + try { + const response = await fetch(url, { + signal: controller.signal, + headers: { + Accept: "application/json", + "User-Agent": + "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36", + }, }); - } - log.debug({ query, category, resultCount: mapped.length }, "SearXNG search OK"); - return mapped; + if (!response.ok) { + log.warn({ status: response.status, query }, "SearXNG search failed"); + return []; + } + + const data = (await response.json()) as { + results?: Array<{ title?: string; url?: string; content?: string }>; + }; + const results = data.results ?? []; + const mapped = results.slice(0, MAX_RESULTS).map((r) => ({ + title: r.title ?? "", + url: r.url ?? "", + snippet: (r.content ?? "").slice(0, 500), + })); + + // Store in cache (fire and forget — don't block on write) + if (redis) { + redis.setex(cacheKey, CACHE_TTL, JSON.stringify(mapped)).catch(() => { + // Cache write failed silently + }); + } + + log.debug({ query, category, resultCount: mapped.length }, "SearXNG search OK"); + return mapped; + } finally { + clear(); + } } catch (err) { log.warn( { error: err instanceof Error ? err.message : String(err), query }, diff --git a/services/discord-gateway/src/modules/ai-moderation/urlFetcher.ts b/services/discord-gateway/src/modules/ai-moderation/urlFetcher.ts index c49c0d9..3b8ff18 100644 --- a/services/discord-gateway/src/modules/ai-moderation/urlFetcher.ts +++ b/services/discord-gateway/src/modules/ai-moderation/urlFetcher.ts @@ -1,6 +1,7 @@ import { resolve } from "node:dns/promises"; import { isIP } from "node:net"; import { createChildLogger } from "@bete/shared/logger"; +import { createAbortControllerWithTimeout } from "@bete/shared/utils"; const log = createChildLogger("urlFetcher"); @@ -112,8 +113,7 @@ export async function fetchUrlSafely( return { url, type: "error", error: "Unsafe URL blocked" }; } - const controller = new AbortController(); - const timeoutId = setTimeout(() => controller.abort(), FETCH_TIMEOUT_MS); + const { controller, clear } = createAbortControllerWithTimeout(FETCH_TIMEOUT_MS); try { const response = await fetch(url, { @@ -190,7 +190,7 @@ export async function fetchUrlSafely( error: err instanceof Error ? err.message : String(err), }; } finally { - clearTimeout(timeoutId); + clear(); } }