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
This commit is contained in:
asepharyana
2026-09-04 15:32:03 +07:00
committed by asepharyana
parent 9682abef82
commit 276062fdd4
5 changed files with 112 additions and 40 deletions
@@ -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<string[]> {
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,
@@ -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,
@@ -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",
});
}
}
@@ -1,5 +1,6 @@
export {
incrementCounter,
incrementCounterBy,
registerCollector,
setGauge,
startMetricsServer,
@@ -37,16 +37,25 @@ export function registerCollector(fn: () => void): void {
export function incrementCounter(
name: string,
labels?: Record<string, string>,
): 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<string, string>,
): 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,
});
}