perf(ai-moderation): speed up analysis queue (ramai + sepi)
- Parallelize per-user reputation/profile fetches in textBatchProcessor (was a serial ~2N DB/Redis round-trip loop per sub-batch; now Promise.all over unique users). Cuts per-batch latency, biggest win on small/quiet batches. - Make the LLM concurrency semaphore dynamic (cached per config value) instead of frozen at import time, so AI_LLM_MAX_CONCURRENT is tunable without code change and reflects current config. - Bump AI_LLM_MAX_CONCURRENT default 5 -> 8 (gemini-flash-lite is cheap; helps throughput when busy). - Lower AI_ANALYSIS_DEBOUNCE_MS 500 -> 250 (snappier first-message analysis when quiet). - Lower AI_ANALYSIS_RECOVERY_INTERVAL_MS 15000 -> 10000 (stuck/errored messages re-analyze sooner). tsc, biome, vitest (129) all clean.
This commit is contained in:
@@ -18,7 +18,20 @@ const log = createChildLogger("llm-client");
|
|||||||
// Concurrency limiter for LLM API calls (inlined from concurrencyLimiter.ts)
|
// Concurrency limiter for LLM API calls (inlined from concurrencyLimiter.ts)
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
const llmSemaphore = pLimit(config.AI_LLM_MAX_CONCURRENT ?? 5);
|
// The limiter is cached per configured concurrency value so it can be tuned
|
||||||
|
// (env / BWS) without a code change and always reflects the current config —
|
||||||
|
// a module-level `pLimit(config.X)` would freeze the cap at import time.
|
||||||
|
let llmSemaphore = pLimit(config.AI_LLM_MAX_CONCURRENT ?? 5);
|
||||||
|
let llmSemaphoreLimit = config.AI_LLM_MAX_CONCURRENT ?? 5;
|
||||||
|
|
||||||
|
function getLlmSemaphore() {
|
||||||
|
const wanted = config.AI_LLM_MAX_CONCURRENT ?? 5;
|
||||||
|
if (wanted !== llmSemaphoreLimit) {
|
||||||
|
llmSemaphore = pLimit(wanted);
|
||||||
|
llmSemaphoreLimit = wanted;
|
||||||
|
}
|
||||||
|
return llmSemaphore;
|
||||||
|
}
|
||||||
|
|
||||||
let activeCount = 0;
|
let activeCount = 0;
|
||||||
let pendingCount = 0;
|
let pendingCount = 0;
|
||||||
@@ -30,7 +43,7 @@ export async function withLlmConcurrency<T>(fn: () => Promise<T>): Promise<T> {
|
|||||||
"Queuing LLM request",
|
"Queuing LLM request",
|
||||||
);
|
);
|
||||||
|
|
||||||
return llmSemaphore(async () => {
|
return getLlmSemaphore()(async () => {
|
||||||
pendingCount--;
|
pendingCount--;
|
||||||
activeCount++;
|
activeCount++;
|
||||||
|
|
||||||
|
|||||||
@@ -214,20 +214,28 @@ export async function runTextOnlyBatch(
|
|||||||
asOf?: number | null;
|
asOf?: number | null;
|
||||||
}
|
}
|
||||||
>();
|
>();
|
||||||
for (const msg of batch) {
|
// ── Per-user reputation + profile context (fetched ONCE per unique user,
|
||||||
if (!userContexts.has(msg.user_id)) {
|
// in parallel — was a serial per-message loop that cost ~2N sequential
|
||||||
const rep = await initializeUserReputation(msg.user_id, msg.guild_id);
|
// DB/Redis round-trips per sub-batch and dominated latency on small
|
||||||
const repAttrs = formatReputationAttrs(rep);
|
// batches). ─────────────────────────────────────────────────────────
|
||||||
const repXml = `<user_reputation ${repAttrs}/>`;
|
const uniqueUserIds = [...new Set(batch.map((m) => m.user_id))];
|
||||||
userContexts.set(msg.user_id, repXml);
|
const batchGuildId = batch[0]?.guild_id ?? "";
|
||||||
}
|
const userFetches = await Promise.all(
|
||||||
if (!userProfiles.has(msg.user_id)) {
|
uniqueUserIds.map(async (uid) => {
|
||||||
const profile = await getUserProfile(msg.user_id);
|
const [rep, profile] = await Promise.all([
|
||||||
userProfiles.set(msg.user_id, {
|
initializeUserReputation(uid, batchGuildId),
|
||||||
text: profile?.profile_summary ?? "",
|
getUserProfile(uid),
|
||||||
asOf: profile?.last_analyzed_at ?? null,
|
]);
|
||||||
});
|
return { uid, rep, profile };
|
||||||
}
|
}),
|
||||||
|
);
|
||||||
|
for (const { uid, rep, profile } of userFetches) {
|
||||||
|
const repAttrs = formatReputationAttrs(rep);
|
||||||
|
userContexts.set(uid, `<user_reputation ${repAttrs}/>`);
|
||||||
|
userProfiles.set(uid, {
|
||||||
|
text: profile?.profile_summary ?? "",
|
||||||
|
asOf: profile?.last_analyzed_at ?? null,
|
||||||
|
});
|
||||||
}
|
}
|
||||||
const userProfilesBlock = buildUserProfilesBlock(userProfiles);
|
const userProfilesBlock = buildUserProfilesBlock(userProfiles);
|
||||||
|
|
||||||
|
|||||||
@@ -177,7 +177,7 @@ export const configSchema = z
|
|||||||
QDRANT_URL: z.string().optional(),
|
QDRANT_URL: z.string().optional(),
|
||||||
QDRANT_COLLECTION: z.string().default("gmw_text_moderation"),
|
QDRANT_COLLECTION: z.string().default("gmw_text_moderation"),
|
||||||
QDRANT_API_KEY: z.string().optional(),
|
QDRANT_API_KEY: z.string().optional(),
|
||||||
AI_LLM_MAX_CONCURRENT: z.coerce.number().int().positive().default(5),
|
AI_LLM_MAX_CONCURRENT: z.coerce.number().int().positive().default(8),
|
||||||
AI_LLM_IMAGE_MAX_DIMENSION: z.coerce
|
AI_LLM_IMAGE_MAX_DIMENSION: z.coerce
|
||||||
.number()
|
.number()
|
||||||
.int()
|
.int()
|
||||||
@@ -226,11 +226,11 @@ export const configSchema = z
|
|||||||
.default(5),
|
.default(5),
|
||||||
|
|
||||||
// ── AI Analysis Timing ──────────────────────────────────────────────
|
// ── AI Analysis Timing ──────────────────────────────────────────────
|
||||||
AI_ANALYSIS_DEBOUNCE_MS: z.coerce.number().positive().default(500),
|
AI_ANALYSIS_DEBOUNCE_MS: z.coerce.number().positive().default(250),
|
||||||
AI_ANALYSIS_RECOVERY_INTERVAL_MS: z.coerce
|
AI_ANALYSIS_RECOVERY_INTERVAL_MS: z.coerce
|
||||||
.number()
|
.number()
|
||||||
.positive()
|
.positive()
|
||||||
.default(15000),
|
.default(10000),
|
||||||
AI_ANALYSIS_ERROR_COOLDOWN_MS: z.coerce.number().positive().default(30000),
|
AI_ANALYSIS_ERROR_COOLDOWN_MS: z.coerce.number().positive().default(30000),
|
||||||
|
|
||||||
// ── AI Analysis Batch ───────────────────────────────────────────────
|
// ── AI Analysis Batch ───────────────────────────────────────────────
|
||||||
|
|||||||
Reference in New Issue
Block a user