diff --git a/.env.example b/.env.example index 84b05306..df09933b 100644 --- a/.env.example +++ b/.env.example @@ -134,4 +134,7 @@ RETENTION_DRY_RUN=true # Dry-run: log but do not delete (defaul AUTO_MIGRATE_ON_STARTUP=true # Run database migrations on startup (default: true) # === Worker Pool === -# PISCINA_MAX_THREADS=4 # Worker thread pool size (optional, defaults to CPU cores) +# Text and media AI-analysis batches run on separate Piscina pools so a slow +# image/vision batch can never queue-block fast text-only batches. +# PISCINA_MAX_THREADS=4 # Text-analysis worker pool size (optional, defaults to CPU cores) +# PISCINA_MEDIA_MAX_THREADS=2 # Media-analysis worker pool size (optional, default 2) diff --git a/.hermes/plans/2026-08-31_video-receive-eager-selfbot-connection-spec.md b/.hermes/plans/2026-08-31_video-receive-eager-selfbot-connection-spec.md new file mode 100644 index 00000000..864147b4 --- /dev/null +++ b/.hermes/plans/2026-08-31_video-receive-eager-selfbot-connection-spec.md @@ -0,0 +1,93 @@ +# Spec: Fix Video Capture — Eagerly Establish the Selfbot Voice Connection at Join Time + +Status: PLANNED +Date: 2026-08-31 +Author: Hermes +Related: `.hermes/plans/2026-08-31_video-receive-phaseC-spec.md` (Phase C build, made Option A this fix) + +## Symptom (from live logs, 2026-08-31 ~12:34) +A user was actively screen-sharing + on camera in the recorded voice channel. +The gateway recorded MANY users' audio (.ogg) fine, but video capture produced +nothing. The only video signal in `journalctl -u gmw-discord-gateway` was: + +``` +[VOICE (guild:2)]: Sending voice state update: {"self_mute":false,...,"flags":2} +[VOICE] received voice state update: {member hunterz ...} # OTHER user, not bot +[VOICE] connection? true, guild session channel +[VOICE (guild:2)]: Setting sessionId (stored as "undefined") +[VOICE (guild:2)]: Authenticated with sessionId # debug print only +[VOICE (guild:2)]: Authenticate failed - VOICE_CONNECTION_TIMEOUT # +15s +video-recorder: userId=..., "Connection not established within 15 seconds." +``` + +## Root cause (verified against discord.js-selfbot-v13 3.7.1 source) +The gateway records audio via `@discordjs/voice` (`joinVoiceChannel` + adapter). +Video receive lives on the SEPARATE selfbot `ClientVoiceManager.connection` +(a singleton `VoiceConnection`). `videoRecorder.ts` currently calls +`client.voice.joinChannel(channel)` LAZILY — only when a `voiceStateUpdate` +shows `newState.streaming === true`. + +At that moment the bot is ALREADY connected to the channel via @discordjs/voice. +A selfbot `joinChannel` then does `VoiceConnection.authenticate()` → +`sendVoiceStateUpdate()`, and waits for a fresh `VOICE_SERVER_UPDATE` +(`setTokenAndEndpoint`) + `VOICE_STATE_UPDATE` (`setSessionId`) to reach +`checkAuthenticated()` (needs token+endpoint+sessionId). Because the bot is +already in an established voice session, Discord does NOT emit a new +`VOICE_SERVER_UPDATE` for the lazy selfbot re-join → token/endpoint never set → +15s `VOICE_CONNECTION_TIMEOUT`. + +This is fatal to video: `joinStreamConnection(userId)` (STREAM_WATCH op 20) and +`receiver.createVideoStream(userId, out)` (Recorder/ffmpeg) BOTH live on the +parent selfbot `VoiceConnection` and require it `CONNECTED` (its own voice +WS+UDP socket feeds `PacketHandler.push`, authenticated with +`authentication.secret_key`). + +## Fix — establish the selfbot connection eagerly, at voice-join time +The selfbot `VoiceConnection` must exist and be `CONNECTED` before any streamer +appears. Establish it once, synchronously alongside the @discordjs/voice join in +`recorder.startRecording`, so it rides the bot's FRESH voice join — when Discord +DOES emit VOICE_SERVER_UPDATE. Then cache it and let `videoRecorder` reuse it. + +Ordering: in `startRecording`, after the @discordjs/voice `joinVoiceChannel` +returns (and retries) — fire `ensureSelfbotVoice(channel)` best-effort: +1. `await client.voice.joinChannel(channel, { selfMute:false, selfDeaf:false, + selfVideo:false })` (rejects ~VOICE_CONNECTION_TIMEOUT on failure → log + + return null; do NOT block audio). +2. Cache the returned selfbot `VoiceConnection` keyed by guildId. +3. Wire teardown: on `recorder` voice stop / destroyed → `untrackChannel` + + destroy the cached selfbot connection (`disconnect()`). + +`videoRecorder.startVideoRecording` then uses the cached selfbot connection: +- If cached & `status === CONNECTED` → use it. +- Else → fall back to a lazy `joinChannel` (still best-effort). + +## The two-connection coexistence risk (must verify live) +@discordjs/voice (audio) and the selfbot `VoiceConnection` (video) each open +their OWN low-level voice WS+UDP on the same session. The spec's original +open-question flagged this. Mitigations: +- Clear logging: `Selfbot voice connected (guild=...)`, plus a periodic + `djs/voice status` log so we can confirm audio stays `READY` while the selfbot + connection is up. +- If Discord kicks/breaks the audio connection, logs will show + @discordjs/voice `Disconnected`/reconnect churn — we detect and pivot. + +## Files touched +- `services/discord-gateway/src/modules/voice-recording/videoRecorder.ts`: + add `ensureSelfbotVoice(channel)` (return cached/connected), use it in + `startVideoRecording`, add `destroyGuildSelfbotVoice(guildId)`, + richer status logging. +- `services/discord-gateway/src/modules/voice-recording/recorder.ts`: call + `ensureSelfbotVoice(channel)` after `joinVoiceChannel` (best-effort); + call `destroyGuildSelfbotVoice` on voice stop/destroy. +- Tests: `tests/videoRecorder.test.ts` (update to assert eager-connection reuse + + status gating). + +## Verification +1. `pnpm typecheck` + `pnpm build` + biome clean (discord-gateway). +2. Tests green. +3. Commit + push; CI `Build & Deploy (Nix)` green, service restarts. +4. LIVE (deploy): join a channel with the bot → journal shows + `Selfbot voice connected` (parent CONNECTED). When a member streams → + `Sender signal screenshare` / `Video recorder ready` + a `.mkv` under + `//video-*.mkv`; playable via ffmpeg. Confirm audio + recording still flows (no djs/voice reconnect churn). diff --git a/services/backend/src/shared/config/index.ts b/services/backend/src/shared/config/index.ts index ee4f3c89..5254f7f8 100644 --- a/services/backend/src/shared/config/index.ts +++ b/services/backend/src/shared/config/index.ts @@ -193,6 +193,7 @@ export const configSchema = z .positive() .default(50), PISCINA_MAX_THREADS: z.coerce.number().int().positive().optional(), + PISCINA_MEDIA_MAX_THREADS: z.coerce.number().int().positive().optional(), // ── Voice Transcription ──────────────────────────────────────────────── AI_VOICE_TRANSCRIPTION_ENABLED: z diff --git a/services/discord-gateway/ARCHITECTURE.md b/services/discord-gateway/ARCHITECTURE.md index 3d3b5c79..22d4a51a 100644 --- a/services/discord-gateway/ARCHITECTURE.md +++ b/services/discord-gateway/ARCHITECTURE.md @@ -77,20 +77,26 @@ handles a whole batch (text + media split internally, parallel paths). - Main thread owns the LLM semaphore (`AI_LLM_MAX_CONCURRENT`, default 5) via `llmClient.withLlmConcurrency`. -- Piscina pool (`PISCINA_MAX_THREADS`, default 4) runs the heavy LLM work off - the event loop; **each worker thread initializes its own pg Pool** (min 0, - grows to `POSTGRES_POOL_MAX`). See "Memory & connections" below. +- Two Piscina pools run the heavy LLM work off the event loop: a text pool + (`PISCINA_MAX_THREADS`, default 4) and a dedicated media pool + (`PISCINA_MEDIA_MAX_THREADS`, default 2). A batch is routed to the media + pool if ANY of its messages carries an attachment/sticker/embed — this + keeps a slow image/vision batch from occupying every thread and blocking + unrelated text-only batches behind it. **Each worker thread (in either + pool) initializes its own pg Pool** (min 0, grows to `POSTGRES_POOL_MAX`). + See "Memory & connections" below. ## Memory & DB connections `MemoryMax=1G` (raised from 512M — live RSS sits at ~500 MiB, peak 508 MiB, so 512M left ~2% headroom and risked an OOM-kill restart). Host has 8 GB free. -`POSTGRES_POOL_MIN=0` (default). The gateway = main process + up to 4 Piscina -worker threads, each with its own pg Pool. With min:0 the pools stay empty -until a query runs and drop idle clients afterward, instead of holding -`(1 main + 4 workers) × 2 = 10` permanently-open idle connections against -PgBouncer. The pool still grows on demand up to `POSTGRES_POOL_MAX`. +`POSTGRES_POOL_MIN=0` (default). The gateway = main process + up to 4 text +Piscina worker threads + up to 2 media Piscina worker threads, each with its +own pg Pool. With min:0 the pools stay empty until a query runs and drop +idle clients afterward, instead of holding `(1 main + 4 text + 2 media) × 2 += 14` permanently-open idle connections against PgBouncer. The pool still +grows on demand up to `POSTGRES_POOL_MAX`. ## Event channels (Redis pub/sub) diff --git a/services/discord-gateway/src/app/bootstrap.ts b/services/discord-gateway/src/app/bootstrap.ts index 55d8fced..b49e8f57 100644 --- a/services/discord-gateway/src/app/bootstrap.ts +++ b/services/discord-gateway/src/app/bootstrap.ts @@ -7,7 +7,10 @@ import { getAnalysisQueueStatus, startPendingAIAnalysisWorker, } from "../modules/ai-moderation/aiAnalyzer.js"; -import { workerPool } from "../modules/ai-moderation/circuitBreaker.js"; +import { + mediaWorkerPool, + textWorkerPool, +} from "../modules/ai-moderation/circuitBreaker.js"; import { registerChannelTopicCapture } from "../modules/channel-topic/index.js"; import { CommandHandler } from "../modules/command-handler/commandHandler.js"; import { @@ -396,12 +399,25 @@ export async function initializeDiscordGateway() { if (typeof status.lastError === "string") { setGauge("ai_analysis_last_error_present", status.lastError ? 1 : 0); } - const pool = workerPool as unknown as { - _poolState?: { size: number; active: number }; - }; - if (pool._poolState) { - setGauge("ai_analysis_worker_threads", pool._poolState.size); - setGauge("ai_analysis_worker_threads_active", pool._poolState.active); + 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"); diff --git a/services/discord-gateway/src/modules/ai-moderation/batchProcessor.ts b/services/discord-gateway/src/modules/ai-moderation/batchProcessor.ts index 1feabd8f..50ca0204 100644 --- a/services/discord-gateway/src/modules/ai-moderation/batchProcessor.ts +++ b/services/discord-gateway/src/modules/ai-moderation/batchProcessor.ts @@ -5,7 +5,7 @@ import { messageStore } from "../message-capture/messageStore.js"; import type { MessageRecord } from "../message-capture/types.js"; import { pickBatchWithinBudget as pickBatchWithinBudgetPure } from "./batchBudget.js"; import { partitionBatchOutcome } from "./batchOutcomeClassifier.js"; -import { workerPool } from "./circuitBreaker.js"; +import { mediaWorkerPool, textWorkerPool } from "./circuitBreaker.js"; import { estimateTokens } from "./conversationContext.js"; import { conversationErrorCooldown, @@ -14,6 +14,7 @@ import { resetConversationBatchFailures, } from "./conversationState.js"; import { enqueueIndividualFallbacks } from "./individualFallbackProcessor.js"; +import { hasMediaContent } from "./mediaAnalysisClient.js"; import { broadcastAnalysisCompleted, LAST_ERROR, @@ -121,29 +122,31 @@ export async function skipAgeRestrictedMessages( // Batch pipeline // --------------------------------------------------------------------------- -export async function processBatch( +/** + * Runs one worker job (either the text-only or the media sub-batch of a + * conversation) end-to-end: dispatch → broadcast/save → fallback routing. + * Returns whether the *caller* should schedule the next debounce pass for + * this conversation (mirrors the old single-job semantics, now evaluated + * per queue). + * + * Broadcasting happens here, inside each queue's own call — NOT after + * waiting on the other queue. That's the actual fix for "text menunggu + * image": previously one mixed conversation batch made ONE worker call + * with both text and media targets, and `runModerationAnalysis` only + * resolves (so results only get saved/broadcast) once BOTH finish — so a + * fast text verdict sat unused until the slow vision/image verdict was + * ready too. Splitting into two independent jobs means the text queue + * saves+broadcasts its rows the moment IT finishes, regardless of how long + * the media queue takes. + */ +async function runQueueBatch( + pool: typeof textWorkerPool, conversationKey: string, messages: MessageRecord[], - processingStartedAt: number, -): Promise { - if (messages.length === 0) { - if (conversationProcessing.get(conversationKey) === processingStartedAt) { - conversationProcessing.delete(conversationKey); - } - return; - } - const cooldownUntil = conversationErrorCooldown.get(conversationKey) ?? 0; - if (Date.now() < cooldownUntil) { - if (conversationProcessing.get(conversationKey) === processingStartedAt) { - conversationProcessing.delete(conversationKey); - } - return; - } - +): Promise { activeRequests++; - let shouldScheduleNext = false; try { - const result = (await workerPool.run({ + const result = (await pool.run({ type: "batch", conversationKey, messages, @@ -191,7 +194,7 @@ export async function processBatch( }, "Batch analysis failed, will retry after cooldown", ); - return; + return false; } // Batch succeeded -- partition per-message outcome (2026-08-25). @@ -281,20 +284,13 @@ export async function processBatch( conversationErrorCooldown.set(conversationKey, newCooldown); } - // Release the processing lock immediately so the cooldown timer controls retry - if (conversationProcessing.get(conversationKey) === processingStartedAt) { - conversationProcessing.delete(conversationKey); - } - // Do NOT schedule next -- let the cooldown gate it - shouldScheduleNext = false; + return false; } - if (apiFailedMessages.length === 0) { - resetConversationBatchFailures(conversationKey); - conversationErrorCooldown.delete(conversationKey); - } - shouldScheduleNext = true; + resetConversationBatchFailures(conversationKey); + conversationErrorCooldown.delete(conversationKey); + return true; } catch (error) { recordConversationBatchFailure(conversationKey); @@ -326,18 +322,63 @@ export async function processBatch( }, "Analysis worker failed, will retry after cooldown", ); + return false; } finally { activeRequests--; + } +} + +export async function processBatch( + conversationKey: string, + messages: MessageRecord[], + processingStartedAt: number, +): Promise { + if (messages.length === 0) { if (conversationProcessing.get(conversationKey) === processingStartedAt) { conversationProcessing.delete(conversationKey); } - if (shouldScheduleNext) { - setImmediate(() => { - // Dynamic import to avoid circular dependency at module scope - import("./batchScheduler.js").then((m) => - m.scheduleConversationAnalysis(conversationKey), - ); - }); + return; + } + const cooldownUntil = conversationErrorCooldown.get(conversationKey) ?? 0; + if (Date.now() < cooldownUntil) { + if (conversationProcessing.get(conversationKey) === processingStartedAt) { + conversationProcessing.delete(conversationKey); } + return; + } + + // Split the batch itself — not just route it — so text and media never + // share one worker call. A conversation batch commonly mixes plain-text + // messages with an image/sticker from someone else; without this split, + // ALL of it (including the plain-text messages) would ride along on the + // media job and wait for vision analysis to finish. Each sub-batch is now + // dispatched to its own pool AND handled independently below, so the text + // queue's results land as soon as text analysis completes, full stop. + const textMessages = messages.filter((m) => !hasMediaContent(m)); + const mediaMessages = messages.filter((m) => hasMediaContent(m)); + + const jobs: Promise[] = []; + if (textMessages.length > 0) { + jobs.push(runQueueBatch(textWorkerPool, conversationKey, textMessages)); + } + if (mediaMessages.length > 0) { + jobs.push(runQueueBatch(mediaWorkerPool, conversationKey, mediaMessages)); + } + + const outcomes = await Promise.allSettled(jobs); + const shouldScheduleNext = outcomes.every( + (o) => o.status === "fulfilled" && o.value, + ); + + if (conversationProcessing.get(conversationKey) === processingStartedAt) { + conversationProcessing.delete(conversationKey); + } + if (shouldScheduleNext) { + setImmediate(() => { + // Dynamic import to avoid circular dependency at module scope + import("./batchScheduler.js").then((m) => + m.scheduleConversationAnalysis(conversationKey), + ); + }); } } diff --git a/services/discord-gateway/src/modules/ai-moderation/circuitBreaker.ts b/services/discord-gateway/src/modules/ai-moderation/circuitBreaker.ts index 7fe1ec12..4a2dc504 100644 --- a/services/discord-gateway/src/modules/ai-moderation/circuitBreaker.ts +++ b/services/discord-gateway/src/modules/ai-moderation/circuitBreaker.ts @@ -6,7 +6,16 @@ import { config } from "../../shared/config/config.js"; import type { MessageRecord } from "../message-capture/types.js"; // --------------------------------------------------------------------------- -// Piscina worker pool (shared by batch + individual pipelines) +// Piscina worker pools (2026-08-31: split text vs media) +// +// Both the batch and individual-fallback pipelines used to share ONE pool. +// A conversation batch containing images/stickers/embeds routes through +// runMediaBatch (download + vision + LLM — tens of seconds per batch), and +// Piscina hands each worker thread exactly one task at a time. With a small +// fixed thread count, a handful of slow media batches could occupy every +// thread and leave fast text-only batches for OTHER conversations queued +// behind them for the whole media duration. Text and media now get their +// own dedicated pools so a media backlog can never starve text analysis. // --------------------------------------------------------------------------- function getAnalysisWorkerUrl(): URL { @@ -25,12 +34,32 @@ function getAnalysisWorkerUrl(): URL { return candidates[2]; } -export const workerPool = new Piscina({ - filename: fileURLToPath(getAnalysisWorkerUrl()), +const analysisWorkerFilename = fileURLToPath(getAnalysisWorkerUrl()); + +/** Dedicated pool for text-only batch/individual analysis jobs. */ +export const textWorkerPool = new Piscina({ + filename: analysisWorkerFilename, execArgv: process.execArgv, maxThreads: config.PISCINA_MAX_THREADS ?? availableParallelism(), }); +/** + * Dedicated pool for jobs whose batch contains at least one message with + * media (attachments/stickers/embeds). Kept separate so slow image/vision + * analysis never blocks the text pool above. + */ +export const mediaWorkerPool = new Piscina({ + filename: analysisWorkerFilename, + execArgv: process.execArgv, + maxThreads: config.PISCINA_MEDIA_MAX_THREADS ?? 2, +}); + +/** + * @deprecated Use `textWorkerPool` or `mediaWorkerPool` directly. Kept as an + * alias to the text pool only for anything not yet migrated. + */ +export const workerPool = textWorkerPool; + /** * Gets the conversation key for a message (thread_id or channel_id). */ diff --git a/services/discord-gateway/src/modules/ai-moderation/individualFallbackProcessor.ts b/services/discord-gateway/src/modules/ai-moderation/individualFallbackProcessor.ts index 8855c35a..4e6bc671 100644 --- a/services/discord-gateway/src/modules/ai-moderation/individualFallbackProcessor.ts +++ b/services/discord-gateway/src/modules/ai-moderation/individualFallbackProcessor.ts @@ -6,9 +6,14 @@ import type { AnalysisResult, MessageRecord, } from "../message-capture/types.js"; -import { getConversationKey, workerPool } from "./circuitBreaker.js"; +import { + getConversationKey, + mediaWorkerPool, + textWorkerPool, +} from "./circuitBreaker.js"; import { fireAlert } from "./conversationState.js"; import { classifyIndividualWorkerResult } from "./fallbackResultClassifier.js"; +import { hasMediaContent } from "./mediaAnalysisClient.js"; import { broadcastAnalysisCompleted, LAST_ERROR, @@ -73,8 +78,12 @@ async function processIndividualFallback( let exhaustedOnIncomplete = false; try { - // Run the LLM-heavy work in the worker thread - const workerResult = (await workerPool.run({ + // Run the LLM-heavy work in the worker thread — same text/media pool + // split as the batch pipeline (see circuitBreaker.ts), so a single + // media message falling back to individual analysis can't queue behind + // (or block) text-only fallbacks, and vice versa. + const pool = hasMediaContent(message) ? mediaWorkerPool : textWorkerPool; + const workerResult = (await pool.run({ type: "individual", message, skipNormalAnalysis: false, diff --git a/services/discord-gateway/src/shared/config/index.ts b/services/discord-gateway/src/shared/config/index.ts index 9c9fde12..0e92e357 100644 --- a/services/discord-gateway/src/shared/config/index.ts +++ b/services/discord-gateway/src/shared/config/index.ts @@ -319,11 +319,20 @@ export const configSchema = z .int() .positive() .default(50), - // Worker pool size. Default 4 (not availableParallelism) because each - // Piscina thread owns its own pLimit(5) semaphore — on big VPSes - // availableParallelism × 5 concurrent LLM calls would overwhelm the - // router. Keep threads modest; concurrency is capped per-thread anyway. + // Text-analysis worker pool size (2026-08-31: split from the media pool + // below so a slow image/vision batch can never occupy every thread and + // starve the far more common text-only batches). Default 4 (not + // availableParallelism) because each Piscina thread owns its own + // pLimit(5) semaphore — on big VPSes availableParallelism × 5 concurrent + // LLM calls would overwhelm the router. Keep threads modest; concurrency + // is capped per-thread anyway. PISCINA_MAX_THREADS: z.coerce.number().int().positive().default(4), + // Media-analysis worker pool size — dedicated threads for batches that + // contain images/stickers/embeds (download + vision + LLM, much slower + // than text). Kept small since media batches are less frequent and each + // one is long-running; sized independently from PISCINA_MAX_THREADS so + // tuning one never starves the other. + PISCINA_MEDIA_MAX_THREADS: z.coerce.number().int().positive().default(2), // ── Voice Transcription ──────────────────────────────────────────────── AI_VOICE_TRANSCRIPTION_ENABLED: z