feat: multimodal video detection + reply/forward/crosspost + batch optimization

- Video frame extraction via ffmpeg (4 key frames per video → vision LLM)
- Video display in FE MessageCard with HTML5 <video> player
- Reply/forward/crosspost indicator in FE + pipeline in DG/BE
- Fix: missing sanitizeAiContent + escapeXml in media path (prompt injection)
- Optimize: text-only batch results saved to DB immediately, no longer wait for media analysis
- BE mapper/schema/repo: add reference fields (is_reply, is_forward, etc.)
This commit is contained in:
MythEclipse
2026-06-22 09:46:06 +07:00
parent dc62b283c0
commit f502918b27
7 changed files with 261 additions and 48 deletions
@@ -134,6 +134,12 @@ export class MessagesRepository {
type: data.type ?? "text",
metadata: null,
ai_status: "pending",
is_reply: data.isReply ?? false,
is_forward: data.isForward ?? false,
is_crosspost: data.isCrosspost ?? false,
reference_message_id: data.referenceMessageId ?? null,
reference_channel_id: data.referenceChannelId ?? null,
reference_guild_id: data.referenceGuildId ?? null,
})
.returning();
@@ -260,6 +266,12 @@ export class MessagesRepository {
ai_severity: pgMessagesTable.ai_severity,
ai_confidence: pgMessagesTable.ai_confidence,
ai_analysis: pgMessagesTable.ai_analysis,
is_reply: pgMessagesTable.is_reply,
is_forward: pgMessagesTable.is_forward,
is_crosspost: pgMessagesTable.is_crosspost,
reference_message_id: pgMessagesTable.reference_message_id,
reference_channel_id: pgMessagesTable.reference_channel_id,
reference_guild_id: pgMessagesTable.reference_guild_id,
})
.from(pgMessagesTable)
.where(and(...conditions))
@@ -19,6 +19,12 @@ export const messageCreateSchema = z.object({
avatarUrl: z.string().optional(),
content: z.string(),
type: z.enum(["text", "edited", "deleted"]).default("text"),
isReply: z.boolean().optional(),
isForward: z.boolean().optional(),
isCrosspost: z.boolean().optional(),
referenceMessageId: z.string().optional(),
referenceChannelId: z.string().optional(),
referenceGuildId: z.string().optional(),
});
export const messageUpdateSchema = z.object({
@@ -25,6 +25,12 @@ export interface MappedMessage {
ai_recommended_action: string | null;
ai_analyzed_at: number | null;
ai_error: string | null;
is_reply: boolean | null;
is_forward: boolean | null;
is_crosspost: boolean | null;
reference_message_id: string | null;
reference_channel_id: string | null;
reference_guild_id: string | null;
}
export function mapMessageRow(row: Record<string, unknown>): MappedMessage {
@@ -53,5 +59,11 @@ export function mapMessageRow(row: Record<string, unknown>): MappedMessage {
ai_recommended_action: (row.ai_recommended_action as string | null) ?? null,
ai_analyzed_at: (row.ai_analyzed_at as number | null) ?? null,
ai_error: (row.ai_error as string | null) ?? null,
is_reply: row.is_reply === null ? null : Boolean(row.is_reply),
is_forward: row.is_forward === null ? null : Boolean(row.is_forward),
is_crosspost: row.is_crosspost === null ? null : Boolean(row.is_crosspost),
reference_message_id: (row.reference_message_id as string | null) ?? null,
reference_channel_id: (row.reference_channel_id as string | null) ?? null,
reference_guild_id: (row.reference_guild_id as string | null) ?? null,
};
}
@@ -10,6 +10,7 @@ import type {
AnalysisResult,
MessageRecord,
} from "../message-capture/types.js";
import { extractMessageMediaEvidence } from "../message-capture/messageMetadata.js";
import { buildConversationContext } from "./conversationContext.js";
import {
runModerationAnalysis,
@@ -159,36 +160,92 @@ async function processBatch(job: {
const allMessageIds = [...targetIds, ...contextIds];
const attachments = await getAttachmentsForMessages(allMessageIds);
const result = await runModerationAnalysis({
targets: messages,
contextText: contextLines.join("\n"),
attachments,
});
// ── Split: text-only vs media ──────────────────────────────────────
// Text-only analysis runs fast (single LLM call, no vision).
// Media analysis is slow (download + vision → LLM).
// By splitting here, text results are saved to DB immediately
// instead of waiting for media downloads to finish.
// ────────────────────────────────────────────────────────────────────
const textOnly: MessageRecord[] = [];
const media: MessageRecord[] = [];
const updates = result.results.map((analysisResult) => ({
messageId: analysisResult.messageId,
result: {
status: analysisResult.status,
flags: JSON.stringify(analysisResult.flags),
score: analysisResult.score,
analysis: analysisResult.analysis,
categories: analysisResult.categories,
severity: analysisResult.severity,
confidence: analysisResult.confidence,
recommendedAction: analysisResult.recommendedAction,
analyzedAt: Date.now(),
error: null,
},
}));
try {
const rows = await updateMessagesAIAnalysisBulk(updates);
return { ok: true, conversationKey, rows };
} catch (dbErr) {
throw new Error(
`Failed to update DB: ${dbErr instanceof Error ? dbErr.message : String(dbErr)}`,
);
for (const msg of messages) {
const meta = msg.metadata ? extractMessageMediaEvidence(msg.metadata) : null;
if (meta && (meta.attachments.length > 0 || meta.stickers.length > 0 || meta.embeds.length > 0)) {
media.push(msg);
} else {
textOnly.push(msg);
}
}
const allRows: MessageRecord[] = [];
// Phase 1: Text-only → save immediately (fast)
if (textOnly.length > 0) {
const textResult = await runModerationAnalysis({
targets: textOnly,
contextText: contextLines.join("\n"),
attachments,
});
const textUpdates = textResult.results.map((analysisResult) => ({
messageId: analysisResult.messageId,
result: {
status: analysisResult.status,
flags: JSON.stringify(analysisResult.flags),
score: analysisResult.score,
analysis: analysisResult.analysis,
categories: analysisResult.categories,
severity: analysisResult.severity,
confidence: analysisResult.confidence,
recommendedAction: analysisResult.recommendedAction,
analyzedAt: Date.now(),
error: null,
},
}));
if (textUpdates.length > 0) {
const rows = await updateMessagesAIAnalysisBulk(textUpdates);
allRows.push(...rows);
logger.info(
{ count: textUpdates.length, conversationKey },
"Text-only batch saved — media analysis still in progress",
);
}
}
// Phase 2: Media → save when done (slow: download + vision)
if (media.length > 0) {
const mediaResult = await runModerationAnalysis({
targets: media,
contextText: contextLines.join("\n"),
attachments,
});
const mediaUpdates = mediaResult.results.map((analysisResult) => ({
messageId: analysisResult.messageId,
result: {
status: analysisResult.status,
flags: JSON.stringify(analysisResult.flags),
score: analysisResult.score,
analysis: analysisResult.analysis,
categories: analysisResult.categories,
severity: analysisResult.severity,
confidence: analysisResult.confidence,
recommendedAction: analysisResult.recommendedAction,
analyzedAt: Date.now(),
error: null,
},
}));
if (mediaUpdates.length > 0) {
const rows = await updateMessagesAIAnalysisBulk(mediaUpdates);
allRows.push(...rows);
}
}
logger.info(
{ total: messages.length, textOnly: textOnly.length, media: media.length, saved: allRows.length },
"Batch analysis complete",
);
return { ok: true, conversationKey, rows: allRows };
}
// ---------------------------------------------------------------------------
@@ -1,4 +1,9 @@
import { execFile } from "node:child_process";
import { createChildLogger } from "@bete/shared/logger";
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, retryWithBackoff } from "@bete/shared/utils";
import { LRUCache } from "lru-cache";
import type { ChatCompletion } from "openai/resources/chat/completions";
@@ -864,7 +869,7 @@ async function runTextOnlyBatch(
.map((url) => {
const fetchedText = urlFetchMap.get(url);
if (!fetchedText) return null;
return `<web_content url="${url}">${fetchedText}</web_content>`;
return `<web_content url="${escapeXml(url)}">${escapeXml(fetchedText)}</web_content>`;
})
.filter(Boolean)
.join("\n");
@@ -876,7 +881,7 @@ async function runTextOnlyBatch(
const profileLine = userProfileCtx ? `\n ${userProfileCtx}` : "";
const refXml = await buildReferenceXml(msg);
const refLine = refXml ? `\n ${refXml}` : "";
return `<message id="${msg.id}" user="${msg.username}">\n ${userCtx}${profileLine}${refLine}\n <content>${content}</content>${webContext}\n</message>`;
return `<message id="${msg.id}" user="${msg.username}">\n ${userCtx}${profileLine}${refLine}\n <content>${escapeXml(content)}</content>${webContext}\n</message>`;
}),
).then((blocks) => blocks.join("\n"));
@@ -1008,7 +1013,7 @@ async function prepareMediaMessage(
(att) =>
att.message_id === targetId &&
(att.uploaded_url ?? att.discord_url ?? null) &&
att.type.startsWith("image/"),
(att.type.startsWith("image/") || att.type.startsWith("video/")),
)
.slice(0, 8);
@@ -1092,12 +1097,12 @@ async function prepareMediaMessage(
const userCtx = `<user_reputation trust_score="${rep.trust_score}" />`;
const profile = await getUserProfile(target.user_id);
const userProfileCtx = profile
? `\n <user_profile>${profile.profile_summary}</user_profile>`
? `\n <user_profile>${sanitizeAiContent(profile.profile_summary)}</user_profile>`
: "";
const refXml = await buildReferenceXml(target);
const refLine = refXml ? `\n ${refXml}` : "";
const messageBlock = `<message id="${target.id}" user="${target.username}">\n ${userCtx}${userProfileCtx}${refLine}\n <content>${content}</content>${mediaContext ? ` ${mediaContext}` : ""}${webContext}${mediaAnalysisContext}\n</message>`;
const messageBlock = `<message id="${escapeXml(target.id)}" user="${escapeXml(target.username)}">\n ${userCtx}${userProfileCtx}${refLine}\n <content>${escapeXml(content)}</content>${mediaContext ? ` ${escapeXml(mediaContext)}` : ""}${webContext}${mediaAnalysisContext}\n</message>`;
return { targetId, messageBlock };
}
@@ -1745,6 +1750,68 @@ async function downloadSingleAttachment(
const imageBytes = Buffer.concat(chunks);
const sniffedMime = sniffImageMimeType(imageBytes);
if (!sniffedMime && att.type.startsWith("video/")) {
// ── Video frame extraction via ffmpeg ──
const execFileAsync = promisify(execFile);
const tmpDir = await mkdtemp(path.join(tmpdir(), "bete-video-"));
const inputPath = path.join(tmpDir, att.filename || "video.mp4");
const outputPattern = path.join(tmpDir, "frame-%03d.jpg");
try {
await writeFile(inputPath, imageBytes);
// Get video duration via ffprobe, then extract 4 evenly-spaced frames
const { stdout: durationStr } = await execFileAsync("/usr/bin/ffprobe", [
"-v", "error",
"-show_entries", "format=duration",
"-of", "csv=p=0",
inputPath,
], { timeout: 10000 });
const duration = parseFloat(durationStr.trim()) || 1;
// fps = 3/duration gives exactly 4 frames at 0, dur/3, 2*dur/3, dur
const fps = (3 / duration).toFixed(6);
await execFileAsync("/usr/bin/ffmpeg", [
"-i", inputPath,
"-vf", `fps=${fps}`,
"-frames:v", "4",
"-vsync", "vfr",
"-q:v", "2",
outputPattern,
], { timeout: 30000 });
// Read extracted frames and add to imageMap
for (let i = 1; i <= 4; i++) {
const framePath = path.join(tmpDir, `frame-${String(i).padStart(3, "0")}.jpg`);
try {
const frameBytes = await readFile(framePath);
const { data: resizedBuffer, mimeType: resizedMime } =
await resizeImageForVision(frameBytes, maxDimension);
const dataUrl = `data:${resizedMime};base64,${resizedBuffer.toString("base64")}`;
const part: MessageImagePart = {
type: "image_url",
image_url: { url: dataUrl },
sourceLabel: `[frame ${i}/4 dari video ${att.filename} (attachment), pesan id=${att.message_id}]`,
};
addImageToMap(imageMap, targetId, part);
} catch {
// Frame may not exist if video is short; skip silently
}
}
log.info({ attachmentId: att.id, frameCount: 4 }, "Extracted video frames for vision analysis");
} catch (ffmpegErr) {
log.warn(
{ attachmentId: att.id, error: ffmpegErr instanceof Error ? ffmpegErr.message : String(ffmpegErr) },
"Failed to extract video frames with ffmpeg — skipping video",
);
} finally {
// Cleanup temp files
try { await unlink(inputPath); } catch { /* ignore */ }
for (let i = 1; i <= 4; i++) {
try { await unlink(path.join(tmpDir, `frame-${String(i).padStart(3, "0")}.jpg`)); } catch { /* ignore */ }
}
try { await rm(tmpDir, { recursive: true, force: true }); } catch { /* ignore */ }
}
return;
}
if (!sniffedMime) {
log.warn(
{ attachmentId: att.id },
@@ -1752,7 +1819,6 @@ async function downloadSingleAttachment(
);
return;
}
const { data: resizedBuffer, mimeType: resizedMime } =
await resizeImageForVision(imageBytes, maxDimension);
File diff suppressed because one or more lines are too long
@@ -10,6 +10,7 @@ import {
RotateCw,
Smile,
Trash2,
Video,
} from "lucide-react";
import { Fragment, useEffect, useMemo, useState } from "react";
import type { MessageRecord } from "../../../entities/message/types.js";
@@ -183,7 +184,13 @@ function MessageRow({
a.contentType?.startsWith("image/") ||
/\.(png|jpe?g|gif|webp)$/i.test(a.name),
);
const videoAttachments = attachments.filter(
(a) =>
a.contentType?.startsWith("video/") ||
/\.(mp4|webm|mov|mkv|avi)$/i.test(a.name),
);
const hasImages = imageAttachments.length > 0;
const hasVideos = videoAttachments.length > 0;
/** Hide the fallback text ("[Attachment: ...]", "[Sticker: ...]", "[Embed]") when the actual media IS already shown visually. */
const isFallbackText =
@@ -365,6 +372,27 @@ function MessageRow({
</div>
)}
{/* Attached videos */}
{hasVideos && (
<div className="flex gap-2 overflow-x-auto">
{videoAttachments.slice(0, 4).map((vid) => (
<video
key={vid.url}
src={vid.url}
controls
className="h-28 w-48 shrink-0 rounded-lg border border-border object-cover bg-black"
preload="metadata"
/>
))}
{videoAttachments.length > 4 && (
<div className="flex h-28 w-16 items-center justify-center rounded-lg border border-border bg-muted text-[11px] text-muted-foreground">
+{videoAttachments.length - 4}
<Video className="ml-0.5 h-3 w-3" />
</div>
)}
</div>
)}
{/* Categories */}
{categories.length > 0 && (
<div className="flex flex-wrap gap-1">