diff --git a/services/discord-gateway/src/app/bootstrap.ts b/services/discord-gateway/src/app/bootstrap.ts index 460bef3d..21853388 100644 --- a/services/discord-gateway/src/app/bootstrap.ts +++ b/services/discord-gateway/src/app/bootstrap.ts @@ -1,36 +1,19 @@ import { Client } from "discord.js-selfbot-v13"; -import { ConfigError, DatabaseError } from "@/shared/errors/index"; -import { createChildLogger } from "@/shared/logger/index"; import { - getAnalysisQueueStatus, - startPendingAIAnalysisWorker, -} from "../modules/ai-moderation/aiAnalyzer.js"; -import { - mediaWorkerPool, - textWorkerPool, -} from "../modules/ai-moderation/circuitBreaker.js"; -import { registerChannelTopicCapture } from "../modules/channel-topic/index.js"; + ConfigError, + DatabaseError, + errorMessage, +} from "@/shared/errors/index.js"; +import { createChildLogger } from "@/shared/logger/index.js"; import { CommandHandler } from "../modules/command-handler/commandHandler.js"; import { EventBroadcaster, RedisEventPublisher, } from "../modules/event-broadcaster/index.js"; import { - registerCollector, - setGauge, startMetricsServer, stopMetricsServer, } 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 { closeDatabase, @@ -38,45 +21,34 @@ import { } from "../shared/database/drizzle.js"; import { runMigrations } from "../shared/database/migrate.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"; const logger = createChildLogger("discord-gateway"); // ─── 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) { throw new ConfigError( "AI_ANALYSIS_ENABLED=true but AI_LLM_API_KEY is missing from environment. AI analysis cannot run without credentials.", ); } +} - 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()); - - // 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, - }); - +/** Run migrations (when enabled) then open the PostgreSQL pool. */ +async function connectDatabase(): Promise { try { if (config.AUTO_MIGRATE_ON_STARTUP) { logger.info( @@ -90,183 +62,79 @@ export async function initializeDiscordGateway() { logger.info("PostgreSQL database initialized"); } catch (err) { logger.error( - { err, errorMsg: err instanceof Error ? err.message : String(err) }, + { err, errorMsg: errorMessage(err) }, "Failed to initialize database", ); 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) => { - if ( - msg.toLowerCase().includes("error") || - msg.toLowerCase().includes("stream") - ) { + const lower = msg.toLowerCase(); + if (lower.includes("error") || lower.includes("stream")) { logger.info({ debugMsg: msg }, "Discord Client Debug"); } else if (config.VERBOSE) { 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"); - setMessageCaptureEventBroadcaster(eventBroadcaster); - setModerationEventBroadcaster(eventBroadcaster); - registerMessageCapture(client); - startPendingAIAnalysisWorker(client, eventBroadcaster); - - // 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(); + startGatewayLifecycle({ + client, + eventBroadcaster, + commandHandler, + logger, + }); }); client.on("error", (err) => { - logger.error( - { err, errorMsg: err instanceof Error ? err.message : String(err) }, - "Client error", - ); + logger.error({ err, errorMsg: errorMessage(err) }, "Client error"); }); - process.on("SIGINT", () => { - gracefulShutdown("SIGINT"); - }); + registerProcessGuards(logger, gracefulShutdown); - process.on("SIGTERM", () => { - gracefulShutdown("SIGTERM"); - }); - - 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 + // Metrics: register live pipeline collectors before starting the server. + registerPipelineMetrics(logger); startMetricsServer(); logger.info("Calling Discord client.login"); - - // Fix: use await + try/catch instead of .then().catch() try { await client.login(token); logger.info("Discord client logged in successfully"); diff --git a/services/discord-gateway/src/app/lifecycle.ts b/services/discord-gateway/src/app/lifecycle.ts new file mode 100644 index 00000000..b14da663 --- /dev/null +++ b/services/discord-gateway/src/app/lifecycle.ts @@ -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(); +} diff --git a/services/discord-gateway/src/app/metrics-collector.ts b/services/discord-gateway/src/app/metrics-collector.ts new file mode 100644 index 00000000..b19a7e23 --- /dev/null +++ b/services/discord-gateway/src/app/metrics-collector.ts @@ -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"); + } + }); +} diff --git a/services/discord-gateway/src/app/process-guards.ts b/services/discord-gateway/src/app/process-guards.ts new file mode 100644 index 00000000..76106cf1 --- /dev/null +++ b/services/discord-gateway/src/app/process-guards.ts @@ -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"); + }); +} diff --git a/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts b/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts index bf51c6fb..ca187ace 100644 --- a/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts +++ b/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts @@ -12,28 +12,14 @@ import { buildSkipAnalysisUserResult, isAgeRestrictedMessage, isSkipAnalysisUser, - skipAgeRestrictedMessages, - skipAnalysisUserMessages, } from "./batchProcessor.js"; import { scheduleConversationAnalysis } from "./batchScheduler.js"; import { getConversationKey } from "./circuitBreaker.js"; -import { - ANALYSIS_LANES, - type AnalysisLane, - clearConversationProcessing, - conversationConsecutiveErrors, - conversationDebounceTimers, - conversationErrorCooldown, - conversationProcessing, - isConversationProcessingLocked, -} from "./conversationState.js"; +import { conversationDebounceTimers } from "./conversationState.js"; import { activeIndividualRequests, - enqueueIndividualFallbacks, individualCooldownUntil, individualInFlight, - individualInFlightByConversation, - individualInFlightLastTouched, } from "./individualFallbackProcessor.js"; import { broadcastAnalysisCompleted, @@ -41,31 +27,20 @@ import { setModerationClient, setSharedEventBroadcaster, } from "./moderationState.js"; -import { deleteExpiredQdrantPoints } from "./qdrantClient.js"; -import { pruneExpiredTexts } from "./textCacheStore.js"; +import { startRecoveryWorker } from "./recovery-worker.js"; const logger = createChildLogger("ai-analyzer"); // --------------------------------------------------------------------------- -// Cache hygiene (expired verdict sweep) -// --------------------------------------------------------------------------- -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 +// Public API — queueing, status, worker startup // --------------------------------------------------------------------------- /** * 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 { if (!config.AI_ANALYSIS_ENABLED) return; @@ -78,13 +53,7 @@ export async function queueMessageAnalysis(messageId: string): Promise { } if (isAgeRestrictedMessage(message)) { - const updated = await messageStore.updateMessageAIAnalysis( - message.id, - buildAgeRestrictedSkipResult(), - ); - if (updated) { - broadcastAnalysisCompleted(updated); - } + await recordSkip(message.id, buildAgeRestrictedSkipResult()); logger.debug( { messageId }, "Skipped AI analysis for age-restricted message", @@ -93,13 +62,7 @@ export async function queueMessageAnalysis(messageId: string): Promise { } if (isSkipAnalysisUser(message)) { - const updated = await messageStore.updateMessageAIAnalysis( - message.id, - buildSkipAnalysisUserResult(), - ); - if (updated) { - broadcastAnalysisCompleted(updated); - } + await recordSkip(message.id, buildSkipAnalysisUserResult()); logger.debug( { messageId, userId: message.user_id }, "Skipped AI analysis for configured skip-list user", @@ -119,6 +82,17 @@ export async function queueMessageAnalysis(messageId: string): Promise { } } +/** Persist a skip verdict and broadcast it so the dashboard reflects it. */ +async function recordSkip( + messageId: string, + result: Parameters[1], +): Promise { + const updated = await messageStore.updateMessageAIAnalysis(messageId, result); + if (updated) { + broadcastAnalysisCompleted(updated); + } +} + /** * 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` - * state (not just `pending`), and skips conversations that already have - * individual fallback work in progress to avoid DB last-write-wins races. + * Also injects the Discord client and event broadcaster into the pipeline + * state so downstream modules can act and publish. */ export function startPendingAIAnalysisWorker( client?: Client, @@ -167,151 +142,5 @@ export function startPendingAIAnalysisWorker( .catch(console.error); } - setInterval(() => { - // [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>. - // 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[] = []; - 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); + startRecoveryWorker(); } diff --git a/services/discord-gateway/src/modules/ai-moderation/cache-prune.ts b/services/discord-gateway/src/modules/ai-moderation/cache-prune.ts new file mode 100644 index 00000000..0dfad3f1 --- /dev/null +++ b/services/discord-gateway/src/modules/ai-moderation/cache-prune.ts @@ -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; +} diff --git a/services/discord-gateway/src/modules/ai-moderation/index.ts b/services/discord-gateway/src/modules/ai-moderation/index.ts new file mode 100644 index 00000000..93dd390d --- /dev/null +++ b/services/discord-gateway/src/modules/ai-moderation/index.ts @@ -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"; diff --git a/services/discord-gateway/src/modules/ai-moderation/recovery-worker.ts b/services/discord-gateway/src/modules/ai-moderation/recovery-worker.ts new file mode 100644 index 00000000..f0d77130 --- /dev/null +++ b/services/discord-gateway/src/modules/ai-moderation/recovery-worker.ts @@ -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[] = []; + 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 { + 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); + } + } +} diff --git a/services/discord-gateway/src/modules/message-capture/index.ts b/services/discord-gateway/src/modules/message-capture/index.ts new file mode 100644 index 00000000..c7d4d584 --- /dev/null +++ b/services/discord-gateway/src/modules/message-capture/index.ts @@ -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";