diff --git a/services/backend/src/modules/dashboard/dashboard.repository.ts b/services/backend/src/modules/dashboard/dashboard.repository.ts index 6fabdc9..b7aef81 100644 --- a/services/backend/src/modules/dashboard/dashboard.repository.ts +++ b/services/backend/src/modules/dashboard/dashboard.repository.ts @@ -210,24 +210,24 @@ export class DashboardRepository { [...params, limit + 1], ); - const data = ((rows as Record[]) || []).slice(0, limit).map((r) => ({ - channel_id: String(r.channel_id), - channel_name: r.channel_name as string | null, - guild_id: r.guild_id as string | null, - total_messages: Number(r.total_messages), - flagged_count: Number(r.flagged_count), - last_message_at: r.last_message_at ? Number(r.last_message_at) : null, - culture_summary: r.culture_summary as string | null, - last_analyzed_at: r.last_analyzed_at - ? Number(r.last_analyzed_at) - : null, - })); + const data = ((rows as Record[]) || []) + .slice(0, limit) + .map((r) => ({ + channel_id: String(r.channel_id), + channel_name: r.channel_name as string | null, + guild_id: r.guild_id as string | null, + total_messages: Number(r.total_messages), + flagged_count: Number(r.flagged_count), + last_message_at: r.last_message_at ? Number(r.last_message_at) : null, + culture_summary: r.culture_summary as string | null, + last_analyzed_at: r.last_analyzed_at + ? Number(r.last_analyzed_at) + : null, + })); const lastRow = rows[limit - 1] as Record | undefined; const nextCursor = - rows.length > limit - ? String(lastRow?.total_messages ?? "") - : null; + rows.length > limit ? String(lastRow?.total_messages ?? "") : null; return { data, nextCursor }; } diff --git a/services/backend/src/modules/dashboard/dashboard.routes.ts b/services/backend/src/modules/dashboard/dashboard.routes.ts index 3030ab3..12b6c59 100644 --- a/services/backend/src/modules/dashboard/dashboard.routes.ts +++ b/services/backend/src/modules/dashboard/dashboard.routes.ts @@ -56,9 +56,7 @@ export function createDashboardRouter(): Router { const search = typeof req.query.search === "string" ? req.query.search : undefined; const guildId = - typeof req.query.guild_id === "string" - ? req.query.guild_id - : undefined; + typeof req.query.guild_id === "string" ? req.query.guild_id : undefined; const result = await dashboardService.listChannels({ limit, diff --git a/services/backend/src/modules/dashboard/dashboard.service.ts b/services/backend/src/modules/dashboard/dashboard.service.ts index 533afdf..5987098 100644 --- a/services/backend/src/modules/dashboard/dashboard.service.ts +++ b/services/backend/src/modules/dashboard/dashboard.service.ts @@ -25,7 +25,11 @@ export class DashboardService { return dashboardRepository.getUserDetail(userId); } - async listChannels(query: { limit: number; search?: string; guildId?: string }) { + async listChannels(query: { + limit: number; + search?: string; + guildId?: string; + }) { logger.debug({ query }, "Listing dashboard channels"); return dashboardRepository.listChannels(query); } diff --git a/services/backend/src/modules/recordings/recordings.routes.ts b/services/backend/src/modules/recordings/recordings.routes.ts index 382d79f..e06211b 100644 --- a/services/backend/src/modules/recordings/recordings.routes.ts +++ b/services/backend/src/modules/recordings/recordings.routes.ts @@ -14,8 +14,13 @@ export function createRecordingsRouter(): Router { "/recordings", asyncHandler(async (req: Request, res: Response) => { const limit = Number(req.query.limit) || 50; - logger.debug({ limit }, "Fetching recordings"); - const result = await recordingsService.getRecent(limit); + const channelId = req.query.channelId as string | undefined; + const userId = req.query.userId as string | undefined; + logger.debug({ limit, channelId, userId }, "Fetching recordings"); + const result = await recordingsService.getRecent(limit, { + channelId, + userId, + }); res.json(result); }), ); diff --git a/services/backend/src/modules/recordings/recordings.service.ts b/services/backend/src/modules/recordings/recordings.service.ts index 91d377b..27e6b07 100644 --- a/services/backend/src/modules/recordings/recordings.service.ts +++ b/services/backend/src/modules/recordings/recordings.service.ts @@ -5,17 +5,37 @@ import { getDatabase } from "../../shared/database/index.js"; const logger = createChildLogger("recordings.service"); export class RecordingsService { - async getRecent(limit = 50) { + async getRecent( + limit = 50, + filters?: { channelId?: string; userId?: string }, + ) { logger.info({ limit }, "getRecent called"); const db = getDatabase(); logger.debug({ limit }, "Fetching recent voice recordings"); + const conditions: string[] = []; + const params: unknown[] = []; + + if (filters?.channelId) { + params.push(filters.channelId); + conditions.push(`channel_id = $${params.length}`); + } + if (filters?.userId) { + params.push(filters.userId); + conditions.push(`user_id = $${params.length}`); + } + + const whereClause = + conditions.length > 0 ? `WHERE ${conditions.join(" AND ")}` : ""; + const { rows } = await db.execute(sql` SELECT id, user_id, username, avatar_url, guild_id, channel_id, channel_name, filename, size_bytes, download_url, - upload_status, upload_error, created_at, uploaded_at + upload_status, upload_error, created_at, uploaded_at, + COALESCE(size_bytes, 0) AS duration_bytes FROM voice_recordings + ${sql.raw(whereClause)} ORDER BY created_at DESC LIMIT ${limit} `); diff --git a/services/backend/src/modules/voice/voice.service.ts b/services/backend/src/modules/voice/voice.service.ts index e32e1ec..3f46c63 100644 --- a/services/backend/src/modules/voice/voice.service.ts +++ b/services/backend/src/modules/voice/voice.service.ts @@ -4,6 +4,7 @@ import { COMMAND_VOICE_CHANNELS, COMMAND_VOICE_CONNECT, COMMAND_VOICE_DISCONNECT, + COMMAND_VOICE_DISCONNECT_GUILD, CommandReply, VOICE_STATUS_KEY, } from "@bete/shared"; @@ -28,11 +29,19 @@ export interface Channel { type: "voice" | "text"; } +export interface GuildVoiceEntry { + guildId: string; + channelId: string; + channelName: string; + connectedAt: number; +} + export interface VoiceStatus { connected: boolean; activeGuildId: string | null; activeChannelId: string | null; activeChannelName: string | null; + connections: GuildVoiceEntry[]; } export const DEFAULT_VOICE_STATUS: VoiceStatus = { @@ -40,6 +49,7 @@ export const DEFAULT_VOICE_STATUS: VoiceStatus = { activeGuildId: null, activeChannelId: null, activeChannelName: null, + connections: [], }; /** @@ -157,3 +167,18 @@ export async function disconnectVoice(): Promise { "disconnectVoice", ); } + +/** + * Disconnect from a specific guild's voice channel. + */ +export async function disconnectVoiceGuild( + guildId: string, +): Promise { + logger.info({ guildId }, "disconnectVoiceGuild called"); + return withFallback( + () => + publishCommand(COMMAND_VOICE_DISCONNECT_GUILD, { guildId }), + () => readVoiceStatusFallback(), + "disconnectVoiceGuild", + ); +} diff --git a/services/backend/src/ws/redis-bridge.ts b/services/backend/src/ws/redis-bridge.ts index 1d2e761..57c55be 100644 --- a/services/backend/src/ws/redis-bridge.ts +++ b/services/backend/src/ws/redis-bridge.ts @@ -2,10 +2,19 @@ import { DISCORD_ANALYSIS_QUEUE_STATUS, DISCORD_ATTACHMENT_CREATED, DISCORD_ATTACHMENT_UPLOADED, + DISCORD_CHANNEL_TOPIC_UPDATED, + DISCORD_GUILD_MEMBER_ADDED, + DISCORD_GUILD_MEMBER_REMOVED, DISCORD_MESSAGE_ANALYZED, DISCORD_MESSAGE_CREATED, DISCORD_MESSAGE_DELETED, DISCORD_MESSAGE_UPDATED, + DISCORD_PRESENCE_UPDATED, + DISCORD_REACTION_ADDED, + DISCORD_REACTION_REMOVED, + DISCORD_THREAD_CREATED, + DISCORD_THREAD_DELETED, + DISCORD_THREAD_UPDATED, DISCORD_VOICE_ACTIVE_USER, DISCORD_VOICE_PCM, DISCORD_VOICE_STARTED, @@ -40,6 +49,18 @@ const SUBSCRIPTIONS: ChannelMapping[] = [ }, { channel: DISCORD_VOICE_ACTIVE_USER, eventType: "voice_active_user" }, { channel: DISCORD_VOICE_PCM, eventType: "voice_pcm_data" }, + { channel: DISCORD_REACTION_ADDED, eventType: "reaction_added" }, + { channel: DISCORD_REACTION_REMOVED, eventType: "reaction_removed" }, + { channel: DISCORD_THREAD_CREATED, eventType: "thread_created" }, + { channel: DISCORD_THREAD_DELETED, eventType: "thread_deleted" }, + { channel: DISCORD_THREAD_UPDATED, eventType: "thread_updated" }, + { + channel: DISCORD_CHANNEL_TOPIC_UPDATED, + eventType: "channel_topic_updated", + }, + { channel: DISCORD_PRESENCE_UPDATED, eventType: "presence_updated" }, + { channel: DISCORD_GUILD_MEMBER_ADDED, eventType: "guild_member_added" }, + { channel: DISCORD_GUILD_MEMBER_REMOVED, eventType: "guild_member_removed" }, ]; let subscriber: Redis | null = null; diff --git a/services/discord-gateway/src/modules/command-handler/commandHandler.ts b/services/discord-gateway/src/modules/command-handler/commandHandler.ts index 4c45b7a..2377293 100644 --- a/services/discord-gateway/src/modules/command-handler/commandHandler.ts +++ b/services/discord-gateway/src/modules/command-handler/commandHandler.ts @@ -9,7 +9,10 @@ import { createChildLogger } from "@bete/shared/logger"; import type { Client } from "discord.js-selfbot-v13"; import Redis from "ioredis"; import { config } from "../../shared/config/config.js"; -import type { VoiceController } from "../voice-recording/voiceController.js"; +import type { + VoiceController, + VoiceStatus, +} from "../voice-recording/voiceController.js"; import { GuildHandler } from "./guild.handler.js"; import { type CommandHandlerFn, @@ -30,6 +33,12 @@ interface VoiceStatusPayload { activeGuildId: string | null; activeChannelId: string | null; activeChannelName: string | null; + connections: Array<{ + guildId: string; + channelId: string; + channelName: string; + connectedAt: number; + }>; } // --------------------------------------------------------------------------- @@ -69,7 +78,11 @@ export class CommandHandler { this.voiceController = voiceController; // Create domain-specific handlers with their dependencies - this.voiceHandler = new VoiceHandler(client, voiceController, this.redisPub); + this.voiceHandler = new VoiceHandler( + client, + voiceController, + this.redisPub, + ); this.mediaHandler = new MediaHandler(); this.guildHandler = new GuildHandler(client); this.moderationHandler = new ModerationHandler(client); @@ -164,14 +177,23 @@ export class CommandHandler { // ---- Status publishing ---- private publishVoiceStatus(): void { - const status: VoiceStatusPayload = this.voiceController + const raw = this.voiceController ? this.voiceController.getStatus() : { + ready: false, connected: false, activeGuildId: null, activeChannelId: null, activeChannelName: null, + connections: [], }; + const status: VoiceStatusPayload = { + connected: raw.connected, + activeGuildId: raw.activeGuildId, + activeChannelId: raw.activeChannelId, + activeChannelName: raw.activeChannelName, + connections: raw.connections ?? [], + }; this.setKey(VOICE_STATUS_KEY, JSON.stringify(status)); } diff --git a/services/discord-gateway/src/modules/command-handler/handler-registry.ts b/services/discord-gateway/src/modules/command-handler/handler-registry.ts index 6d20cb6..52fc1e0 100644 --- a/services/discord-gateway/src/modules/command-handler/handler-registry.ts +++ b/services/discord-gateway/src/modules/command-handler/handler-registry.ts @@ -9,6 +9,7 @@ import { COMMAND_VOICE_CHANNELS, COMMAND_VOICE_CONNECT, COMMAND_VOICE_DISCONNECT, + COMMAND_VOICE_DISCONNECT_GUILD, COMMAND_VOICE_TRANSMIT_START, COMMAND_VOICE_TRANSMIT_STOP, type CommandMessage, @@ -46,6 +47,9 @@ export function createHandlerRegistry( registry.set(COMMAND_VOICE_DISCONNECT, (cmd) => voiceHandler.handleVoiceDisconnect(cmd), ); + registry.set(COMMAND_VOICE_DISCONNECT_GUILD, (cmd) => + voiceHandler.handleVoiceDisconnectGuild(cmd), + ); registry.set(COMMAND_VOICE_CHANNELS, (cmd) => voiceHandler.handleVoiceChannels(cmd), ); diff --git a/services/discord-gateway/src/modules/command-handler/voice.handler.ts b/services/discord-gateway/src/modules/command-handler/voice.handler.ts index b586864..fc3693a 100644 --- a/services/discord-gateway/src/modules/command-handler/voice.handler.ts +++ b/services/discord-gateway/src/modules/command-handler/voice.handler.ts @@ -1,4 +1,8 @@ -import { type CommandMessage, type CommandReply } from "@bete/shared"; +import { + COMMAND_VOICE_DISCONNECT_GUILD, + type CommandMessage, + type CommandReply, +} from "@bete/shared"; import { createChildLogger } from "@bete/shared/logger"; import type { Client } from "discord.js-selfbot-v13"; import type Redis from "ioredis"; @@ -72,6 +76,33 @@ export class VoiceHandler { return { id: cmd.id, success: true, data: status }; } + async handleVoiceDisconnectGuild( + cmd: CommandMessage, + ): Promise> { + if (!this.voiceController) { + return { + id: cmd.id, + success: false, + data: null, + error: "Gateway not initialized", + }; + } + + const guildId = String(cmd.payload.guildId ?? ""); + if (!guildId) { + return { + id: cmd.id, + success: false, + data: null, + error: "guildId is required", + }; + } + + await this.voiceController.disconnectGuild(guildId); + const status = this.voiceController.getStatus(); + return { id: cmd.id, success: true, data: status }; + } + async handleVoiceChannels( cmd: CommandMessage, ): Promise> { @@ -126,7 +157,9 @@ export class VoiceHandler { try { // Reuse shared Redis connection from CommandHandler - const transmitRedis = this.sharedRedis ?? new (await import("ioredis")).default(config.REDIS_URL); + const transmitRedis = + this.sharedRedis ?? + new (await import("ioredis")).default(config.REDIS_URL); await voiceTransmitter.start(transmitRedis); const status = voiceTransmitter.getStatus(); diff --git a/services/discord-gateway/src/modules/voice-recording/mediaSource.ts b/services/discord-gateway/src/modules/voice-recording/mediaSource.ts index 4bd2ced..e3969ce 100644 --- a/services/discord-gateway/src/modules/voice-recording/mediaSource.ts +++ b/services/discord-gateway/src/modules/voice-recording/mediaSource.ts @@ -307,7 +307,9 @@ export async function extractMediaInfo(url: string): Promise { if (proc.stderr) { proc.stderr.on("data", (chunk: Buffer) => { if (stderrBuf.length < MAX_STDERR) { - stderrBuf += chunk.toString("utf8").slice(0, MAX_STDERR - stderrBuf.length); + stderrBuf += chunk + .toString("utf8") + .slice(0, MAX_STDERR - stderrBuf.length); } }); } diff --git a/services/discord-gateway/src/modules/voice-recording/recorder.ts b/services/discord-gateway/src/modules/voice-recording/recorder.ts index 90e87ad..83c1d54 100644 --- a/services/discord-gateway/src/modules/voice-recording/recorder.ts +++ b/services/discord-gateway/src/modules/voice-recording/recorder.ts @@ -194,29 +194,33 @@ export function stopRecording(guildId: string): void { const snapshot = session.snapshot(Date.now()); const stoppedAt = Date.now(); - _eventBroadcaster.voiceRecordingStopped({ - guild_id: guildId, - session_id: session.sessionId, - duration_ms: snapshot.durationMs, - participants: snapshot.participants.length, - segment_count: snapshot.segments.length, - status: snapshot.status, - stopped_at: stoppedAt, - }).catch(() => {}); + _eventBroadcaster + .voiceRecordingStopped({ + guild_id: guildId, + session_id: session.sessionId, + duration_ms: snapshot.durationMs, + participants: snapshot.participants.length, + segment_count: snapshot.segments.length, + status: snapshot.status, + stopped_at: stoppedAt, + }) + .catch(() => {}); // Auto-enqueue muxer job if there are multiple segments const segments = snapshot.segments; if (segments.length >= 2) { const outputFile = `${config.RECORDINGS_DIR}/merged/${session.sessionId}.ogg`; - import("./muxer.js").then(({ enqueueMuxerJob }) => { - enqueueMuxerJob({ - inputs: segments.map((s) => s.oggPath), - output: outputFile, - guildId, - channelId: snapshot.channelId, - sessionId: session.sessionId, - }).catch(() => {}); - }).catch(() => {}); + import("./muxer.js") + .then(({ enqueueMuxerJob }) => { + enqueueMuxerJob({ + inputs: segments.map((s) => s.oggPath), + output: outputFile, + guildId, + channelId: snapshot.channelId, + sessionId: session.sessionId, + }).catch(() => {}); + }) + .catch(() => {}); } } } diff --git a/services/discord-gateway/src/modules/voice-recording/voiceController.ts b/services/discord-gateway/src/modules/voice-recording/voiceController.ts index 441ca8b..5b74ed4 100644 --- a/services/discord-gateway/src/modules/voice-recording/voiceController.ts +++ b/services/discord-gateway/src/modules/voice-recording/voiceController.ts @@ -7,34 +7,49 @@ import { startRecording, stopRecording } from "./recorder.js"; const logger = createChildLogger("voice-controller"); +// ─── Types ─────────────────────────────────────────────────────────────── + +export interface GuildVoiceState { + guildId: string; + channelId: string; + channelName: string; + connectedAt: number; +} + export interface VoiceStatus { ready: boolean; connected: boolean; activeGuildId: string | null; activeChannelId: string | null; activeChannelName: string | null; + /** Multi-guild: list of all active connections */ + connections: GuildVoiceState[]; } +// ─── VoiceController ───────────────────────────────────────────────────── + export class VoiceController { - private activeGuildId: string | null = null; - private activeChannelId: string | null = null; - private activeChannelName: string | null = null; - private connecting = false; + private connections = new Map(); + private connecting = new Set(); constructor(private readonly client: Client) {} getStatus(): VoiceStatus { logger.debug("getStatus called"); - const connection = this.activeGuildId - ? getVoiceConnection(this.activeGuildId) + + // Primary connection (legacy compat — first entry or explicitly set) + const primaryGuildId = this.connections.keys().next().value ?? null; + const primary = primaryGuildId + ? this.connections.get(primaryGuildId) : undefined; return { ready: this.client.isReady(), - connected: Boolean(connection), - activeGuildId: this.activeGuildId, - activeChannelId: this.activeChannelId, - activeChannelName: this.activeChannelName, + connected: this.connections.size > 0, + activeGuildId: primary?.guildId ?? null, + activeChannelId: primary?.channelId ?? null, + activeChannelName: primary?.channelName ?? null, + connections: Array.from(this.connections.values()), }; } @@ -48,18 +63,21 @@ export class VoiceController { ); } - if (this.connecting) { + if (this.connecting.has(guildId)) { throw new AppError( - "Voice connection is already in progress", + `Voice connection for guild ${guildId} is already in progress`, "CONNECT_IN_PROGRESS", 409, ); } - this.connecting = true; + this.connecting.add(guildId); try { - await this.disconnect(); + // Disconnect existing connection for this guild first + if (this.connections.has(guildId)) { + await this.disconnectGuild(guildId); + } const guild = this.getGuild(guildId); const channel = @@ -94,10 +112,18 @@ export class VoiceController { ); } - discordPlayer.setConnection(connection as VoiceConnection); - this.activeGuildId = guildId; - this.activeChannelId = channelId; - this.activeChannelName = channel.name; + // If this is the first connection, set it as the player's connection + if (this.connections.size === 0) { + discordPlayer.setConnection(connection as VoiceConnection); + } + + const state: GuildVoiceState = { + guildId, + channelId, + channelName: channel.name, + connectedAt: Date.now(), + }; + this.connections.set(guildId, state); logger.info( { guildId, channelId, channelName: channel.name }, @@ -106,24 +132,31 @@ export class VoiceController { return this.getStatus(); } finally { - this.connecting = false; + this.connecting.delete(guildId); } } async disconnect(): Promise { logger.info("disconnect called"); - if (this.activeGuildId) { - stopRecording(this.activeGuildId); + + // Disconnect all guilds + const guildIds = Array.from(this.connections.keys()); + for (const gid of guildIds) { + await this.disconnectGuild(gid); } discordPlayer.stop(); - this.activeGuildId = null; - this.activeChannelId = null; - this.activeChannelName = null; - return this.getStatus(); } + async disconnectGuild(guildId: string): Promise { + logger.info({ guildId }, "disconnectGuild called"); + if (this.connections.has(guildId)) { + stopRecording(guildId); + this.connections.delete(guildId); + } + } + private getGuild(guildId: string): Guild { const guild = this.client.guilds.cache.get(guildId); diff --git a/services/frontend/src/features/dashboard/components/ChannelProfileDetail.tsx b/services/frontend/src/features/dashboard/components/ChannelProfileDetail.tsx index cdbb6ed..c14886d 100644 --- a/services/frontend/src/features/dashboard/components/ChannelProfileDetail.tsx +++ b/services/frontend/src/features/dashboard/components/ChannelProfileDetail.tsx @@ -1,10 +1,7 @@ import { motion } from "framer-motion"; import { AlertCircle, ArrowLeft, Hash, RefreshCw } from "lucide-react"; import type { DashboardChannelDetail } from "../../../shared/api/client"; -import { - cardItem, - cardStagger, -} from "../../../shared/hooks/useFramerStagger"; +import { cardItem, cardStagger } from "../../../shared/hooks/useFramerStagger"; import { cn } from "../../../shared/lib/utils"; import { Badge, diff --git a/services/frontend/src/features/dashboard/components/ChannelSummaryList.tsx b/services/frontend/src/features/dashboard/components/ChannelSummaryList.tsx index 720634b..d74bc46 100644 --- a/services/frontend/src/features/dashboard/components/ChannelSummaryList.tsx +++ b/services/frontend/src/features/dashboard/components/ChannelSummaryList.tsx @@ -1,20 +1,9 @@ import { motion } from "framer-motion"; -import { - AlertCircle, - Hash, - Loader2, - RefreshCw, - Search, -} from "lucide-react"; +import { AlertCircle, Hash, Loader2, RefreshCw, Search } from "lucide-react"; import type { DashboardChannel } from "../../../shared/api/client"; import { cardItem, cardStagger } from "../../../shared/hooks/useFramerStagger"; import { cn } from "../../../shared/lib/utils"; -import { - Card, - CardContent, - Input, - Skeleton, -} from "../../../shared/ui"; +import { Card, CardContent, Input, Skeleton } from "../../../shared/ui"; interface ChannelSummaryListProps { channels: DashboardChannel[]; diff --git a/services/frontend/src/shared/api/client.ts b/services/frontend/src/shared/api/client.ts index 511ceae..005b949 100644 --- a/services/frontend/src/shared/api/client.ts +++ b/services/frontend/src/shared/api/client.ts @@ -281,7 +281,11 @@ export interface DashboardStats { today_messages: number; today_flagged: number; active_users_24h: number; - top_channels: Array<{ channel_id: string; channel_name: string | null; message_count: number }>; + top_channels: Array<{ + channel_id: string; + channel_name: string | null; + message_count: number; + }>; moderation_overview: { pending: number; processing: number; diff --git a/services/frontend/src/shared/ws/events.ts b/services/frontend/src/shared/ws/events.ts index 04547e6..102a2b3 100644 --- a/services/frontend/src/shared/ws/events.ts +++ b/services/frontend/src/shared/ws/events.ts @@ -16,6 +16,15 @@ export interface WsEventMap { voice_active_user: { data: unknown }; attachment_created: { data: unknown }; analysis_queue_status: { data: unknown }; + reaction_added: { data: unknown }; + reaction_removed: { data: unknown }; + thread_created: { data: unknown }; + thread_deleted: { data: unknown }; + thread_updated: { data: unknown }; + channel_topic_updated: { data: unknown }; + presence_updated: { data: unknown }; + guild_member_added: { data: unknown }; + guild_member_removed: { data: unknown }; } export type WsEventType = keyof WsEventMap;