From 276062fdd442ddbd750c409fa936f68ee45bf7bb Mon Sep 17 00:00:00 2001 From: asepharyana Date: Fri, 4 Sep 2026 15:32:03 +0700 Subject: [PATCH] feat(ai): scope persistence, parallel tools, Prometheus token/cache counters C2 full scope persistence: - getRecentConversationContext now returns per-turn guildId/channelId - processMessage merges historical scope when current request is unscoped - chatbot remembers server context across all 8 history exchanges (not just 3) A4+Metrics: - Add incrementCounterBy(name, delta, labels) to gateway-metrics - Token-usage counters: llm_tokens_total{model, type, label} per batch - Cache hit counters: moderation_cache_hits{type} exact/semantic-qdrant/semantic-pg - Cache miss counters: moderation_cache_misses per batch - Prometheus /metrics now exposes cost + cache hit-rate for dashboards Performance: - Parallel tool execution within each chatbot round (Promise.all) - All tool results collected before sending to model (ordering preserved) - Tool failure now logged with structured warning (chatbot.tools.ts) Verified: backend tsc+biome 37/37, discord-gateway tsc+biome 210/210 --- .../src/modules/chatbot/chatbot.service.ts | 111 ++++++++++++------ .../src/modules/ai-moderation/llmCaller.ts | 18 +++ .../ai-moderation/moderationOrchestrator.ts | 9 ++ .../src/modules/gateway-metrics/index.ts | 1 + .../src/modules/gateway-metrics/metrics.ts | 13 +- 5 files changed, 112 insertions(+), 40 deletions(-) diff --git a/services/backend/src/modules/chatbot/chatbot.service.ts b/services/backend/src/modules/chatbot/chatbot.service.ts index 31c28a90..d7745e93 100644 --- a/services/backend/src/modules/chatbot/chatbot.service.ts +++ b/services/backend/src/modules/chatbot/chatbot.service.ts @@ -26,13 +26,25 @@ class ChatbotService { // longer bake server stats into the prompt — the model must pull current // data via tools (see buildSystemPrompt), so it always answers from live // numbers instead of a stale snapshot. - const scope = { + // + // If the current request doesn't carry a guild/channel scope (e.g. a + // cross-server dashboard surface), fall back to the most recent scope the + // user was chatting in from history, so the agent doesn't drift to the + // wrong server between turns. + const explicitScope = { guildId: context?.guildId, channelId: context?.channelId, }; + const scope = + explicitScope.guildId || explicitScope.channelId + ? explicitScope + : (recentContext.lastScope ?? { + guildId: undefined, + channelId: undefined, + }); const systemPrompt = this.buildSystemPrompt(scope); - const conversationHistory = this.buildHistoryMessages(recentContext); + const conversationHistory = this.buildHistoryMessages(recentContext.turns); const llmResponse = await this.callLLM( systemPrompt, conversationHistory, @@ -61,14 +73,28 @@ class ChatbotService { await chatbotRepository.clearChatHistory(userId); } - private async getRecentConversationContext( - userId: string, - ): Promise { + private async getRecentConversationContext(userId: string): Promise<{ + turns: string[]; + lastScope: { guildId?: string; channelId?: string } | null; + }> { const history = await chatbotRepository.getChatHistory(userId, 8); - return history.flatMap((row) => [ - `User: ${row.user_message}`, - `Bot: ${row.bot_response}`, - ]); + const turns: string[] = []; + let lastScope: { guildId?: string; channelId?: string } | null = null; + for (const row of history) { + turns.push(`User: ${row.user_message}`, `Bot: ${row.bot_response}`); + // Capture the most recent scope this user was chatting in, so the + // agent keeps server context across turns even if the current request + // doesn't carry one. + const g = row.context?.guildId; + const c = row.context?.channelId; + if (g || c) { + lastScope = { + guildId: g ?? undefined, + channelId: c ?? undefined, + }; + } + } + return { turns, lastScope }; } private buildSystemPrompt(scope: { @@ -213,35 +239,44 @@ Gaya ngobrol: ); if (toolCalls.length > 0) { - // Execute each tool, append tool results, continue loop. - for (const tc of toolCalls) { - messages.push({ - role: "assistant", - content: null, - tool_calls: [ - { - id: tc.id, - type: "function", - function: { name: tc.name, arguments: tc.arguments }, - }, - ], - }); - // Auto-scope: if the model omitted guildId/channelId, fill them - // from the request scope so tools query the right server without - // the model having to guess IDs. - const scopedArgs = { ...tc.args }; - if (scope.guildId && scopedArgs.guildId == null) { - scopedArgs.guildId = scope.guildId; - } - if (scope.channelId && scopedArgs.channelId == null) { - scopedArgs.channelId = scope.channelId; - } - let result = ""; - try { - result = await executeTool(tc.name, scopedArgs); - } catch (e) { - result = `Tool error: ${(e as Error).message}`; - } + // Executor for a single tool call: pushes the assistant tool_call, + // runs the tool, returns { tc, result } so results can be appended + // in order after all tools execute in parallel. + const executions = await Promise.all( + toolCalls.map(async (tc) => { + messages.push({ + role: "assistant", + content: null, + tool_calls: [ + { + id: tc.id, + type: "function", + function: { name: tc.name, arguments: tc.arguments }, + }, + ], + }); + // Auto-scope: if the model omitted guildId/channelId, fill them + // from the request scope so tools query the right server without + // the model having to guess IDs. + const scopedArgs = { ...tc.args }; + if (scope.guildId && scopedArgs.guildId == null) { + scopedArgs.guildId = scope.guildId; + } + if (scope.channelId && scopedArgs.channelId == null) { + scopedArgs.channelId = scope.channelId; + } + let result = ""; + try { + result = await executeTool(tc.name, scopedArgs); + } catch (e) { + result = `Tool error: ${(e as Error).message}`; + } + return { tc, result }; + }), + ); + // Append tool results in the SAME order as the tool_calls so the + // API's function-calling contract isn't violated by reordering. + for (const { tc, result } of executions) { messages.push({ role: "tool", tool_call_id: tc.id, diff --git a/services/discord-gateway/src/modules/ai-moderation/llmCaller.ts b/services/discord-gateway/src/modules/ai-moderation/llmCaller.ts index d63b4f01..33717f35 100644 --- a/services/discord-gateway/src/modules/ai-moderation/llmCaller.ts +++ b/services/discord-gateway/src/modules/ai-moderation/llmCaller.ts @@ -14,6 +14,7 @@ import type { ChatCompletion } from "openai/resources/chat/completions"; import { createChildLogger } from "@/shared/logger/index"; import { delay, retryWithBackoff } from "@/shared/utils/index"; import { config } from "../../shared/config/config.js"; +import { incrementCounterBy } from "../gateway-metrics/index.js"; import type { AnalysisResult } from "../message-capture/types.js"; import { llmChat } from "./llmClient.js"; import { parseModerationResponse } from "./moderationResponseParser.js"; @@ -170,6 +171,23 @@ export async function callModerationLLM( // so cost per channel/guild can be tracked (routers bill per token). const usage = result?.usage; if (usage && (usage.prompt_tokens || usage.completion_tokens)) { + // Prometheus token-usage counters for cost tracking per model + phase. + // Counters accumulate across scrapes so total spend per model/label is + // queryable (rate(...) dashboards). + if (usage.prompt_tokens) { + incrementCounterBy("llm_tokens_total", usage.prompt_tokens, { + model: config.AI_LLM_MODEL, + type: "prompt", + label, + }); + } + if (usage.completion_tokens) { + incrementCounterBy("llm_tokens_total", usage.completion_tokens, { + model: config.AI_LLM_MODEL, + type: "completion", + label, + }); + } log.info( { label, diff --git a/services/discord-gateway/src/modules/ai-moderation/moderationOrchestrator.ts b/services/discord-gateway/src/modules/ai-moderation/moderationOrchestrator.ts index 2d0d104c..46bdf4ab 100644 --- a/services/discord-gateway/src/modules/ai-moderation/moderationOrchestrator.ts +++ b/services/discord-gateway/src/modules/ai-moderation/moderationOrchestrator.ts @@ -8,6 +8,7 @@ import { LRUCache } from "lru-cache"; import { createChildLogger } from "@/shared/logger/index"; import { config } from "../../shared/config/config.js"; +import { incrementCounterBy } from "../gateway-metrics/index.js"; import { extractMessageMediaEvidence } from "../message-capture/messageMetadata.js"; import type { AnalysisResult, @@ -243,10 +244,12 @@ export async function runModerationAnalysis( const representative = hitByKey.get(candidate.scopedKey); if (representative) { cacheHits.push({ ...representative, messageId: candidate.target.id }); + incrementCounterBy("moderation_cache_hits", 1, { type: "exact" }); } else { uncachedTargets.push(candidate.target); } } + incrementCounterBy("moderation_cache_misses", uncachedTargets.length); // ── Phase 2: semantic cache — batched (one embed call + one Qdrant // batch search for ALL uncached text targets) ───────────────────────── @@ -321,6 +324,9 @@ export async function runModerationAnalysis( hitByKey.set(cacheKey, hit); servedCacheKeys.add(cacheKey); // bump hit_count for metrics logCacheEvent("hit", cacheKey, "text"); + incrementCounterBy("moderation_cache_hits", 1, { + type: "semantic-qdrant", + }); } } else { // Legacy Postgres fallback path (no Qdrant): per-candidate scan. @@ -360,6 +366,9 @@ export async function runModerationAnalysis( hitByKey.set(cacheKey, hit); servedCacheKeys.add(cacheKey); // bump hit_count for metrics logCacheEvent("hit", cacheKey, "text"); + incrementCounterBy("moderation_cache_hits", 1, { + type: "semantic-pg", + }); } } diff --git a/services/discord-gateway/src/modules/gateway-metrics/index.ts b/services/discord-gateway/src/modules/gateway-metrics/index.ts index cfd838aa..00b585a1 100644 --- a/services/discord-gateway/src/modules/gateway-metrics/index.ts +++ b/services/discord-gateway/src/modules/gateway-metrics/index.ts @@ -1,5 +1,6 @@ export { incrementCounter, + incrementCounterBy, registerCollector, setGauge, startMetricsServer, diff --git a/services/discord-gateway/src/modules/gateway-metrics/metrics.ts b/services/discord-gateway/src/modules/gateway-metrics/metrics.ts index 61a64bad..518fab65 100644 --- a/services/discord-gateway/src/modules/gateway-metrics/metrics.ts +++ b/services/discord-gateway/src/modules/gateway-metrics/metrics.ts @@ -37,16 +37,25 @@ export function registerCollector(fn: () => void): void { export function incrementCounter( name: string, labels?: Record, +): void { + incrementCounterBy(name, 1, labels); +} + +/** Increment a counter by an explicit delta (e.g. token counts per batch). */ +export function incrementCounterBy( + name: string, + delta: number, + labels?: Record, ): void { const k = key(`bete_${name}`, labels); const existing = metrics.get(k); if (existing) { - existing.value += 1; + existing.value += delta; } else { metrics.set(k, { help: `Counter: ${name}`, type: "counter", - value: 1, + value: delta, labels: labels ? { ...labels } : undefined, }); }