refactor(gateway): split bootstrap + aiAnalyzer, add module barrels

app/:
- bootstrap.ts 277 -> 145 lines: config guard, DB connect, client debug
  logging and startup order are now named steps with a comment header
- lifecycle.ts (new): everything wired on the Discord 'ready' hook, in
  explicit order (inject broadcaster -> register listeners -> start workers)
- process-guards.ts (new): SIGINT/SIGTERM/uncaughtException/unhandledRejection
  in ONE place, using isTransientStreamError() instead of two duplicated
  inline code lists
- metrics-collector.ts (new): AI pipeline Prometheus gauges

modules/:
- ai-moderation/index.ts + message-capture/index.ts (new): public facades so
  app/ never reaches into internal files
- aiAnalyzer.ts 317 -> 146 lines: pure entry API; skip-verdict recording
  extracted into recordSkip()
- recovery-worker.ts (new): stranded-message recovery + stale lane/CB pruning
- cache-prune.ts (new): 6h expired-verdict sweep, throttled + resettable
- drop 3 dead re-exports (pickBatchWithinBudget/onCircuitBreakerAlert/
  getConversationKey) whose consumers import the origin files directly

No behavior change. typecheck + lint + 138 tests green; nix build OK.
This commit is contained in:
asepharyana
2026-09-24 15:04:53 +07:00
parent 6bf3b40cc7
commit 494e16b3b3
9 changed files with 597 additions and 400 deletions
+71 -203
View File
@@ -1,36 +1,19 @@
import { Client } from "discord.js-selfbot-v13"; import { Client } from "discord.js-selfbot-v13";
import { ConfigError, DatabaseError } from "@/shared/errors/index";
import { createChildLogger } from "@/shared/logger/index";
import { import {
getAnalysisQueueStatus, ConfigError,
startPendingAIAnalysisWorker, DatabaseError,
} from "../modules/ai-moderation/aiAnalyzer.js"; errorMessage,
import { } from "@/shared/errors/index.js";
mediaWorkerPool, import { createChildLogger } from "@/shared/logger/index.js";
textWorkerPool,
} from "../modules/ai-moderation/circuitBreaker.js";
import { registerChannelTopicCapture } from "../modules/channel-topic/index.js";
import { CommandHandler } from "../modules/command-handler/commandHandler.js"; import { CommandHandler } from "../modules/command-handler/commandHandler.js";
import { import {
EventBroadcaster, EventBroadcaster,
RedisEventPublisher, RedisEventPublisher,
} from "../modules/event-broadcaster/index.js"; } from "../modules/event-broadcaster/index.js";
import { import {
registerCollector,
setGauge,
startMetricsServer, startMetricsServer,
stopMetricsServer, stopMetricsServer,
} from "../modules/gateway-metrics/index.js"; } from "../modules/gateway-metrics/index.js";
import { registerGuildMemberEvents } from "../modules/guild-member-events/index.js";
import {
registerMessageCapture,
setEventBroadcaster as setMessageCaptureEventBroadcaster,
} from "../modules/message-capture/messageCapture.js";
import { setModerationEventBroadcaster } from "../modules/message-capture/moderationActionsDb.js";
import { startDigestScheduler } from "../modules/monitor/digestScheduler.js";
import { registerReactionCapture } from "../modules/reaction-tracking/index.js";
import { registerThreadCapture } from "../modules/thread-tracking/index.js";
import { registerPresenceCapture } from "../modules/user-presence/index.js";
import { config } from "../shared/config/index.js"; import { config } from "../shared/config/index.js";
import { import {
closeDatabase, closeDatabase,
@@ -38,45 +21,34 @@ import {
} from "../shared/database/drizzle.js"; } from "../shared/database/drizzle.js";
import { runMigrations } from "../shared/database/migrate.js"; import { runMigrations } from "../shared/database/migrate.js";
import { createDiscordClientOptions } from "../shared/discord/clientOptions.js"; import { createDiscordClientOptions } from "../shared/discord/clientOptions.js";
import { startRetentionCleanup } from "./retention.js"; import { startGatewayLifecycle } from "./lifecycle.js";
import { registerPipelineMetrics } from "./metrics-collector.js";
import { registerProcessGuards } from "./process-guards.js";
import { createGracefulShutdown } from "./shutdown.js"; import { createGracefulShutdown } from "./shutdown.js";
const logger = createChildLogger("discord-gateway"); const logger = createChildLogger("discord-gateway");
// ─── Bootstrap ───────────────────────────────────────────────────────────── // ─── Bootstrap ─────────────────────────────────────────────────────────────
//
// Startup order:
// 1. validate config (fail fast on missing AI credentials)
// 2. connect infrastructure (migrations → DB pool)
// 3. build long-lived services (Discord client, Redis publisher, command
// handler) + install shutdown/process guards
// 4. start observability (pipeline gauges → metrics server)
// 5. log in (ready-hook wires listeners via lifecycle.ts)
export async function initializeDiscordGateway() { /** Refuse to start when AI analysis is on but no LLM credentials exist. */
function assertConfigIsUsable(): void {
if (config.AI_ANALYSIS_ENABLED && !config.AI_LLM_API_KEY) { if (config.AI_ANALYSIS_ENABLED && !config.AI_LLM_API_KEY) {
throw new ConfigError( throw new ConfigError(
"AI_ANALYSIS_ENABLED=true but AI_LLM_API_KEY is missing from environment. AI analysis cannot run without credentials.", "AI_ANALYSIS_ENABLED=true but AI_LLM_API_KEY is missing from environment. AI analysis cannot run without credentials.",
); );
} }
}
const token = config.DISCORD_TOKEN; /** Run migrations (when enabled) then open the PostgreSQL pool. */
logger.info( async function connectDatabase(): Promise<void> {
{ hasToken: token.length > 0, tokenLength: token.length },
"Config loaded",
);
logger.info("Creating Discord client");
const client = new Client(createDiscordClientOptions());
// Initialize Redis event broadcaster
const redisPublisher = new RedisEventPublisher(config.REDIS_URL, logger);
const eventBroadcaster = new EventBroadcaster(redisPublisher);
// Initialize Redis command handler for backend→gateway commands
const commandHandler = new CommandHandler();
const gracefulShutdown = createGracefulShutdown({
logger,
closeDatabase,
client,
eventBroadcaster,
commandHandler,
stopMetricsServer,
});
try { try {
if (config.AUTO_MIGRATE_ON_STARTUP) { if (config.AUTO_MIGRATE_ON_STARTUP) {
logger.info( logger.info(
@@ -90,183 +62,79 @@ export async function initializeDiscordGateway() {
logger.info("PostgreSQL database initialized"); logger.info("PostgreSQL database initialized");
} catch (err) { } catch (err) {
logger.error( logger.error(
{ err, errorMsg: err instanceof Error ? err.message : String(err) }, { err, errorMsg: errorMessage(err) },
"Failed to initialize database", "Failed to initialize database",
); );
throw new DatabaseError( throw new DatabaseError(
`Database initialization failed: ${err instanceof Error ? err.message : String(err)}`, `Database initialization failed: ${errorMessage(err)}`,
); );
} }
}
/** Log only client debug lines that carry signal (errors/streams, or VERBOSE). */
function registerClientDebugLogging(client: Client): void {
client.on("debug", (msg) => { client.on("debug", (msg) => {
if ( const lower = msg.toLowerCase();
msg.toLowerCase().includes("error") || if (lower.includes("error") || lower.includes("stream")) {
msg.toLowerCase().includes("stream")
) {
logger.info({ debugMsg: msg }, "Discord Client Debug"); logger.info({ debugMsg: msg }, "Discord Client Debug");
} else if (config.VERBOSE) { } else if (config.VERBOSE) {
logger.debug({ debugMsg: msg }, "Discord Client Debug"); logger.debug({ debugMsg: msg }, "Discord Client Debug");
} }
}); });
}
client.on("ready", async () => { export async function initializeDiscordGateway() {
assertConfigIsUsable();
const token = config.DISCORD_TOKEN;
logger.info(
{ hasToken: token.length > 0, tokenLength: token.length },
"Config loaded",
);
logger.info("Creating Discord client");
const client = new Client(createDiscordClientOptions());
// Long-lived services: Redis event broadcaster (gateway → backend) and the
// Redis command handler (backend → gateway).
const redisPublisher = new RedisEventPublisher(config.REDIS_URL, logger);
const eventBroadcaster = new EventBroadcaster(redisPublisher);
const commandHandler = new CommandHandler();
const gracefulShutdown = createGracefulShutdown({
logger,
closeDatabase,
client,
eventBroadcaster,
commandHandler,
stopMetricsServer,
});
await connectDatabase();
registerClientDebugLogging(client);
client.on("ready", () => {
logger.info({ user: client.user?.tag }, "Bot logged in"); logger.info({ user: client.user?.tag }, "Bot logged in");
setMessageCaptureEventBroadcaster(eventBroadcaster); startGatewayLifecycle({
setModerationEventBroadcaster(eventBroadcaster); client,
registerMessageCapture(client); eventBroadcaster,
startPendingAIAnalysisWorker(client, eventBroadcaster); commandHandler,
logger,
// Register new event captures });
registerReactionCapture(client, eventBroadcaster);
registerThreadCapture(client, eventBroadcaster);
registerPresenceCapture(client, eventBroadcaster);
registerChannelTopicCapture(client, eventBroadcaster);
registerGuildMemberEvents(client, eventBroadcaster);
// Start command handler after Discord is ready
commandHandler.start(client);
logger.info("Command handler started");
// Start retention cleanup scheduler
startRetentionCleanup();
// Start weekly moderation digest (public, automated)
startDigestScheduler();
}); });
client.on("error", (err) => { client.on("error", (err) => {
logger.error( logger.error({ err, errorMsg: errorMessage(err) }, "Client error");
{ err, errorMsg: err instanceof Error ? err.message : String(err) },
"Client error",
);
}); });
process.on("SIGINT", () => { registerProcessGuards(logger, gracefulShutdown);
gracefulShutdown("SIGINT");
});
process.on("SIGTERM", () => { // Metrics: register live pipeline collectors before starting the server.
gracefulShutdown("SIGTERM"); registerPipelineMetrics(logger);
});
process.on("uncaughtException", (err) => {
const code =
typeof (err as NodeJS.ErrnoException).code === "string"
? (err as NodeJS.ErrnoException).code
: "";
// Transient stream-teardown errors (voice stop/disconnect races, child
// process stdin closed while we still write) are NOT fatal — crashing the
// gateway on EPIPE takes the whole bot offline mid-music. Log + continue.
if (
code === "EPIPE" ||
code === "ERR_STREAM_DESTROYED" ||
code === "ERR_STREAM_WRITE_AFTER_END" ||
code === "ECONNRESET"
) {
logger.warn(
{ error: err },
"Uncaught transient stream error — continuing",
);
return;
}
logger.error(
{
err,
errorMsg: err instanceof Error ? err.message : String(err),
stack: err?.stack,
},
"Uncaught exception",
);
gracefulShutdown("uncaughtException");
});
process.on("unhandledRejection", (reason) => {
const err =
reason instanceof Error ? reason : new Error(String(reason ?? "unknown"));
const code = (err as NodeJS.ErrnoException).code ?? "";
// Same transient-teardown policy as uncaughtException: a rejection that
// fires while a stream is being torn down (EPIPE after ffmpeg stdin
// closes, write-after-destroy, socket reset) must NOT take the whole
// gateway offline. Log detail + continue. Everything else still shuts
// down so real bugs surface.
if (
code === "EPIPE" ||
code === "ERR_STREAM_DESTROYED" ||
code === "ERR_STREAM_WRITE_AFTER_END" ||
code === "ECONNRESET"
) {
logger.warn(
{ error: err },
"Unhandled rejection transient stream error — continuing",
);
return;
}
logger.error({ error: err, reason: String(reason) }, "Unhandled rejection");
gracefulShutdown("unhandledRejection");
});
// ── Metrics: register live pipeline collectors before starting server ──
// These refresh on every scrape so Prometheus sees real AI-analysis
// queue depth, concurrency, and DB pool state instead of an empty stub.
registerCollector(() => {
if (!config.AI_ANALYSIS_ENABLED) return;
try {
const status = getAnalysisQueueStatus();
setGauge("ai_analysis_queued_conversations", status.queuedConversations);
setGauge("ai_analysis_active_batch_requests", status.activeRequests);
setGauge(
"ai_analysis_active_text_requests",
status.activeTextRequests ?? status.activeRequests,
);
setGauge(
"ai_analysis_active_media_requests",
status.activeMediaRequests ?? 0,
);
setGauge(
"ai_analysis_active_individual_requests",
status.activeIndividualRequests,
);
setGauge(
"ai_analysis_individual_in_flight",
status.individualInFlightCount,
);
setGauge(
"ai_analysis_individual_circuit_breaker_active",
status.individualCircuitBreakerActive ? 1 : 0,
);
if (typeof status.lastError === "string") {
setGauge("ai_analysis_last_error_present", status.lastError ? 1 : 0);
}
type PoolState = { _poolState?: { size: number; active: number } };
const textPool = textWorkerPool as unknown as PoolState;
const mediaPool = mediaWorkerPool as unknown as PoolState;
// Reported per queue (2026-08-31 text/media pool split) so the text
// and media backlogs are distinguishable in dashboards/alerts instead
// of one combined "worker threads" number.
if (textPool._poolState) {
setGauge("ai_analysis_worker_threads_text", textPool._poolState.size);
setGauge(
"ai_analysis_worker_threads_active_text",
textPool._poolState.active,
);
}
if (mediaPool._poolState) {
setGauge("ai_analysis_worker_threads_media", mediaPool._poolState.size);
setGauge(
"ai_analysis_worker_threads_active_media",
mediaPool._poolState.active,
);
}
} catch (err) {
logger.warn({ error: String(err) }, "AI metrics collector failed");
}
});
// Start metrics server
startMetricsServer(); startMetricsServer();
logger.info("Calling Discord client.login"); logger.info("Calling Discord client.login");
// Fix: use await + try/catch instead of .then().catch()
try { try {
await client.login(token); await client.login(token);
logger.info("Discord client logged in successfully"); logger.info("Discord client logged in successfully");
@@ -0,0 +1,61 @@
import type { Client } from "discord.js-selfbot-v13";
import type { Logger } from "@/shared/logger/index.js";
import { startPendingAIAnalysisWorker } from "../modules/ai-moderation/index.js";
import { registerChannelTopicCapture } from "../modules/channel-topic/index.js";
import type { CommandHandler } from "../modules/command-handler/commandHandler.js";
import type { EventBroadcaster } from "../modules/event-broadcaster/index.js";
import { registerGuildMemberEvents } from "../modules/guild-member-events/index.js";
import {
registerMessageCapture,
setEventBroadcaster as setMessageCaptureEventBroadcaster,
setModerationEventBroadcaster,
} from "../modules/message-capture/index.js";
import { startDigestScheduler } from "../modules/monitor/digestScheduler.js";
import { registerReactionCapture } from "../modules/reaction-tracking/index.js";
import { registerThreadCapture } from "../modules/thread-tracking/index.js";
import { registerPresenceCapture } from "../modules/user-presence/index.js";
import { startRetentionCleanup } from "./retention.js";
export interface GatewayLifecycleOptions {
client: Client;
eventBroadcaster: EventBroadcaster;
commandHandler: CommandHandler;
logger: Logger;
}
/**
* Wires everything that must start once Discord is connected.
*
* Ordering matters:
* 1. Inject the event broadcaster into the modules that publish events —
* they must be able to publish before their listeners are registered.
* 2. Register the Discord event listeners (capture modules).
* 3. Start the background workers/schedulers.
*/
export function startGatewayLifecycle({
client,
eventBroadcaster,
commandHandler,
logger,
}: GatewayLifecycleOptions): void {
// 1. Inject broadcaster first so no captured event is dropped.
setMessageCaptureEventBroadcaster(eventBroadcaster);
setModerationEventBroadcaster(eventBroadcaster);
// 2. Discord event listeners.
registerMessageCapture(client);
registerReactionCapture(client, eventBroadcaster);
registerThreadCapture(client, eventBroadcaster);
registerPresenceCapture(client, eventBroadcaster);
registerChannelTopicCapture(client, eventBroadcaster);
registerGuildMemberEvents(client, eventBroadcaster);
// 3. Background workers + schedulers.
startPendingAIAnalysisWorker(client, eventBroadcaster);
commandHandler.start(client);
logger.info("Command handler started");
startRetentionCleanup();
// Weekly moderation digest (public, automated)
startDigestScheduler();
}
@@ -0,0 +1,77 @@
import type { Logger } from "@/shared/logger/index.js";
import {
getAnalysisQueueStatus,
mediaWorkerPool,
textWorkerPool,
} from "../modules/ai-moderation/index.js";
import {
registerCollector,
setGauge,
} from "../modules/gateway-metrics/index.js";
import { config } from "../shared/config/index.js";
/** Piscina exposes its live thread counters on `_poolState`. */
type PoolState = { _poolState?: { size: number; active: number } };
/**
* Registers the AI-pipeline Prometheus gauges.
*
* The collector refreshes on every scrape, so Prometheus sees real queue
* depth / concurrency / worker-thread state instead of an empty stub.
* Registered before the metrics server starts.
*/
export function registerPipelineMetrics(logger: Logger): void {
registerCollector(() => {
if (!config.AI_ANALYSIS_ENABLED) return;
try {
const status = getAnalysisQueueStatus();
setGauge("ai_analysis_queued_conversations", status.queuedConversations);
setGauge("ai_analysis_active_batch_requests", status.activeRequests);
setGauge(
"ai_analysis_active_text_requests",
status.activeTextRequests ?? status.activeRequests,
);
setGauge(
"ai_analysis_active_media_requests",
status.activeMediaRequests ?? 0,
);
setGauge(
"ai_analysis_active_individual_requests",
status.activeIndividualRequests,
);
setGauge(
"ai_analysis_individual_in_flight",
status.individualInFlightCount,
);
setGauge(
"ai_analysis_individual_circuit_breaker_active",
status.individualCircuitBreakerActive ? 1 : 0,
);
if (typeof status.lastError === "string") {
setGauge("ai_analysis_last_error_present", status.lastError ? 1 : 0);
}
// Reported per queue (2026-08-31 text/media pool split) so the text
// and media backlogs are distinguishable in dashboards/alerts instead
// of one combined "worker threads" number.
const textPool = textWorkerPool as unknown as PoolState;
const mediaPool = mediaWorkerPool as unknown as PoolState;
if (textPool._poolState) {
setGauge("ai_analysis_worker_threads_text", textPool._poolState.size);
setGauge(
"ai_analysis_worker_threads_active_text",
textPool._poolState.active,
);
}
if (mediaPool._poolState) {
setGauge("ai_analysis_worker_threads_media", mediaPool._poolState.size);
setGauge(
"ai_analysis_worker_threads_active_media",
mediaPool._poolState.active,
);
}
} catch (err) {
logger.warn({ error: String(err) }, "AI metrics collector failed");
}
});
}
@@ -0,0 +1,61 @@
import { errorMessage, isTransientStreamError } from "@/shared/errors/index.js";
import type { Logger } from "@/shared/logger/index.js";
import type { GracefulShutdown } from "./shutdown.js";
/**
* Process-level signal + error guards.
*
* Extracted from bootstrap so the "what keeps the gateway alive vs what
* shuts it down" policy lives in exactly one place.
*
* Policy: transient stream-teardown failures (EPIPE / ERR_STREAM_DESTROYED /
* ERR_STREAM_WRITE_AFTER_END / ECONNRESET) are logged and IGNORED — crashing
* the gateway on them (voice stop/disconnect races, a child process stdin
* closed while we still write) takes the whole bot offline mid-operation.
* Anything else is a real bug: log with stack and shut down cleanly.
*/
export function registerProcessGuards(
logger: Logger,
gracefulShutdown: GracefulShutdown,
): void {
process.on("SIGINT", () => {
gracefulShutdown("SIGINT");
});
process.on("SIGTERM", () => {
gracefulShutdown("SIGTERM");
});
process.on("uncaughtException", (err) => {
if (isTransientStreamError(err)) {
logger.warn(
{ error: err },
"Uncaught transient stream error — continuing",
);
return;
}
logger.error(
{
err,
errorMsg: errorMessage(err),
stack: err?.stack,
},
"Uncaught exception",
);
gracefulShutdown("uncaughtException");
});
process.on("unhandledRejection", (reason) => {
const err =
reason instanceof Error ? reason : new Error(String(reason ?? "unknown"));
if (isTransientStreamError(err)) {
logger.warn(
{ error: err },
"Unhandled rejection transient stream error — continuing",
);
return;
}
logger.error({ error: err, reason: String(reason) }, "Unhandled rejection");
gracefulShutdown("unhandledRejection");
});
}
@@ -12,28 +12,14 @@ import {
buildSkipAnalysisUserResult, buildSkipAnalysisUserResult,
isAgeRestrictedMessage, isAgeRestrictedMessage,
isSkipAnalysisUser, isSkipAnalysisUser,
skipAgeRestrictedMessages,
skipAnalysisUserMessages,
} from "./batchProcessor.js"; } from "./batchProcessor.js";
import { scheduleConversationAnalysis } from "./batchScheduler.js"; import { scheduleConversationAnalysis } from "./batchScheduler.js";
import { getConversationKey } from "./circuitBreaker.js"; import { getConversationKey } from "./circuitBreaker.js";
import { import { conversationDebounceTimers } from "./conversationState.js";
ANALYSIS_LANES,
type AnalysisLane,
clearConversationProcessing,
conversationConsecutiveErrors,
conversationDebounceTimers,
conversationErrorCooldown,
conversationProcessing,
isConversationProcessingLocked,
} from "./conversationState.js";
import { import {
activeIndividualRequests, activeIndividualRequests,
enqueueIndividualFallbacks,
individualCooldownUntil, individualCooldownUntil,
individualInFlight, individualInFlight,
individualInFlightByConversation,
individualInFlightLastTouched,
} from "./individualFallbackProcessor.js"; } from "./individualFallbackProcessor.js";
import { import {
broadcastAnalysisCompleted, broadcastAnalysisCompleted,
@@ -41,31 +27,20 @@ import {
setModerationClient, setModerationClient,
setSharedEventBroadcaster, setSharedEventBroadcaster,
} from "./moderationState.js"; } from "./moderationState.js";
import { deleteExpiredQdrantPoints } from "./qdrantClient.js"; import { startRecoveryWorker } from "./recovery-worker.js";
import { pruneExpiredTexts } from "./textCacheStore.js";
const logger = createChildLogger("ai-analyzer"); const logger = createChildLogger("ai-analyzer");
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// Cache hygiene (expired verdict sweep) // Public API — queueing, status, worker startup
// ---------------------------------------------------------------------------
const CACHE_PRUNE_INTERVAL_MS = 6 * 60 * 60 * 1000; // every 6 hours
let lastCachePruneAt = 0;
// ---------------------------------------------------------------------------
// Re-exports from sub-modules (preserving original public API)
// ---------------------------------------------------------------------------
export { pickBatchWithinBudget } from "./batchProcessor.js";
export { getConversationKey } from "./circuitBreaker.js";
export { onCircuitBreakerAlert } from "./conversationState.js";
// ---------------------------------------------------------------------------
// Public API
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
/** /**
* Queues a message for analysis (debounced by conversation). * Queues a message for analysis (debounced by conversation).
*
* Messages that never need an LLM call are short-circuited here and recorded
* with their skip verdict: age-restricted messages and configured skip-list
* users.
*/ */
export async function queueMessageAnalysis(messageId: string): Promise<void> { export async function queueMessageAnalysis(messageId: string): Promise<void> {
if (!config.AI_ANALYSIS_ENABLED) return; if (!config.AI_ANALYSIS_ENABLED) return;
@@ -78,13 +53,7 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
} }
if (isAgeRestrictedMessage(message)) { if (isAgeRestrictedMessage(message)) {
const updated = await messageStore.updateMessageAIAnalysis( await recordSkip(message.id, buildAgeRestrictedSkipResult());
message.id,
buildAgeRestrictedSkipResult(),
);
if (updated) {
broadcastAnalysisCompleted(updated);
}
logger.debug( logger.debug(
{ messageId }, { messageId },
"Skipped AI analysis for age-restricted message", "Skipped AI analysis for age-restricted message",
@@ -93,13 +62,7 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
} }
if (isSkipAnalysisUser(message)) { if (isSkipAnalysisUser(message)) {
const updated = await messageStore.updateMessageAIAnalysis( await recordSkip(message.id, buildSkipAnalysisUserResult());
message.id,
buildSkipAnalysisUserResult(),
);
if (updated) {
broadcastAnalysisCompleted(updated);
}
logger.debug( logger.debug(
{ messageId, userId: message.user_id }, { messageId, userId: message.user_id },
"Skipped AI analysis for configured skip-list user", "Skipped AI analysis for configured skip-list user",
@@ -119,6 +82,17 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
} }
} }
/** Persist a skip verdict and broadcast it so the dashboard reflects it. */
async function recordSkip(
messageId: string,
result: Parameters<typeof messageStore.updateMessageAIAnalysis>[1],
): Promise<void> {
const updated = await messageStore.updateMessageAIAnalysis(messageId, result);
if (updated) {
broadcastAnalysisCompleted(updated);
}
}
/** /**
* Queues a conversation for analysis (debounced). * Queues a conversation for analysis (debounced).
*/ */
@@ -144,11 +118,12 @@ export function getAnalysisQueueStatus(): AnalysisQueueStatus {
} }
/** /**
* Starts the periodic recovery worker. * Starts the background workers behind the analysis pipeline:
* - the recovery worker (stranded pending / incomplete messages + cache prune)
* - the optional culture and user-profile learners.
* *
* Now also recovers messages stuck in `error/analysis_incomplete` * Also injects the Discord client and event broadcaster into the pipeline
* state (not just `pending`), and skips conversations that already have * state so downstream modules can act and publish.
* individual fallback work in progress to avoid DB last-write-wins races.
*/ */
export function startPendingAIAnalysisWorker( export function startPendingAIAnalysisWorker(
client?: Client, client?: Client,
@@ -167,151 +142,5 @@ export function startPendingAIAnalysisWorker(
.catch(console.error); .catch(console.error);
} }
setInterval(() => { startRecoveryWorker();
// [D] Periodic cache hygiene: purge expired moderation verdicts from
// Postgres and Qdrant. Expired entries are never reused (filters check
// expires_at) but accumulate forever without this sweep.
const now = Date.now();
if (now - lastCachePruneAt >= CACHE_PRUNE_INTERVAL_MS) {
lastCachePruneAt = now;
Promise.all([pruneExpiredTexts(), deleteExpiredQdrantPoints()])
.then(([pgDeleted, qdDeleted]) => {
if (pgDeleted > 0 || qdDeleted > 0) {
logger.info(
{ pgDeleted, qdDeleted },
"Expired moderation cache pruned",
);
}
})
.catch((err: unknown) => {
logger.warn({ error: String(err) }, "Moderation cache prune failed");
});
}
// Only revert stuck processing messages if there's active processing.
// Avoids a DB query every recovery interval when the pipeline is idle.
if (conversationProcessing.size > 0) {
messageStore
.revertStuckProcessingMessages(300000)
.catch((err: unknown) => {
logger.error(
{ error: String(err) },
"Failed to run stuck processing recovery",
);
});
}
Promise.all([
messageStore.getPendingConversationKeys(500),
messageStore.getConversationKeysWithIncompleteAnalysis(200),
])
.then(([pendingKeys, incompleteKeys]) => {
const now = Date.now();
for (const [key, expiry] of conversationErrorCooldown) {
if (now >= expiry) conversationErrorCooldown.delete(key);
}
// conversationProcessing is now a Partial<Record<lane, startedAt>>.
// Prune stale lane slots individually so one stale lane never clears
// the other lane's healthy lock.
for (const [key, record] of conversationProcessing) {
for (const lane of ANALYSIS_LANES as readonly AnalysisLane[]) {
const startedAt = record?.[lane];
if (
startedAt &&
now - startedAt >= config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS
) {
clearConversationProcessing(key, lane);
}
}
}
const staleThreshold = config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS * 2;
for (const [key, lastTouched] of individualInFlightLastTouched) {
if (now - lastTouched >= staleThreshold) {
individualInFlightLastTouched.delete(key);
individualInFlightByConversation.delete(key);
logger.warn(
{ key },
"Pruned stale individualInFlightByConversation entry",
);
}
}
// Also prune stale per-conversation CB error counts that have cooled
// down so old conversations can be retried.
for (const [key] of conversationConsecutiveErrors) {
const cbExpire = conversationErrorCooldown.get(key) ?? 0;
if (cbExpire && now >= cbExpire) {
conversationConsecutiveErrors.delete(key);
}
}
const incompleteKeySet = new Set(incompleteKeys);
// --- Batch recovery for pending messages ---
for (const key of pendingKeys) {
if (
ANALYSIS_LANES.some((lane) =>
conversationDebounceTimers.has(`${key}::${lane}`),
)
) {
continue;
}
// Batch recovery must not race ANY in-flight batch lane, so the
// lock check is lane-agnostic here (individual fallback handles
// error rows separately).
if (isConversationProcessingLocked(key)) continue;
if (individualInFlightByConversation.has(key)) continue;
if (incompleteKeySet.has(key)) continue;
const cooldownUntil = conversationErrorCooldown.get(key);
if (cooldownUntil && now < cooldownUntil) continue;
// No lane specified → schedule BOTH lanes; each fetches its own
// pending subset from the DB.
scheduleConversationAnalysis(key);
}
// --- Individual recovery for error/analysis_incomplete messages ---
// Circuit breaker check: no point iterating if individual CB is active.
if (now >= individualCooldownUntil) {
const promises: Promise<void>[] = [];
for (const key of incompleteKeys) {
// Skip if individual work is already running for this conversation.
if (individualInFlightByConversation.has(key)) continue;
// Skip if batch processing is running.
if (isConversationProcessingLocked(key)) continue;
promises.push(
messageStore
.getIncompleteMessagesByConversation(key, 500)
.then(async (msgs) => {
const processableMessages = await skipAnalysisUserMessages(
await skipAgeRestrictedMessages(msgs),
);
return processableMessages;
})
.then((msgs) => {
if (msgs.length > 0) {
enqueueIndividualFallbacks(msgs);
}
})
.catch((err: unknown) => {
logger.error(
{ key, error: String(err) },
"Failed to fetch incomplete messages for recovery",
);
}),
);
}
// Errors are handled per-key; return the combined promise for observability.
return Promise.all(promises);
}
})
.catch((err: unknown) => {
logger.error(
{ error: err instanceof Error ? err.message : String(err) },
"Pending AI analysis recovery worker failed",
);
});
}, config.AI_ANALYSIS_RECOVERY_INTERVAL_MS);
} }
@@ -0,0 +1,41 @@
import { createChildLogger } from "@/shared/logger/index";
import { config } from "../../shared/config/index.js";
import { deleteExpiredQdrantPoints } from "./qdrantClient.js";
import { pruneExpiredTexts } from "./textCacheStore.js";
const logger = createChildLogger("cache-prune");
/** Expired-verdict sweep cadence. */
const CACHE_PRUNE_INTERVAL_MS = 6 * 60 * 60 * 1000; // every 6 hours
let lastCachePruneAt = 0;
/**
* Cache hygiene: purge expired moderation verdicts from Postgres and Qdrant.
*
* Expired entries are never reused (read filters check `expires_at`) but they
* accumulate forever without a sweep. Called from the recovery interval; the
* 6-hour throttle keeps it to one sweep per window.
*/
export function runCachePruneIfDue(now: number = Date.now()): void {
if (now - lastCachePruneAt < CACHE_PRUNE_INTERVAL_MS) return;
lastCachePruneAt = now;
Promise.all([pruneExpiredTexts(), deleteExpiredQdrantPoints()])
.then(([pgDeleted, qdDeleted]) => {
if (pgDeleted > 0 || qdDeleted > 0) {
logger.info(
{ pgDeleted, qdDeleted },
"Expired moderation cache pruned",
);
}
})
.catch((err: unknown) => {
logger.warn({ error: String(err) }, "Moderation cache prune failed");
});
}
/** Reset the throttle window (tests). */
export function resetCachePruneState(): void {
lastCachePruneAt = 0;
}
@@ -0,0 +1,35 @@
/**
* Public surface of the AI-moderation module.
*
* The module has ~50 internal files; callers outside it (app/, tests, other
* modules) should import from THIS barrel so internal files can be moved
* without touching call sites.
*
* Deep imports remain valid inside the module itself.
*/
export type {
AIRecommendedAction,
AISeverity,
AIStatus,
AnalysisQueueStatus,
AnalysisResult,
} from "../../shared/moderation-types.js";
// ── Entry API: queueing, status, recovery worker ──────────────────────────
export {
getAnalysisQueueStatus,
queueConversationAnalysis,
queueMessageAnalysis,
startPendingAIAnalysisWorker,
} from "./aiAnalyzer.js";
// ── Worker pools (app/metrics-collector reads their live thread counters) ──
export {
getConversationKey,
mediaWorkerPool,
textWorkerPool,
} from "./circuitBreaker.js";
// ── Pipeline state hooks the bootstrap injects into ───────────────────────
export {
setModerationClient,
setSharedEventBroadcaster,
} from "./moderationState.js";
@@ -0,0 +1,184 @@
import { createChildLogger } from "@/shared/logger/index";
import { config } from "../../shared/config/index.js";
import { messageStore } from "../message-capture/messageStore.js";
import {
skipAgeRestrictedMessages,
skipAnalysisUserMessages,
} from "./batchProcessor.js";
import { scheduleConversationAnalysis } from "./batchScheduler.js";
import { runCachePruneIfDue } from "./cache-prune.js";
import {
ANALYSIS_LANES,
type AnalysisLane,
clearConversationProcessing,
conversationConsecutiveErrors,
conversationDebounceTimers,
conversationErrorCooldown,
conversationProcessing,
isConversationProcessingLocked,
} from "./conversationState.js";
import {
enqueueIndividualFallbacks,
individualCooldownUntil,
individualInFlightByConversation,
individualInFlightLastTouched,
} from "./individualFallbackProcessor.js";
const logger = createChildLogger("ai-recovery");
/** Revert messages stuck in `processing` for longer than this. */
const STUCK_PROCESSING_AGE_MS = 300_000;
/**
* Starts the periodic recovery worker.
*
* Recovers two classes of stranded work:
* - `pending` messages → re-scheduled through the normal per-lane debounce.
* - `error` / `analysis_incomplete` messages → individual fallback queue.
*
* Also prunes stale in-memory bookkeeping (lane locks, per-conversation
* circuit-breaker counters, individual-fallback in-flight markers) so a
* crashed batch cannot wedge a conversation forever, and triggers the
* throttled cache-prune sweep.
*
* Skips conversations that already have individual fallback work in progress
* to avoid DB last-write-wins races.
*/
export function startRecoveryWorker(): void {
setInterval(() => {
runCachePruneIfDue();
// Only revert stuck processing messages if there's active processing.
// Avoids a DB query every recovery interval when the pipeline is idle.
if (conversationProcessing.size > 0) {
messageStore
.revertStuckProcessingMessages(STUCK_PROCESSING_AGE_MS)
.catch((err: unknown) => {
logger.error(
{ error: String(err) },
"Failed to run stuck processing recovery",
);
});
}
Promise.all([
messageStore.getPendingConversationKeys(500),
messageStore.getConversationKeysWithIncompleteAnalysis(200),
])
.then(([pendingKeys, incompleteKeys]) => {
const now = Date.now();
pruneStaleConversationState(now);
const incompleteKeySet = new Set(incompleteKeys);
// --- Batch recovery for pending messages ---
for (const key of pendingKeys) {
if (
ANALYSIS_LANES.some((lane) =>
conversationDebounceTimers.has(`${key}::${lane}`),
)
) {
continue;
}
// Batch recovery must not race ANY in-flight batch lane, so the
// lock check is lane-agnostic here (individual fallback handles
// error rows separately).
if (isConversationProcessingLocked(key)) continue;
if (individualInFlightByConversation.has(key)) continue;
if (incompleteKeySet.has(key)) continue;
const cooldownUntil = conversationErrorCooldown.get(key);
if (cooldownUntil && now < cooldownUntil) continue;
// No lane specified → schedule BOTH lanes; each fetches its own
// pending subset from the DB.
scheduleConversationAnalysis(key);
}
// --- Individual recovery for error/analysis_incomplete messages ---
// Circuit breaker check: no point iterating if individual CB is active.
if (now >= individualCooldownUntil) {
const promises: Promise<void>[] = [];
for (const key of incompleteKeys) {
// Skip if individual work is already running for this conversation.
if (individualInFlightByConversation.has(key)) continue;
// Skip if batch processing is running.
if (isConversationProcessingLocked(key)) continue;
promises.push(
recoverIncompleteConversation(key).catch((err: unknown) => {
logger.error(
{ key, error: String(err) },
"Failed to fetch incomplete messages for recovery",
);
}),
);
}
// Errors are handled per-key; return the combined promise for observability.
return Promise.all(promises);
}
})
.catch((err: unknown) => {
logger.error(
{ error: err instanceof Error ? err.message : String(err) },
"Pending AI analysis recovery worker failed",
);
});
}, config.AI_ANALYSIS_RECOVERY_INTERVAL_MS);
}
/** Fetch one conversation's incomplete messages and queue them individually. */
async function recoverIncompleteConversation(key: string): Promise<void> {
const msgs = await messageStore.getIncompleteMessagesByConversation(key, 500);
const processable = await skipAnalysisUserMessages(
await skipAgeRestrictedMessages(msgs),
);
if (processable.length > 0) {
enqueueIndividualFallbacks(processable);
}
}
/**
* Drop stale in-memory bookkeeping:
* - per-lane processing locks past the timeout (pruned PER LANE so one stale
* lane never clears the other lane's healthy lock),
* - individual-fallback in-flight markers that stopped being touched,
* - per-conversation circuit-breaker error counts whose cooldown has lapsed.
*/
function pruneStaleConversationState(now: number): void {
for (const [key, expiry] of conversationErrorCooldown) {
if (now >= expiry) conversationErrorCooldown.delete(key);
}
for (const [key, record] of conversationProcessing) {
for (const lane of ANALYSIS_LANES as readonly AnalysisLane[]) {
const startedAt = record?.[lane];
if (
startedAt &&
now - startedAt >= config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS
) {
clearConversationProcessing(key, lane);
}
}
}
const staleThreshold = config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS * 2;
for (const [key, lastTouched] of individualInFlightLastTouched) {
if (now - lastTouched >= staleThreshold) {
individualInFlightLastTouched.delete(key);
individualInFlightByConversation.delete(key);
logger.warn(
{ key },
"Pruned stale individualInFlightByConversation entry",
);
}
}
// Also prune stale per-conversation CB error counts that have cooled
// down so old conversations can be retried.
for (const [key] of conversationConsecutiveErrors) {
const cbExpire = conversationErrorCooldown.get(key) ?? 0;
if (cbExpire && now >= cbExpire) {
conversationConsecutiveErrors.delete(key);
}
}
}
@@ -0,0 +1,41 @@
/**
* Public surface of the message-capture module.
*
* Callers outside this module import from here instead of reaching into
* messageCapture.ts / moderationActionsDb.ts / messageStore.ts directly, so
* the internal file layout can change without touching call sites.
*
* Deep imports remain valid inside the module itself.
*/
// ── Discord listener registration (called from app/lifecycle.ts) ───────────
export {
registerMessageCapture,
setEventBroadcaster,
} from "./messageCapture.js";
// ── Message persistence facade ────────────────────────────────────────────
export { messageStore } from "./messageStore.js";
// ── Live moderation-action publishing ─────────────────────────────────────
export { setModerationEventBroadcaster } from "./moderationActionsDb.js";
// ── Domain types ──────────────────────────────────────────────────────────
export type {
AIRecommendedAction,
AISeverity,
AIStatus,
AnalysisQueueStatus,
AnalysisResult,
AttachmentRecord,
DashboardMessage,
MessageQuery,
MessageRecord,
MessageReview,
ModerationAction,
ModerationActionType,
ModerationWsEvent,
PageResult,
RetentionPolicy,
ReviewStatus,
RoleMetadata,
UserMetadata,
} from "./types.js";