fix(moderation): remove manual reanalyze triggers — auto-recovery only

Manual per-message and batch reanalyze buttons/endpoints let anyone
re-queue arbitrary messages for LLM analysis, burning AI credits on
spam. Removed:
- FE: Reanalyze buttons in message list, search panel, and messages page
- FE: useReanalyze/useReanalyzeBatch hooks + messagesApi methods
- BE: POST /api/messages/:id/reanalyze and /reanalyze-batch endpoints
- BE: markForReanalysis/reanalyzeErrorBatch service+repository methods

Recovery of failed messages is fully automatic: the discord-gateway
startPendingAIAnalysisWorker retries 'pending' (batch path) and
'error/analysis_incomplete' (individual path) messages on
AI_ANALYSIS_RECOVERY_INTERVAL_MS.
This commit is contained in:
asepharyana
2026-08-04 15:37:56 +07:00
parent 88484f12a9
commit 2f51f94610
11 changed files with 6 additions and 243 deletions
@@ -6,7 +6,6 @@ import {
isNull, isNull,
like, like,
lt, lt,
ne,
notInArray, notInArray,
or, or,
type SQL, type SQL,
@@ -223,58 +222,6 @@ export class MessagesRepository {
return mapMessageRow(row as Record<string, unknown>); return mapMessageRow(row as Record<string, unknown>);
} }
/**
* Bulk-reset ai_status from 'error' to 'pending' so the DG recovery worker
* picks them up on its next poll cycle.
*
* Accepts optional scope filters (guildId, channelId) or a list of explicit
* message IDs. Returns the count of rows that were actually updated.
*/
async reanalyzeErrorBatch(opts: {
guildId?: string;
channelId?: string;
messageIds?: string[];
}): Promise<number> {
const db = getDatabase();
const conditions: SQL[] = [eq(pgMessagesTable.ai_status, "error")];
if (opts.messageIds && opts.messageIds.length > 0) {
conditions.push(inArray(pgMessagesTable.id, opts.messageIds));
}
if (opts.guildId) {
conditions.push(eq(pgMessagesTable.guild_id, opts.guildId));
}
if (opts.channelId) {
conditions.push(eq(pgMessagesTable.channel_id, opts.channelId));
}
const result = await db
.update(pgMessagesTable)
.set({ ai_status: "pending" })
.where(and(...conditions));
const count = result.rowCount ?? 0;
logger.info({ count, ...opts }, "Batch reanalyze triggered");
return count;
}
/**
* Mark a single message for re-analysis by resetting ai_status to 'pending'.
* Skips messages already in 'pending' state to avoid write amplification.
*/
async markForReanalysis(id: string): Promise<void> {
const db = getDatabase();
await db
.update(pgMessagesTable)
.set({ ai_status: "pending" })
.where(
and(
eq(pgMessagesTable.id, id),
ne(pgMessagesTable.ai_status, "pending"),
),
);
}
/** /**
* Retrieve messages flagged for review (ai_status IN ('warn', 'flagged')). * Retrieve messages flagged for review (ai_status IN ('warn', 'flagged')).
* Optionally filtered by channelId, with configurable limit. * Optionally filtered by channelId, with configurable limit.
@@ -1,7 +1,7 @@
import type { Request, Response, Router } from "express"; import type { Request, Response, Router } from "express";
import express from "express"; import express from "express";
import { createChildLogger } from "@/shared/logger/index"; import { createChildLogger } from "@/shared/logger/index";
import { asyncHandler, validateBody } from "../../shared/middlewares/index.js"; import { asyncHandler } from "../../shared/middlewares/index.js";
import { import {
handleGetAttachmentsByChannel, handleGetAttachmentsByChannel,
handleGetImageMessages, handleGetImageMessages,
@@ -9,27 +9,10 @@ import {
handleGetMessagesByChannel, handleGetMessagesByChannel,
handleListMessages, handleListMessages,
} from "./messages.controller.js"; } from "./messages.controller.js";
import { reanalyzeBatchSchema } from "./messages.schema.js";
import { messagesService } from "./messages.service.js"; import { messagesService } from "./messages.service.js";
const logger = createChildLogger("messages.routes"); const logger = createChildLogger("messages.routes");
/**
* Per-message in-flight guard for the single reanalyze endpoint.
* Prevents concurrent spam-clicks from issuing duplicate UPDATE + recovery
* worker triggers for the same message.
*/
const reanalyzeInFlight = new Set<string>();
/**
* Per-scope in-flight guard for the batch reanalyze endpoint.
* Scope key = "guildId:channelId" (empty string used for undefined parts).
* Two concurrent batch-reanalyze requests for the same scope are rejected
* with 409 so the recovery worker is not triggered multiple times for the
* same set of error messages.
*/
const reanalyzeBatchInFlight = new Set<string>();
export function createMessagesRouter(): Router { export function createMessagesRouter(): Router {
const router = express.Router(); const router = express.Router();
@@ -51,74 +34,6 @@ export function createMessagesRouter(): Router {
// (uses /detail/ prefix to avoid collision with :channelId route above) // (uses /detail/ prefix to avoid collision with :channelId route above)
router.get("/messages/detail/:id", handleGetMessageById); router.get("/messages/detail/:id", handleGetMessageById);
// POST /api/messages/reanalyze-batch — Bulk retry all errored messages
// MUST be registered BEFORE /messages/:id/reanalyze so "reanalyze-batch"
// is not captured as an :id param.
router.post(
"/messages/reanalyze-batch",
validateBody(reanalyzeBatchSchema),
asyncHandler(async (req: Request, res: Response) => {
const { guildId, channelId, messageIds } = req.body as {
guildId?: string;
channelId?: string;
messageIds?: string[];
};
// Idempotency guard: one concurrent batch-reanalyze per scope.
// Prevents two admin sessions clicking simultaneously from each
// triggering the recovery worker for the same set of messages.
const scopeKey = `${guildId ?? ""}:${channelId ?? ""}`;
if (reanalyzeBatchInFlight.has(scopeKey)) {
res
.status(409)
.json({ error: "REANALYZE_BATCH_IN_PROGRESS", scope: scopeKey });
return;
}
reanalyzeBatchInFlight.add(scopeKey);
let count = 0;
try {
count = await messagesService.reanalyzeErrorBatch({
guildId,
channelId,
messageIds,
});
} finally {
reanalyzeBatchInFlight.delete(scopeKey);
}
logger.info({ count, guildId, channelId }, "Batch reanalyze completed");
res.status(200).json({ ok: true, count });
}),
);
// POST /api/messages/:id/reanalyze - Mark single message for re-analysis
router.post(
"/messages/:id/reanalyze",
asyncHandler(async (req: Request, res: Response) => {
const id = String(req.params.id ?? "");
if (!id) {
res.status(400).json({ error: "MISSING_ID" });
return;
}
// Idempotency guard: reject concurrent duplicate requests for the same ID.
if (reanalyzeInFlight.has(id)) {
res.status(409).json({ error: "REANALYZE_IN_PROGRESS", messageId: id });
return;
}
reanalyzeInFlight.add(id);
try {
await messagesService.markForReanalysis(id);
} finally {
reanalyzeInFlight.delete(id);
}
res.status(200).json({ ok: true });
}),
);
// GET /api/review - Get flagged/warned messages for review // GET /api/review - Get flagged/warned messages for review
router.get( router.get(
"/review", "/review",
@@ -38,13 +38,6 @@ export const messageUpdateSchema = z.object({
aiConfidence: z.number().optional(), aiConfidence: z.number().optional(),
}); });
export const reanalyzeBatchSchema = z.object({
guildId: z.string().optional(),
channelId: z.string().optional(),
messageIds: z.array(z.string()).optional(),
});
export type MessageQuery = z.infer<typeof messageQuerySchema>; export type MessageQuery = z.infer<typeof messageQuerySchema>;
export type MessageCreate = z.infer<typeof messageCreateSchema>; export type MessageCreate = z.infer<typeof messageCreateSchema>;
export type MessageUpdate = z.infer<typeof messageUpdateSchema>; export type MessageUpdate = z.infer<typeof messageUpdateSchema>;
export type ReanalyzeBatchInput = z.infer<typeof reanalyzeBatchSchema>;
@@ -58,15 +58,6 @@ export class MessagesService {
return messagesRepository.getImageMessages(guildId, limit); return messagesRepository.getImageMessages(guildId, limit);
} }
async markForReanalysis(id: string): Promise<void> {
if (!id) {
throw new ValidationError("message ID is required");
}
logger.debug({ id }, "Marking message for re-analysis");
await messagesRepository.markForReanalysis(id);
}
async getReviewMessages( async getReviewMessages(
channelId?: string, channelId?: string,
limit?: number, limit?: number,
@@ -74,25 +65,6 @@ export class MessagesService {
logger.debug({ channelId, limit }, "Getting review messages"); logger.debug({ channelId, limit }, "Getting review messages");
return messagesRepository.getReviewMessages(channelId, limit); return messagesRepository.getReviewMessages(channelId, limit);
} }
async reanalyzeErrorBatch(opts: {
guildId?: string;
channelId?: string;
messageIds?: string[];
}) {
if (
!opts.guildId &&
!opts.channelId &&
(!opts.messageIds || opts.messageIds.length === 0)
) {
throw new ValidationError(
"At least one of guildId, channelId, or messageIds[] is required",
);
}
logger.info(opts, "Batch reanalyzing errored messages");
return messagesRepository.reanalyzeErrorBatch(opts);
}
} }
export const messagesService = new MessagesService(); export const messagesService = new MessagesService();
@@ -1,6 +1,6 @@
"use client"; "use client";
import { Flag, Image, Loader2, RefreshCw, Search } from "lucide-react"; import { Flag, Image, Loader2, Search } from "lucide-react";
import { useRouter, useSearchParams } from "next/navigation"; import { useRouter, useSearchParams } from "next/navigation";
import { useCallback, useEffect, useState } from "react"; import { useCallback, useEffect, useState } from "react";
import { GlassCard } from "@/components/glass/card"; import { GlassCard } from "@/components/glass/card";
@@ -13,7 +13,6 @@ import { MessageList } from "@/components/messages/message-list";
import { SearchOverlay } from "@/components/messages/search-overlay"; import { SearchOverlay } from "@/components/messages/search-overlay";
import { EmptyState, ErrorState, LoadingSkeleton } from "@/components/shared"; import { EmptyState, ErrorState, LoadingSkeleton } from "@/components/shared";
import { GuildSelector } from "@/components/shared/guild-selector"; import { GuildSelector } from "@/components/shared/guild-selector";
import { Button } from "@/components/ui/button";
import { import {
Select, Select,
SelectContent, SelectContent,
@@ -28,8 +27,6 @@ import {
useMessages, useMessages,
useMessagesHasMore, useMessagesHasMore,
useMessagesWsSync, useMessagesWsSync,
useReanalyze,
useReanalyzeBatch,
useReview, useReview,
useTextChannels, useTextChannels,
} from "@/hooks"; } from "@/hooks";
@@ -74,8 +71,6 @@ export default function MessagesPage() {
const loadMoreMut = useLoadMore(); const loadMoreMut = useLoadMore();
const { data: images } = useImages(guildId); const { data: images } = useImages(guildId);
const { data: reviews } = useReview(selectedChannel || undefined); const { data: reviews } = useReview(selectedChannel || undefined);
const reanalyzeMut = useReanalyze();
const reanalyzeBatchMut = useReanalyzeBatch();
const { const {
message: detailMessage, message: detailMessage,
@@ -164,14 +159,6 @@ export default function MessagesPage() {
&#8984;K &#8984;K
</span> </span>
</button> </button>
<Button
variant="outline"
size="sm"
onClick={() => reanalyzeBatchMut.mutate(guildId)}
className="h-8 text-xs"
>
<RefreshCw className="mr-1 size-3" /> Reanalyze
</Button>
</div> </div>
{/* ── Sub navigation ── */} {/* ── Sub navigation ── */}
@@ -197,7 +184,6 @@ export default function MessagesPage() {
messages={currentMessages} messages={currentMessages}
selectedId={detailId} selectedId={detailId}
onSelect={setDetailId} onSelect={setDetailId}
onReanalyze={(id) => reanalyzeMut.mutate(id)}
hasMore={cursorData?.hasMore} hasMore={cursorData?.hasMore}
onLoadMore={handleLoadMore} onLoadMore={handleLoadMore}
isLoadingMore={loadMoreMut.isPending} isLoadingMore={loadMoreMut.isPending}
@@ -1,6 +1,6 @@
"use client"; "use client";
import { Loader2, RefreshCw, Search, Sparkles } from "lucide-react"; import { Loader2, Search, Sparkles } from "lucide-react";
import { useCallback, useState } from "react"; import { useCallback, useState } from "react";
import { EmptyState, LoadingSkeleton } from "@/components/shared"; import { EmptyState, LoadingSkeleton } from "@/components/shared";
@@ -10,14 +10,13 @@ import { Button } from "@/components/ui/button";
import { Card, CardContent } from "@/components/ui/card"; import { Card, CardContent } from "@/components/ui/card";
import { Input } from "@/components/ui/input"; import { Input } from "@/components/ui/input";
import { Progress } from "@/components/ui/progress"; import { Progress } from "@/components/ui/progress";
import { useMessageSearch, useReanalyze } from "@/hooks"; import { useMessageSearch } from "@/hooks";
import { renderMessageContent, safeParseJsonArray } from "@/lib/format"; import { renderMessageContent, safeParseJsonArray } from "@/lib/format";
import { cn } from "@/lib/utils"; import { cn } from "@/lib/utils";
export function SearchPanel() { export function SearchPanel() {
const [query, setQuery] = useState(""); const [query, setQuery] = useState("");
const [enabled, setEnabled] = useState(false); const [enabled, setEnabled] = useState(false);
const reanalyzeMut = useReanalyze();
const { data: results, isValidating: isFetching } = useMessageSearch( const { data: results, isValidating: isFetching } = useMessageSearch(
query, query,
@@ -134,14 +133,6 @@ export function SearchPanel() {
</span> </span>
</div> </div>
)} )}
<Button
variant="ghost"
size="xs"
onClick={() => reanalyzeMut.mutate(msg.id)}
>
<RefreshCw className="size-3 mr-1" />
Reanalyze
</Button>
</div> </div>
</div> </div>
</CardContent> </CardContent>
@@ -1,9 +1,8 @@
"use client"; "use client";
import { Hash, RefreshCw } from "lucide-react"; import { Hash } from "lucide-react";
import { Avatar, AvatarFallback, AvatarImage } from "@/components/ui/avatar"; import { Avatar, AvatarFallback, AvatarImage } from "@/components/ui/avatar";
import { Badge } from "@/components/ui/badge"; import { Badge } from "@/components/ui/badge";
import { Button } from "@/components/ui/button";
import { Card, CardContent } from "@/components/ui/card"; import { Card, CardContent } from "@/components/ui/card";
import { Progress } from "@/components/ui/progress"; import { Progress } from "@/components/ui/progress";
import { import {
@@ -18,11 +17,9 @@ import { AiStatusBadge } from "./ai-status-badge";
export function MessageCard({ export function MessageCard({
message: msg, message: msg,
onClick, onClick,
onReanalyze,
}: { }: {
message: MessageRecord; message: MessageRecord;
onClick: (id: string) => void; onClick: (id: string) => void;
onReanalyze: (id: string) => void;
}) { }) {
const severity = ( const severity = (
{ {
@@ -156,16 +153,6 @@ export function MessageCard({
</span> </span>
</div> </div>
)} )}
<Button
variant="ghost"
size="xs"
onClick={(e) => {
e.stopPropagation();
onReanalyze(msg.id);
}}
>
<RefreshCw className="size-3 mr-1" /> Reanalyze
</Button>
</div> </div>
</div> </div>
</CardContent> </CardContent>
@@ -9,7 +9,6 @@ interface MessageListProps {
messages: MessageRecord[]; messages: MessageRecord[];
selectedId: string | null; selectedId: string | null;
onSelect: (id: string) => void; onSelect: (id: string) => void;
onReanalyze?: (id: string) => void;
hasMore?: boolean; hasMore?: boolean;
onLoadMore?: () => void; onLoadMore?: () => void;
isLoadingMore?: boolean; isLoadingMore?: boolean;
@@ -19,7 +18,6 @@ export function MessageList({
messages, messages,
selectedId: _selectedId, selectedId: _selectedId,
onSelect, onSelect,
onReanalyze,
hasMore, hasMore,
onLoadMore, onLoadMore,
isLoadingMore, isLoadingMore,
@@ -27,12 +25,7 @@ export function MessageList({
return ( return (
<> <>
{messages.map((msg) => ( {messages.map((msg) => (
<MessageCard <MessageCard key={msg.id} message={msg} onClick={onSelect} />
key={msg.id}
message={msg}
onClick={onSelect}
onReanalyze={(id) => onReanalyze?.(id)}
/>
))} ))}
{hasMore && ( {hasMore && (
<div className="flex justify-center py-4"> <div className="flex justify-center py-4">
-2
View File
@@ -24,8 +24,6 @@ export {
useMessages, useMessages,
useMessagesHasMore, useMessagesHasMore,
useMessagesWsSync, useMessagesWsSync,
useReanalyze,
useReanalyzeBatch,
useReview, useReview,
useTextChannels, useTextChannels,
} from "./use-messages"; } from "./use-messages";
@@ -155,16 +155,6 @@ export function useMessageDetail(id: string | null) {
}; };
} }
// ── Mutations ────────────────────────────────────
export function useReanalyze() {
return useAction((id: string) => messagesApi.reanalyze(id));
}
export function useReanalyzeBatch() {
return useAction((guildId: string) => messagesApi.reanalyzeBatch(guildId));
}
// ── Search ─────────────────────────────────────── // ── Search ───────────────────────────────────────
export function useMessageSearch(query: string, enabled: boolean) { export function useMessageSearch(query: string, enabled: boolean) {
@@ -65,15 +65,6 @@ export const messagesApi = {
); );
}, },
reanalyze: (id: string) =>
api.post<{ ok: boolean }>(`/api/messages/${id}/reanalyze`, {}),
reanalyzeBatch: (guildId?: string, channelId?: string) =>
api.post<{ ok: boolean; count: number }>("/api/messages/reanalyze-batch", {
guildId,
channelId,
}),
search: (query: string, limit?: number) => { search: (query: string, limit?: number) => {
const params = new URLSearchParams({ q: query }); const params = new URLSearchParams({ q: query });
if (limit) params.set("limit", String(limit)); if (limit) params.set("limit", String(limit));