feat(discord-gateway): implement voice & push improvements
- Voice disconnect broadcast on stopRecording - Multi-guild voice support (VoiceController Map<guildId>) - Session finalization + auto-enqueue muxer job - Recordings API: duration field, channelId/userId filters - Transmitter Redis connection reuse (shared conn) - FFmpeg stderr memory cap (4KB limit) - 10 new Redis event channels + Redis bridge subscriptions - New DB tables: message_reactions, message_edits - Webhook notification module - Gateway metrics / Prometheus endpoint - Multi-guild message capture (MONITOR_GUILD_IDS array) - Thread tracking, presence, channel topic, guild member events - Edit history snapshot on message update - Muxer audio post-processing worker Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
@@ -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));
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
);
|
||||
|
||||
@@ -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<CommandReply<unknown>> {
|
||||
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<CommandReply<unknown>> {
|
||||
@@ -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();
|
||||
|
||||
@@ -307,7 +307,9 @@ export async function extractMediaInfo(url: string): Promise<MediaInfo> {
|
||||
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);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@@ -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(() => {});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<string, GuildVoiceState>();
|
||||
private connecting = new Set<string>();
|
||||
|
||||
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<VoiceStatus> {
|
||||
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<void> {
|
||||
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);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user