2026-05-14 02:31:16 +07:00
import { config } from "../config" ;
import { createChildLogger } from "../logger" ;
import type { SqliteDatabase } from "../muxer-queue" ;
import { retryWithBackoff } from "../retry" ;
2026-05-14 15:02:23 +07:00
import {
getMessageById ,
getPendingAIAnalysisMessages ,
updateMessageAIAnalysis ,
} from "./messageStore" ;
2026-05-14 02:31:16 +07:00
import type { MessageRecord } from "./types" ;
const logger = createChildLogger ( "ai-analyzer" );
const queuedMessageIds = new Set < string >();
let isProcessing = false ;
2026-05-14 03:54:12 +07:00
let activeRequests = 0 ;
const MAX_CONCURRENT_REQUESTS = 1 ;
2026-05-14 04:48:20 +07:00
const MAX_AI_REQUEST_TOKENS = 12 _000 ;
const AI_PROMPT_TOKEN_RESERVE = 3 _000 ;
const MAX_AI_BATCH_MESSAGES = 80 ;
2026-05-14 02:31:16 +07:00
interface ChatCompletionResponse {
choices? : Array < {
message ?: {
content? : string ;
};
} > ;
}
2026-05-14 02:44:26 +07:00
interface LLMAnalysis {
2026-05-14 04:23:11 +07:00
status : "clean" | "warn" | "flagged" ;
2026-05-14 02:44:26 +07:00
flags : string [];
score : number ;
analysis : string ;
}
2026-05-14 02:31:16 +07:00
function getAnalysisText ( message : MessageRecord ) : string {
return ( message . edited_content || message . content || "" ). trim ();
}
2026-05-14 04:42:28 +07:00
function estimateTokens ( text : string ) : number {
return Math . ceil ( text . length / 4 );
}
2026-05-14 15:02:23 +07:00
function formatMessageForAnalysis (
message : MessageRecord ,
index : number ,
) : string {
2026-05-14 04:42:28 +07:00
const text = getAnalysisText ( message );
const time = new Date ( message . created_at ). toISOString ();
return ` ${ index + 1 } . id= ${ message . id } time= ${ time } user= ${ message . username } : ${ text } ` ;
}
function estimateMessageTokens ( message : MessageRecord ) : number {
return estimateTokens ( formatMessageForAnalysis ( message , 0 )) + 16 ;
}
2026-05-14 02:31:16 +07:00
async function fetchJson ( url : string , init : RequestInit ) : Promise < unknown > {
const controller = new AbortController ();
2026-05-14 15:02:23 +07:00
const timeout = setTimeout (
() => controller . abort (),
config . AI_ANALYSIS_TIMEOUT_MS ,
);
2026-05-14 02:31:16 +07:00
try {
const response = await fetch ( url , { ... init , signal : controller.signal });
2026-05-14 04:00:31 +07:00
const text = await response . text ();
2026-05-14 02:31:16 +07:00
if ( ! response . ok ) {
2026-05-14 04:00:31 +07:00
const message = text . includes ( "{" )
? JSON . stringify ( JSON . parse ( text . substring ( text . indexOf ( "{" ))))
: text ;
2026-05-14 02:31:16 +07:00
throw new Error ( `AI request failed ( ${ response . status } ): ${ message } ` );
}
2026-05-14 04:00:31 +07:00
// Handle streaming response: extract JSON from response text
const jsonStart = text . indexOf ( "{" );
const jsonEnd = text . lastIndexOf ( "}" );
if ( jsonStart >= 0 && jsonEnd > jsonStart ) {
try {
return JSON . parse ( text . substring ( jsonStart , jsonEnd + 1 ));
} catch {
// Fall through to parse full text
}
}
return JSON . parse ( text );
2026-05-14 02:31:16 +07:00
} finally {
clearTimeout ( timeout );
}
}
2026-05-14 02:44:26 +07:00
function parseLLMAnalysis ( content : string ) : LLMAnalysis {
const jsonStart = content . indexOf ( "{" );
const jsonEnd = content . lastIndexOf ( "}" );
if ( jsonStart >= 0 && jsonEnd > jsonStart ) {
try {
const parsed = JSON . parse ( content . slice ( jsonStart , jsonEnd + 1 ));
2026-05-14 15:02:23 +07:00
const status =
parsed . status === "flagged"
? "flagged"
: parsed . status === "warn"
? "warn"
: "clean" ;
2026-05-14 02:44:26 +07:00
const flags = Array . isArray ( parsed . flags ) ? parsed . flags . map ( String ) : [];
const score = Math . max ( 0 , Math . min ( 1 , Number ( parsed . score ) || 0 ));
2026-05-14 15:02:23 +07:00
const analysis =
typeof parsed . analysis === "string" ? parsed.analysis : content ;
2026-05-14 02:44:26 +07:00
return { status , flags , score , analysis };
} catch {
// Fall through to text-only parsing.
}
}
2026-05-14 02:31:16 +07:00
return {
2026-05-14 15:02:23 +07:00
status :
/flagged|bahaya|berisiko|toxic|hate|harassment|violence|sexual|self-harm|illegal|scam|hacking/i . test (
content ,
)
? "flagged"
: /warn|provokasi|hinaan|menyerang/i . test ( content )
? "warn"
: "clean" ,
2026-05-14 02:44:26 +07:00
flags : [],
score : 0 ,
analysis : content.trim () || "Tidak ada analisis dari LLM." ,
2026-05-14 02:31:16 +07:00
};
}
2026-05-14 15:02:23 +07:00
async function runLLMAnalysis (
messages : MessageRecord [],
) : Promise < { results : LLMAnalysis []; raw : unknown } > {
const response = ( await retryWithBackoff (
() =>
fetchJson ( ` ${ config . AI_LLM_BASE_URL } /chat/completions` , {
method : "POST" ,
headers : {
Authorization : `Bearer ${ config . AI_LLM_API_KEY } ` ,
"Content-Type" : "application/json" ,
},
body : JSON.stringify ({
model : config.AI_LLM_MODEL ,
messages : [
{
role : "system" ,
content : `Kamu moderator Discord komunitas. Analisis setiap pesan dengan 3 kategori:
2026-05-14 04:23:11 +07:00
- CLEAN: Pesan normal, tidak melanggar aturan
2026-05-14 04:36:59 +07:00
- WARN: Melanggar aturan minor yang menarget orang lain (tone menyerang, hinaan ringan, konflik kecil) - butuh peringatan tapi tidak dihapus
2026-05-14 04:24:19 +07:00
- FLAGGED: Melanggar aturan berat (NSFW, ilegal, hacking, scam, harassment, violence, SARA, gore, spam, promosi judi) - butuh review moderator untuk penghapusan
2026-05-14 04:23:11 +07:00
2026-05-14 04:24:19 +07:00
ATURAN KOMUNITAS LENGKAP:
1. JAGA SIKAP DAN HORMATI SESAMA
- Gunakan bahasa yang sopan dan menghormati semua anggota
- Tanpa memandang latar belakang, usia, gender, atau pandangan
- Dilarang keras: pelecehan, rasisme, seksisme, diskriminasi
2. HINDARI KONFLIK
- Dilarang memancing keributan atau drama
- Jika ada masalah personal, selesaikan secara pribadi
- Jangan melibatkan anggota lain di channel umum
2026-05-14 04:28:21 +07:00
3. KONTEN EKSPLISIT DILARANG
2026-05-14 04:24:19 +07:00
- Dilarang keras: NSFW, ilegal, pornografi, kekerasan (gore), SARA
- Tidak ada tempat untuk penyimpangan atau LGBT
- Tidak ada promosi aktivitas atau ideologi LGBT
2026-05-14 04:28:21 +07:00
4. JAGA PRIVASI
2026-05-14 04:24:19 +07:00
- Dilarang menyebarkan informasi pribadi milik anggota lain tanpa izin
2026-05-14 04:28:21 +07:00
5. PROFIL YANG SOPAN
2026-05-14 04:24:19 +07:00
- Username, foto profil, dan server tag harus pantas
- Jangan gunakan unsur ofensif atau vulgar
2026-05-14 04:28:21 +07:00
6. DILARANG SPAM DAN PENIPUAN
2026-05-14 04:24:19 +07:00
- Dilarang: hoaks, link berbahaya (phishing/scam), spam
- Dilarang: promosi, judi, link referral
2026-05-14 04:36:59 +07:00
7. DISKUSI BERKUALITAS
2026-05-14 04:24:19 +07:00
- Berikan jawaban yang relevan, akurat, dan tidak menyesatkan
- Di channel "Area Serius", pertahankan standar tinggi
2026-05-14 04:36:59 +07:00
KONTEKS KOMUNITAS:
- Ini grup bercanda/santai, jadi slang, candaan ringan, kata kasar ringan tanpa target, pesan pendek seperti "." atau "P", dan pertanyaan tidak jelas tetap CLEAN
- Jangan beri WARN hanya karena pesan singkat, informal, ambigu, low-quality, atau kurang konteks
2026-05-14 04:42:28 +07:00
- Pahami alur pembahasan antar pesan: pesan yang sendiri terlihat normal bisa WARN/FLAGGED jika dalam konteks percakapan sedang memancing konflik, menormalisasi pelanggaran, atau melanjutkan provokasi
- Jangan menghukum orang yang sedang menasehati, menjelaskan bahaya, mengutip, atau menolak tindakan buruk; nilai maksud dan konteksnya
2026-05-14 04:36:59 +07:00
- WARN hanya jika ada orang/kelompok yang diserang, dihina, diprovokasi, atau konflik mulai dipancing
2026-05-14 04:24:19 +07:00
PENENTUAN STATUS:
2026-05-14 04:36:59 +07:00
- WARN jika: hinaan ringan yang menarget orang/kelompok, provokasi konflik kecil, username/profil kurang pantas
2026-05-14 04:24:19 +07:00
- FLAGGED jika: profanity berat, harassment, threats, violence, illegal activity, hacking, scam, NSFW, SARA, gore, spam, judi, LGBT content
2026-05-14 04:23:11 +07:00
Balas JSON array dengan schema: [{"status":"clean|warn|flagged","flags":["..."],"score":0..1,"analysis":"ringkasan Bahasa Indonesia + alasan + aksi disarankan"}]
Satu JSON object per pesan dalam array.` ,
2026-05-14 15:02:23 +07:00
},
{
role : "user" ,
content : `Analisis ${ messages . length } pesan berikut sebagai satu alur percakapan. Tetap kembalikan satu hasil per pesan dengan urutan yang sama: \ n ${ messages . map ( formatMessageForAnalysis ). join ( "\n" ) } ` ,
},
],
temperature : 0.2 ,
}),
signal : AbortSignal.timeout ( config . AI_ANALYSIS_TIMEOUT_MS ),
2026-05-14 02:31:16 +07:00
}),
{ retries : 2 , logger },
2026-05-14 15:02:23 +07:00
)) as ChatCompletionResponse ;
2026-05-14 02:31:16 +07:00
2026-05-14 02:44:26 +07:00
const content = response . choices ? .[ 0 ] ? . message ? . content ? . trim () || "" ;
2026-05-14 04:08:41 +07:00
// Extract JSON array from response
const jsonStart = content . indexOf ( "[" );
const jsonEnd = content . lastIndexOf ( "]" );
let results : LLMAnalysis [] = [];
if ( jsonStart >= 0 && jsonEnd > jsonStart ) {
try {
const parsed = JSON . parse ( content . substring ( jsonStart , jsonEnd + 1 ));
if ( Array . isArray ( parsed )) {
2026-05-14 04:23:11 +07:00
results = parsed . map (( item : any ) => {
2026-05-14 15:02:23 +07:00
const status =
item . status === "flagged"
? "flagged"
: item . status === "warn"
? "warn"
: "clean" ;
2026-05-14 04:23:11 +07:00
return {
status ,
flags : Array.isArray ( item . flags ) ? item . flags . map ( String ) : [],
score : Math.max ( 0 , Math . min ( 1 , Number ( item . score ) || 0 )),
2026-05-14 15:02:23 +07:00
analysis :
typeof item . analysis === "string" ? item.analysis : content ,
2026-05-14 04:23:11 +07:00
};
});
2026-05-14 04:08:41 +07:00
}
} catch {
// Fall through to individual parsing
}
}
// If batch parsing failed, parse as individual responses
if ( results . length === 0 ) {
2026-05-14 04:42:28 +07:00
results = messages . map (() => parseLLMAnalysis ( content ));
2026-05-14 04:08:41 +07:00
}
return { results , raw : response };
2026-05-14 02:31:16 +07:00
}
2026-05-14 15:02:23 +07:00
async function analyzeAndStoreBatch (
db : SqliteDatabase ,
messages : MessageRecord [],
) : Promise < void > {
2026-05-14 04:08:41 +07:00
if ( messages . length === 0 ) return ;
2026-05-14 15:02:23 +07:00
const analyzableMessages = messages . filter (
( message ) => getAnalysisText ( message ). length > 0 ,
);
2026-05-14 04:42:28 +07:00
if ( analyzableMessages . length === 0 ) return ;
2026-05-14 02:31:16 +07:00
2026-05-14 03:54:12 +07:00
activeRequests ++ ;
2026-05-14 02:31:16 +07:00
try {
2026-05-14 04:42:28 +07:00
const { results , raw } = await runLLMAnalysis ( analyzableMessages );
2026-05-14 04:08:41 +07:00
2026-05-14 04:42:28 +07:00
for ( let i = 0 ; i < analyzableMessages . length ; i ++ ) {
const message = analyzableMessages [ i ];
2026-05-14 04:08:41 +07:00
const result = results [ i ] || parseLLMAnalysis ( "" );
const row = updateMessageAIAnalysis ( db , message . id , {
2026-05-14 15:02:23 +07:00
status : result.status as
| "pending"
| "clean"
| "warn"
| "flagged"
| "error" ,
2026-05-14 04:08:41 +07:00
flags : JSON.stringify ( result . flags ),
score : result.score ,
raw : JSON.stringify ( raw ),
analysis : result.analysis ,
analyzedAt : Date.now (),
error : null ,
});
if ( row ) ( globalThis as any ). broadcastMessageAnalyzed ? .( row );
}
2026-05-14 02:31:16 +07:00
} catch ( error ) {
2026-05-14 04:48:20 +07:00
if ( analyzableMessages . length > 1 ) {
const midpoint = Math . ceil ( analyzableMessages . length / 2 );
logger . warn (
2026-05-14 15:02:23 +07:00
{
count : analyzableMessages.length ,
nextBatchSizes : [ midpoint , analyzableMessages . length - midpoint ],
error ,
},
2026-05-14 04:48:20 +07:00
"AI batch failed, splitting into smaller batches" ,
);
await analyzeAndStoreBatch ( db , analyzableMessages . slice ( 0 , midpoint ));
await analyzeAndStoreBatch ( db , analyzableMessages . slice ( midpoint ));
return ;
}
2026-05-14 04:08:41 +07:00
const errorMsg = error instanceof Error ? error.message : String ( error );
2026-05-14 04:42:28 +07:00
for ( const message of analyzableMessages ) {
2026-05-14 04:08:41 +07:00
const row = updateMessageAIAnalysis ( db , message . id , {
status : "error" ,
flags : null ,
score : null ,
raw : null ,
analysis : null ,
analyzedAt : Date.now (),
error : errorMsg ,
});
if ( row ) ( globalThis as any ). broadcastMessageAnalyzed ? .( row );
}
logger . warn ({ count : messages.length , error }, "AI batch analysis failed" );
2026-05-14 03:54:12 +07:00
} finally {
activeRequests -- ;
2026-05-14 02:31:16 +07:00
}
}
async function drainQueue ( db : SqliteDatabase ) : Promise < void > {
if ( isProcessing ) return ;
isProcessing = true ;
try {
2026-05-14 04:42:28 +07:00
const batchTokenLimit = MAX_AI_REQUEST_TOKENS - AI_PROMPT_TOKEN_RESERVE ;
2026-05-14 04:08:41 +07:00
2026-05-14 02:31:16 +07:00
while ( queuedMessageIds . size > 0 ) {
2026-05-14 03:54:12 +07:00
while ( activeRequests >= MAX_CONCURRENT_REQUESTS ) {
await new Promise (( resolve ) => setTimeout ( resolve , 100 ));
}
2026-05-14 04:08:41 +07:00
const batch : MessageRecord [] = [];
2026-05-14 04:42:28 +07:00
let tokenEstimate = 0 ;
for ( const messageId of Array . from ( queuedMessageIds )) {
2026-05-14 04:08:41 +07:00
const message = getMessageById ( db , messageId );
2026-05-14 04:42:28 +07:00
queuedMessageIds . delete ( messageId );
if ( ! message ) continue ;
const messageTokens = estimateMessageTokens ( message );
2026-05-14 15:02:23 +07:00
if (
batch . length > 0 &&
( batch . length >= MAX_AI_BATCH_MESSAGES ||
tokenEstimate + messageTokens > batchTokenLimit )
) {
2026-05-14 04:42:28 +07:00
queuedMessageIds . add ( messageId );
break ;
}
batch . push ( message );
tokenEstimate += messageTokens ;
2026-05-14 04:08:41 +07:00
}
if ( batch . length > 0 ) {
2026-05-14 15:02:23 +07:00
logger . info (
{ count : batch.length , tokenEstimate },
"Processing AI analysis batch" ,
);
2026-05-14 04:08:41 +07:00
await analyzeAndStoreBatch ( db , batch );
}
2026-05-14 02:31:16 +07:00
}
} finally {
isProcessing = false ;
}
}
2026-05-14 15:02:23 +07:00
export function queueMessageAnalysis (
db : SqliteDatabase ,
messageId : string ,
) : void {
2026-05-14 02:31:16 +07:00
if ( ! config . AI_ANALYSIS_ENABLED ) return ;
2026-05-14 02:44:26 +07:00
logger . debug ({ messageId }, "Queueing AI analysis" );
2026-05-14 02:31:16 +07:00
queuedMessageIds . add ( messageId );
setImmediate (() => {
2026-05-14 15:02:23 +07:00
drainQueue ( db ). catch (( error ) =>
logger . error ({ error }, "AI analysis queue failed" ),
);
2026-05-14 02:31:16 +07:00
});
}
2026-05-14 02:44:26 +07:00
export function startPendingAIAnalysisWorker ( db : SqliteDatabase ) : void {
if ( ! config . AI_ANALYSIS_ENABLED ) {
logger . info ( "AI analysis disabled" );
return ;
}
logger . info ( "AI analysis worker started" );
setInterval (() => {
if ( isProcessing ) return ;
2026-05-14 04:42:28 +07:00
const pendingMessages = getPendingAIAnalysisMessages ( db , 500 );
2026-05-14 02:44:26 +07:00
if ( pendingMessages . length === 0 ) return ;
2026-05-14 15:02:23 +07:00
logger . info (
{ count : pendingMessages.length },
"Queueing pending AI analysis messages" ,
);
2026-05-14 02:44:26 +07:00
for ( const message of pendingMessages ) {
queuedMessageIds . add ( message . id );
}
2026-05-14 15:02:23 +07:00
drainQueue ( db ). catch (( error ) =>
logger . error ({ error }, "Pending AI analysis worker failed" ),
);
2026-05-14 02:44:26 +07:00
}, 15000 );
}