refactor: remove voice, recording, and music/media features (gateway + backend)
- Gateway: delete voice-recording/, voice-pcm-ws/, voice/video/media handlers, vendor/discord-voice-fork/, voice DB repos, 2 voice migrations - Gateway: cut voice wiring from bootstrap/shutdown/commandHandler/handler-registry, eventBroadcaster/eventTypes, redis-channels, moderation-types, message-capture, config keys, and deps (@discordjs/voice, opus, libsodium, prism-media, davey) - Backend: delete voice/, recordings/, media/ modules + @discordjs/voice dep - Backend: cut voice/media/recordings oRPC routers, WS gateway-PCM auth + voice handlers, Redis bridge voice aggregation, redis-channels voice/media constants, config keys, moderation-types voice items, e2e voice/recordings suites - Keep voice DB tables (destructive migration avoided); chatbot voiceRecordings tool + dashboard count remain as read-only historical data access
This commit is contained in:
@@ -41,20 +41,6 @@ describe("API Dashboard", () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe("API Recordings", () => {
|
||||
it("GET /recordings returns items with pagination", async () => {
|
||||
const { status, body } = await api("/recordings?limit=5");
|
||||
expect(status).toBe(200);
|
||||
expect(body).toHaveProperty("items");
|
||||
expect(Array.isArray(body.items)).toBe(true);
|
||||
if (body.items.length > 0) {
|
||||
expect(body.items[0]).toHaveProperty("id");
|
||||
expect(body.items[0]).toHaveProperty("username");
|
||||
expect(body.items[0]).toHaveProperty("created_at");
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
describe("API Config", () => {
|
||||
it("GET /config returns 200", async () => {
|
||||
const { status } = await api("/config");
|
||||
@@ -62,13 +48,6 @@ describe("API Config", () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe("API Voice", () => {
|
||||
it("GET /guilds returns 200", async () => {
|
||||
const { status } = await api("/guilds");
|
||||
expect(status).toBe(200);
|
||||
});
|
||||
});
|
||||
|
||||
describe("API Negative", () => {
|
||||
it("GET /nonexistent returns 404", async () => {
|
||||
const { status } = await api("/nonexistent");
|
||||
|
||||
@@ -1,13 +0,0 @@
|
||||
import { z } from "zod";
|
||||
|
||||
export const mediaQueueSchema = z.object({
|
||||
source: z.string().min(1, "source is required"),
|
||||
mode: z.enum(["music", "screen"]).default("music"),
|
||||
});
|
||||
|
||||
export const mediaLoopSchema = z.object({
|
||||
loop: z.boolean().default(false),
|
||||
});
|
||||
|
||||
export type MediaQueueInput = z.infer<typeof mediaQueueSchema>;
|
||||
export type MediaLoopInput = z.infer<typeof mediaLoopSchema>;
|
||||
@@ -1,171 +0,0 @@
|
||||
import {
|
||||
createChildLogger,
|
||||
tryCommandThenFallback,
|
||||
} from "../../shared/commandHelper.js";
|
||||
import {
|
||||
COMMAND_MEDIA_LOOP,
|
||||
COMMAND_MEDIA_QUEUE,
|
||||
COMMAND_MEDIA_SKIP,
|
||||
COMMAND_MEDIA_STOP,
|
||||
MEDIA_STATUS_KEY,
|
||||
} from "../../shared/index.js";
|
||||
import { publishCommand, readRedisStatus } from "../../shared/redis/index.js";
|
||||
|
||||
const logger = createChildLogger("media.service");
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Types — match frontend exactly
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export interface MediaItem {
|
||||
id?: string;
|
||||
source: string;
|
||||
title: string;
|
||||
mode?: "music" | "screen";
|
||||
durationMs?: number | null;
|
||||
thumbnailUrl?: string | null;
|
||||
}
|
||||
|
||||
export interface MediaState {
|
||||
playing: boolean;
|
||||
/** null/absent when idle; "music" | "screen" while a track is active. */
|
||||
activeMode?: "music" | "screen" | null;
|
||||
musicVolume: number;
|
||||
loop: boolean;
|
||||
current: MediaItem | null;
|
||||
queue: MediaItem[];
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Defaults
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
const DEFAULT_COMMAND_TIMEOUT_MS = 5000;
|
||||
|
||||
const DEFAULT_STATE: MediaState = {
|
||||
playing: false,
|
||||
activeMode: null,
|
||||
musicVolume: 0.3,
|
||||
loop: false,
|
||||
current: null,
|
||||
queue: [],
|
||||
};
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Normalisation — handle both boolean (new) and string (legacy) playing values
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
function normalizeMediaState(raw: Record<string, unknown>): MediaState {
|
||||
const rawPlaying = raw.playing;
|
||||
const playing =
|
||||
rawPlaying === true ||
|
||||
rawPlaying === "playing" ||
|
||||
rawPlaying === "buffering";
|
||||
const mode = raw.activeMode;
|
||||
const activeMode: "music" | "screen" | null =
|
||||
mode === "music" || mode === "screen" ? mode : null;
|
||||
return {
|
||||
playing,
|
||||
activeMode,
|
||||
musicVolume: Number(raw.musicVolume ?? 0.3),
|
||||
loop: Boolean(raw.loop ?? false),
|
||||
current: (raw.current as MediaItem | null) ?? null,
|
||||
queue: (raw.queue as MediaItem[]) ?? [],
|
||||
};
|
||||
}
|
||||
|
||||
type MediaReplyData = Record<string, unknown> | MediaState;
|
||||
|
||||
function _fromReply(data: MediaReplyData): MediaState {
|
||||
return normalizeMediaState(data as Record<string, unknown>);
|
||||
}
|
||||
|
||||
async function readStatusFallback(): Promise<MediaState> {
|
||||
const cached = await readRedisStatus(MEDIA_STATUS_KEY);
|
||||
return cached ? normalizeMediaState(cached) : DEFAULT_STATE;
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Service methods
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Read media status from Redis key `media:status` set by discord-gateway.
|
||||
*/
|
||||
export async function getStatus(): Promise<MediaState> {
|
||||
logger.debug("getStatus called");
|
||||
return readStatusFallback();
|
||||
}
|
||||
|
||||
/**
|
||||
* Queue a media source via Redis command to discord-gateway.
|
||||
*/
|
||||
export async function queue(
|
||||
source: string,
|
||||
mode: "music" | "screen" = "music",
|
||||
): Promise<MediaState> {
|
||||
logger.info({ source, mode }, "queue called");
|
||||
return tryCommandThenFallback(
|
||||
() =>
|
||||
publishCommand<MediaState>(
|
||||
COMMAND_MEDIA_QUEUE,
|
||||
// NOTE: gateway MediaHandler reads `payload.url` (not `source`) —
|
||||
// keep the field name aligned or playback silently no-ops.
|
||||
{ url: source, mode },
|
||||
DEFAULT_COMMAND_TIMEOUT_MS,
|
||||
),
|
||||
() => readStatusFallback(),
|
||||
"queue",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Skip current track via Redis command to discord-gateway.
|
||||
*/
|
||||
export async function skip(): Promise<MediaState> {
|
||||
logger.info("skip called");
|
||||
return tryCommandThenFallback(
|
||||
() =>
|
||||
publishCommand<MediaState>(
|
||||
COMMAND_MEDIA_SKIP,
|
||||
{},
|
||||
DEFAULT_COMMAND_TIMEOUT_MS,
|
||||
),
|
||||
() => readStatusFallback(),
|
||||
"skip",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Stop playback via Redis command to discord-gateway.
|
||||
*/
|
||||
export async function stop(): Promise<MediaState> {
|
||||
logger.info("stop called");
|
||||
return tryCommandThenFallback(
|
||||
() =>
|
||||
publishCommand<MediaState>(
|
||||
COMMAND_MEDIA_STOP,
|
||||
{},
|
||||
DEFAULT_COMMAND_TIMEOUT_MS,
|
||||
),
|
||||
() => readStatusFallback(),
|
||||
"stop",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Toggle loop mode (replay current track on natural end) via Redis command.
|
||||
*/
|
||||
export async function setLoop(loop: boolean): Promise<MediaState> {
|
||||
logger.info({ loop }, "setLoop called");
|
||||
return tryCommandThenFallback(
|
||||
() =>
|
||||
publishCommand<MediaState>(
|
||||
COMMAND_MEDIA_LOOP,
|
||||
{ loop },
|
||||
DEFAULT_COMMAND_TIMEOUT_MS,
|
||||
),
|
||||
() => readStatusFallback(),
|
||||
"setLoop",
|
||||
);
|
||||
}
|
||||
@@ -1,184 +0,0 @@
|
||||
import {
|
||||
and,
|
||||
desc,
|
||||
eq,
|
||||
gte,
|
||||
ilike,
|
||||
lt,
|
||||
lte,
|
||||
or,
|
||||
type SQL,
|
||||
sql,
|
||||
} from "drizzle-orm";
|
||||
import { getDatabase } from "../../shared/database/index.js";
|
||||
import { pgVoiceRecordingsTable } from "../../shared/index.js";
|
||||
import { createChildLogger } from "../../shared/logger/index.js";
|
||||
|
||||
const logger = createChildLogger("recordings.service");
|
||||
|
||||
export interface RecordingRow {
|
||||
id: string;
|
||||
user_id: string;
|
||||
username: string;
|
||||
avatar_url: string | null;
|
||||
guild_id: string | null;
|
||||
channel_id: string | null;
|
||||
channel_name: string | null;
|
||||
filename: string;
|
||||
size_bytes: number;
|
||||
download_url: string | null;
|
||||
upload_status: string;
|
||||
upload_error: string | null;
|
||||
created_at: number;
|
||||
uploaded_at: number | null;
|
||||
transcription: string | null;
|
||||
}
|
||||
|
||||
export interface PaginatedRecordings {
|
||||
items: RecordingRow[];
|
||||
nextCursor: string | null;
|
||||
hasMore: boolean;
|
||||
}
|
||||
|
||||
export interface RecordingFilters {
|
||||
channelId?: string;
|
||||
userId?: string;
|
||||
cursor?: string;
|
||||
/** keyword search against transcription + username (ILIKE) */
|
||||
q?: string;
|
||||
/** created_at lower bound (epoch ms) */
|
||||
startDate?: number;
|
||||
/** created_at upper bound (epoch ms) */
|
||||
endDate?: number;
|
||||
}
|
||||
|
||||
export interface SpeakerSummary {
|
||||
user_id: string;
|
||||
username: string;
|
||||
avatar_url: string | null;
|
||||
clips: number;
|
||||
/** summed MP3 bytes; ~128kbps → est duration in seconds */
|
||||
est_duration_s: number;
|
||||
words: number;
|
||||
transcribed: number;
|
||||
last_at: number;
|
||||
}
|
||||
|
||||
export class RecordingsService {
|
||||
async getRecent(
|
||||
limit = 50,
|
||||
filters: RecordingFilters = {},
|
||||
): Promise<PaginatedRecordings> {
|
||||
logger.info({ limit, filters }, "getRecent called");
|
||||
const db = getDatabase();
|
||||
|
||||
const conditions: SQL[] = [];
|
||||
|
||||
if (filters.cursor) {
|
||||
conditions.push(
|
||||
lt(pgVoiceRecordingsTable.created_at, Number(filters.cursor)),
|
||||
);
|
||||
}
|
||||
if (filters.channelId) {
|
||||
conditions.push(eq(pgVoiceRecordingsTable.channel_id, filters.channelId));
|
||||
}
|
||||
if (filters.userId) {
|
||||
conditions.push(eq(pgVoiceRecordingsTable.user_id, filters.userId));
|
||||
}
|
||||
if (filters.q) {
|
||||
const like = `%${filters.q}%`;
|
||||
const qOr = or(
|
||||
ilike(pgVoiceRecordingsTable.transcription, like),
|
||||
ilike(pgVoiceRecordingsTable.username, like),
|
||||
);
|
||||
if (qOr) conditions.push(qOr);
|
||||
}
|
||||
if (filters.startDate) {
|
||||
conditions.push(
|
||||
gte(pgVoiceRecordingsTable.created_at, filters.startDate),
|
||||
);
|
||||
}
|
||||
if (filters.endDate) {
|
||||
conditions.push(lte(pgVoiceRecordingsTable.created_at, filters.endDate));
|
||||
}
|
||||
|
||||
const where = conditions.length > 0 ? and(...conditions) : undefined;
|
||||
|
||||
const allRows = await db
|
||||
.select({
|
||||
id: pgVoiceRecordingsTable.id,
|
||||
user_id: pgVoiceRecordingsTable.user_id,
|
||||
username: pgVoiceRecordingsTable.username,
|
||||
avatar_url: pgVoiceRecordingsTable.avatar_url,
|
||||
guild_id: pgVoiceRecordingsTable.guild_id,
|
||||
channel_id: pgVoiceRecordingsTable.channel_id,
|
||||
channel_name: pgVoiceRecordingsTable.channel_name,
|
||||
filename: pgVoiceRecordingsTable.filename,
|
||||
size_bytes: pgVoiceRecordingsTable.size_bytes,
|
||||
download_url: pgVoiceRecordingsTable.download_url,
|
||||
upload_status: pgVoiceRecordingsTable.upload_status,
|
||||
upload_error: pgVoiceRecordingsTable.upload_error,
|
||||
created_at: pgVoiceRecordingsTable.created_at,
|
||||
uploaded_at: pgVoiceRecordingsTable.uploaded_at,
|
||||
transcription: pgVoiceRecordingsTable.transcription,
|
||||
})
|
||||
.from(pgVoiceRecordingsTable)
|
||||
.where(where)
|
||||
.orderBy(desc(pgVoiceRecordingsTable.created_at))
|
||||
.limit(limit + 1);
|
||||
|
||||
const items = allRows.slice(0, limit) as unknown as RecordingRow[];
|
||||
const hasMore = allRows.length > limit;
|
||||
const nextCursor = hasMore
|
||||
? String(items[items.length - 1]?.created_at)
|
||||
: null;
|
||||
|
||||
return { items, nextCursor, hasMore };
|
||||
}
|
||||
|
||||
async deleteById(id: string): Promise<void> {
|
||||
const db = getDatabase();
|
||||
await db
|
||||
.delete(pgVoiceRecordingsTable)
|
||||
.where(eq(pgVoiceRecordingsTable.id, id));
|
||||
}
|
||||
|
||||
/**
|
||||
* Speaker leaderboard: aggregate clips / est. duration / transcribed words
|
||||
* per user. `est_duration_s` derives from summed MP3 bytes at the uploader's
|
||||
* 128 kbps transcode rate; `words` sums transcription word counts.
|
||||
*/
|
||||
async getSummary(): Promise<SpeakerSummary[]> {
|
||||
const db = getDatabase();
|
||||
const t = pgVoiceRecordingsTable;
|
||||
|
||||
const rows = await db
|
||||
.select({
|
||||
user_id: t.user_id,
|
||||
username: t.username,
|
||||
avatar_url: t.avatar_url,
|
||||
clips: sql<number>`count(*)::int`,
|
||||
total_bytes: sql<number>`coalesce(sum(${t.size_bytes}), 0)`,
|
||||
words: sql<number>`coalesce(sum(array_length(string_to_array(${t.transcription}, ' '), 1)), 0)::int`,
|
||||
transcribed: sql<number>`count(${t.transcription})::int`,
|
||||
last_at: sql<number>`max(${t.created_at})`,
|
||||
})
|
||||
.from(t)
|
||||
.groupBy(t.user_id, t.username, t.avatar_url)
|
||||
.orderBy(sql`clips desc, last_at desc`);
|
||||
|
||||
return rows.map((r) => ({
|
||||
user_id: r.user_id,
|
||||
username: r.username,
|
||||
avatar_url: r.avatar_url,
|
||||
clips: r.clips,
|
||||
// 128 kbps = 128000 bps → bytes*8/128000 = seconds
|
||||
est_duration_s: Math.round((r.total_bytes * 8) / 128000),
|
||||
words: r.words,
|
||||
transcribed: r.transcribed,
|
||||
last_at: r.last_at,
|
||||
}));
|
||||
}
|
||||
}
|
||||
|
||||
export const recordingsService = new RecordingsService();
|
||||
@@ -1,111 +0,0 @@
|
||||
/**
|
||||
* Authoritative live-voice store.
|
||||
*
|
||||
* Single source of truth for who is present / speaking in voice. The backend
|
||||
* WebSocket server is the one relay every frontend client connects to, so it
|
||||
* is the correct place to aggregate the gateway's `voice_active_user` deltas
|
||||
* into a shared snapshot. A late-joining browser must be able to see the same
|
||||
* state as everyone else — this store makes that possible (seeded into the WS
|
||||
* initial states and served via GET /api/voice/status).
|
||||
*/
|
||||
|
||||
export interface LiveSpeaker {
|
||||
userId: string;
|
||||
username: string;
|
||||
avatar?: string | null;
|
||||
speaking: boolean;
|
||||
/** Epoch ms of the most recent activity (start OR end of speech). */
|
||||
lastActiveAt: number;
|
||||
}
|
||||
|
||||
const speakers = new Map<string, LiveSpeaker>();
|
||||
|
||||
const MAX_SPEAKERS = 200;
|
||||
|
||||
/**
|
||||
* Speakers inactive for longer than this are auto-expired from the snapshot.
|
||||
* This handles the case where the gateway disconnects abruptly and never
|
||||
* sends `speaking: false` for active users.
|
||||
*/
|
||||
const SPEAKER_TTL_MS = 30_000;
|
||||
|
||||
/** Purge speakers that haven't been active recently. */
|
||||
function purgeStale(): void {
|
||||
const cutoff = Date.now() - SPEAKER_TTL_MS;
|
||||
for (const [id, s] of speakers) {
|
||||
if (!s.speaking && s.lastActiveAt < cutoff) {
|
||||
speakers.delete(id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Record a voice_active_user event. `speaking: true` upserts the speaker as
|
||||
* active; `speaking: false` marks them inactive while keeping them briefly
|
||||
* for the activity timeline (until TTL expiry).
|
||||
*/
|
||||
export function recordSpeaker(data: {
|
||||
userId: string;
|
||||
username?: string;
|
||||
avatar?: string | null;
|
||||
speaking: boolean;
|
||||
}): void {
|
||||
const { userId, speaking } = data;
|
||||
const existing = speakers.get(userId);
|
||||
const speaker: LiveSpeaker = {
|
||||
userId,
|
||||
username: data.username ?? existing?.username ?? "Unknown",
|
||||
avatar: data.avatar ?? existing?.avatar ?? null,
|
||||
speaking,
|
||||
lastActiveAt: Date.now(),
|
||||
};
|
||||
|
||||
if (speakers.size >= MAX_SPEAKERS && !existing) {
|
||||
let oldestId: string | null = null;
|
||||
let oldestTs = Infinity;
|
||||
for (const [id, s] of speakers) {
|
||||
if (!s.speaking && s.lastActiveAt < oldestTs) {
|
||||
oldestTs = s.lastActiveAt;
|
||||
oldestId = id;
|
||||
}
|
||||
}
|
||||
if (oldestId) speakers.delete(oldestId);
|
||||
else return;
|
||||
}
|
||||
|
||||
speakers.set(userId, speaker);
|
||||
}
|
||||
|
||||
/**
|
||||
* All recently-active speakers, most recently active first.
|
||||
* Stale (non-speaking + old) entries are auto-purged.
|
||||
*/
|
||||
export function getActiveSpeakers(): LiveSpeaker[] {
|
||||
purgeStale();
|
||||
return [...speakers.values()].sort((a, b) => b.lastActiveAt - a.lastActiveAt);
|
||||
}
|
||||
|
||||
/** Only speakers currently flagged as speaking. */
|
||||
export function getSpeakingSpeakers(): LiveSpeaker[] {
|
||||
purgeStale();
|
||||
return [...speakers.values()]
|
||||
.filter((s) => s.speaking)
|
||||
.sort((a, b) => b.lastActiveAt - a.lastActiveAt);
|
||||
}
|
||||
|
||||
/**
|
||||
* Mark ALL tracked speakers as not-speaking and purge stale ones.
|
||||
* Called when the gateway disconnects from voice — ensures the authoritative
|
||||
* snapshot doesn't carry ghost speakers.
|
||||
*/
|
||||
export function clearAllSpeakers(): void {
|
||||
for (const [, s] of speakers) {
|
||||
s.speaking = false;
|
||||
}
|
||||
purgeStale();
|
||||
}
|
||||
|
||||
/** Drop all tracked speakers (used on backend restart). */
|
||||
export function resetLiveSpeakers(): void {
|
||||
speakers.clear();
|
||||
}
|
||||
@@ -1,13 +0,0 @@
|
||||
import { z } from "zod";
|
||||
|
||||
export const connectVoiceSchema = z.object({
|
||||
guildId: z.string().min(1, "guildId is required"),
|
||||
channelId: z.string().min(1, "channelId is required"),
|
||||
});
|
||||
|
||||
export const voiceCommandSchema = z.object({
|
||||
command: z.string().min(1, "command is required"),
|
||||
});
|
||||
|
||||
export type ConnectVoiceInput = z.infer<typeof connectVoiceSchema>;
|
||||
export type VoiceCommandInput = z.infer<typeof voiceCommandSchema>;
|
||||
@@ -1,193 +0,0 @@
|
||||
import { eq } from "drizzle-orm";
|
||||
import {
|
||||
createChildLogger,
|
||||
tryCommandThenFallback,
|
||||
} from "../../shared/commandHelper.js";
|
||||
import { getDatabase } from "../../shared/database/index.js";
|
||||
import {
|
||||
COMMAND_GUILDS_LIST,
|
||||
COMMAND_GUILDS_TEXT_CHANNELS,
|
||||
COMMAND_VOICE_CHANNELS,
|
||||
COMMAND_VOICE_CONNECT,
|
||||
COMMAND_VOICE_DISCONNECT,
|
||||
type CommandReply,
|
||||
pgMessagesTable,
|
||||
VOICE_STATUS_KEY,
|
||||
} from "../../shared/index.js";
|
||||
import { publishCommand, readRedisStatus } from "../../shared/redis/index.js";
|
||||
import { getActiveSpeakers, type LiveSpeaker } from "./live-speaker.js";
|
||||
|
||||
const logger = createChildLogger("voice.service");
|
||||
|
||||
export interface Guild {
|
||||
id: string;
|
||||
name: string;
|
||||
icon: string | null;
|
||||
}
|
||||
|
||||
export interface Channel {
|
||||
id: string;
|
||||
name: string;
|
||||
type: "voice" | "text";
|
||||
/** Whether the selfbot account can actually join this voice channel. */
|
||||
joinable?: boolean;
|
||||
}
|
||||
|
||||
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[];
|
||||
/**
|
||||
* Authoritative shared voice snapshot — who is present / speaking right
|
||||
* now, aggregated server-side from the gateway's `voice_active_user`
|
||||
* deltas. All browsers converge on this same list.
|
||||
*/
|
||||
activeSpeakers: LiveSpeaker[];
|
||||
}
|
||||
|
||||
export const DEFAULT_VOICE_STATUS: VoiceStatus = {
|
||||
connected: false,
|
||||
activeGuildId: null,
|
||||
activeChannelId: null,
|
||||
activeChannelName: null,
|
||||
connections: [],
|
||||
activeSpeakers: [],
|
||||
};
|
||||
|
||||
/** Attach the live speaker snapshot to any voice status payload. */
|
||||
function withActiveSpeakers<T extends Partial<VoiceStatus>>(
|
||||
status: T,
|
||||
): T & { activeSpeakers: LiveSpeaker[] } {
|
||||
return { ...status, activeSpeakers: getActiveSpeakers() };
|
||||
}
|
||||
|
||||
/**
|
||||
* Wraps tryCommandThenFallback with a cleaner signature for use within this module.
|
||||
* Attempts a Redis command first; on failure, falls back to the provided function.
|
||||
*/
|
||||
async function withFallback<T>(
|
||||
commandFn: () => Promise<CommandReply<T> | null>,
|
||||
fallbackFn: () => Promise<T>,
|
||||
name: string,
|
||||
): Promise<T> {
|
||||
return tryCommandThenFallback(commandFn, fallbackFn, name);
|
||||
}
|
||||
|
||||
function readVoiceStatusFallback(): Promise<VoiceStatus> {
|
||||
return readRedisStatus(VOICE_STATUS_KEY).then((cached) =>
|
||||
withActiveSpeakers(
|
||||
(cached as unknown as VoiceStatus) ?? DEFAULT_VOICE_STATUS,
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Get guilds — query from discord-gateway via Redis command for real names.
|
||||
* Falls back to database (distinct guild_id from messages) if gateway unreachable.
|
||||
*/
|
||||
export async function getGuilds(): Promise<Guild[]> {
|
||||
logger.info("getGuilds called");
|
||||
return withFallback(
|
||||
() => publishCommand<Guild[]>(COMMAND_GUILDS_LIST, {}),
|
||||
async () => {
|
||||
const db = getDatabase();
|
||||
const rows = await db
|
||||
.selectDistinct({ guild_id: pgMessagesTable.guild_id })
|
||||
.from(pgMessagesTable)
|
||||
.orderBy(pgMessagesTable.guild_id);
|
||||
return rows.map((row) => ({
|
||||
id: String(row.guild_id ?? ""),
|
||||
name: `Guild ${String(row.guild_id).slice(0, 8)}`,
|
||||
icon: null,
|
||||
}));
|
||||
},
|
||||
"getGuilds",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Get text channels — query from discord-gateway via Redis command for real names.
|
||||
* Falls back to database if gateway unreachable.
|
||||
*/
|
||||
export async function getTextChannels(guildId: string): Promise<Channel[]> {
|
||||
logger.info({ guildId }, "getTextChannels called");
|
||||
return withFallback(
|
||||
() => publishCommand<Channel[]>(COMMAND_GUILDS_TEXT_CHANNELS, { guildId }),
|
||||
async () => {
|
||||
const db = getDatabase();
|
||||
const rows = await db
|
||||
.selectDistinct({ channel_id: pgMessagesTable.channel_id })
|
||||
.from(pgMessagesTable)
|
||||
.where(eq(pgMessagesTable.guild_id, guildId))
|
||||
.orderBy(pgMessagesTable.channel_id);
|
||||
return rows.map((row) => ({
|
||||
id: String(row.channel_id ?? ""),
|
||||
name: `Channel ${String(row.channel_id).slice(0, 8)}`,
|
||||
type: "text" as const,
|
||||
}));
|
||||
},
|
||||
"getTextChannels",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Get voice channels — query from discord-gateway via Redis command.
|
||||
*/
|
||||
export async function getVoiceChannels(guildId: string): Promise<Channel[]> {
|
||||
logger.info({ guildId }, "getVoiceChannels called");
|
||||
const reply = await publishCommand<Channel[]>(COMMAND_VOICE_CHANNELS, {
|
||||
guildId,
|
||||
});
|
||||
return reply?.success && reply.data ? reply.data : [];
|
||||
}
|
||||
|
||||
/**
|
||||
* Get current voice connection status from Redis cache set by discord-gateway.
|
||||
*/
|
||||
export async function getVoiceStatus(): Promise<VoiceStatus> {
|
||||
logger.debug("getVoiceStatus called");
|
||||
const cached = await readRedisStatus(VOICE_STATUS_KEY);
|
||||
return withActiveSpeakers(
|
||||
(cached as unknown as VoiceStatus) ?? DEFAULT_VOICE_STATUS,
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Connect to a voice channel via Redis command to discord-gateway.
|
||||
*/
|
||||
export async function connectVoice(
|
||||
guildId: string,
|
||||
channelId: string,
|
||||
): Promise<VoiceStatus> {
|
||||
logger.info({ guildId, channelId }, "connectVoice called");
|
||||
return withFallback(
|
||||
() =>
|
||||
publishCommand<VoiceStatus>(COMMAND_VOICE_CONNECT, {
|
||||
guildId,
|
||||
channelId,
|
||||
}),
|
||||
() => readVoiceStatusFallback(),
|
||||
"connectVoice",
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Disconnect from voice via Redis command to discord-gateway.
|
||||
*/
|
||||
export async function disconnectVoice(): Promise<VoiceStatus> {
|
||||
logger.info("disconnectVoice called");
|
||||
return withFallback(
|
||||
() => publishCommand<VoiceStatus>(COMMAND_VOICE_DISCONNECT, {}),
|
||||
() => readVoiceStatusFallback(),
|
||||
"disconnectVoice",
|
||||
);
|
||||
}
|
||||
@@ -5,33 +5,13 @@ import { chatRequestSchema } from "../modules/chatbot/chatbot.schema";
|
||||
import { chatbotService } from "../modules/chatbot/chatbot.service";
|
||||
import { dashboardService } from "../modules/dashboard/dashboard.service";
|
||||
import { knowledgeService } from "../modules/knowledge/knowledge.service";
|
||||
import {
|
||||
mediaLoopSchema,
|
||||
mediaQueueSchema,
|
||||
} from "../modules/media/media.schema";
|
||||
import {
|
||||
getStatus,
|
||||
queue,
|
||||
setLoop,
|
||||
skip,
|
||||
stop,
|
||||
} from "../modules/media/media.service";
|
||||
import {
|
||||
messageQuerySchema,
|
||||
semanticSearchSchema,
|
||||
} from "../modules/messages/messages.schema";
|
||||
import { messagesService } from "../modules/messages/messages.service";
|
||||
import { moderationService } from "../modules/moderation/moderation.service";
|
||||
import { recordingsService } from "../modules/recordings/recordings.service";
|
||||
import { uiStateService } from "../modules/ui-state/ui-state.service";
|
||||
import {
|
||||
connectVoice,
|
||||
disconnectVoice,
|
||||
getGuilds,
|
||||
getTextChannels,
|
||||
getVoiceChannels,
|
||||
getVoiceStatus,
|
||||
} from "../modules/voice/voice.service";
|
||||
import { config } from "../shared/config/index";
|
||||
import { publishCommandNoReply } from "../shared/redis/index";
|
||||
|
||||
@@ -238,86 +218,8 @@ const moderationRouter = {
|
||||
.handler(({ input }) => moderationService.getCoverage(input.days)),
|
||||
};
|
||||
|
||||
// ── Media ────────────────────────────────────────────────────────
|
||||
const mediaRouter = {
|
||||
status: os.handler(() => getStatus()),
|
||||
queue: os.input(mediaQueueSchema).handler(async ({ input }) => {
|
||||
await queue(input.source, input.mode);
|
||||
return getStatus();
|
||||
}),
|
||||
skip: os.handler(async () => {
|
||||
await skip();
|
||||
return getStatus();
|
||||
}),
|
||||
stop: os.handler(async () => {
|
||||
await stop();
|
||||
return getStatus();
|
||||
}),
|
||||
loop: os.input(mediaLoopSchema).handler(async ({ input }) => {
|
||||
await setLoop(input.loop);
|
||||
return getStatus();
|
||||
}),
|
||||
};
|
||||
|
||||
// ── Voice ─────────────────────────────────────────────────────────
|
||||
const voiceRouter = {
|
||||
guilds: os.handler(() => getGuilds()),
|
||||
textChannels: os
|
||||
.input(z.object({ guildId: z.string() }))
|
||||
.handler(({ input }) => getTextChannels(input.guildId)),
|
||||
voiceChannels: os
|
||||
.input(z.object({ guildId: z.string() }))
|
||||
.handler(({ input }) => getVoiceChannels(input.guildId)),
|
||||
status: os.handler(() => getVoiceStatus()),
|
||||
connect: os
|
||||
.input(z.object({ guildId: z.string(), channelId: z.string() }))
|
||||
.handler(async ({ input }) => {
|
||||
await connectVoice(input.guildId, input.channelId);
|
||||
return getVoiceStatus();
|
||||
}),
|
||||
disconnect: os.handler(async () => {
|
||||
await disconnectVoice();
|
||||
return getVoiceStatus();
|
||||
}),
|
||||
command: os
|
||||
.input(z.object({ command: z.string().min(1) }))
|
||||
.handler(async ({ input }) => {
|
||||
await publishCommandNoReply(input.command);
|
||||
return { success: true, command: input.command };
|
||||
}),
|
||||
};
|
||||
|
||||
// ── Recordings ───────────────────────────────────────────────────
|
||||
const recordingsRouter = {
|
||||
list: os
|
||||
.input(
|
||||
z.object({
|
||||
limit: z.coerce.number().int().positive().default(50),
|
||||
channelId: z.string().optional(),
|
||||
userId: z.string().optional(),
|
||||
cursor: z.string().optional(),
|
||||
q: z.string().optional(),
|
||||
startDate: z.coerce.number().int().optional(),
|
||||
endDate: z.coerce.number().int().optional(),
|
||||
}),
|
||||
)
|
||||
.handler(({ input }) =>
|
||||
recordingsService.getRecent(input.limit, {
|
||||
channelId: input.channelId,
|
||||
userId: input.userId,
|
||||
cursor: input.cursor,
|
||||
q: input.q,
|
||||
startDate: input.startDate,
|
||||
endDate: input.endDate,
|
||||
}),
|
||||
),
|
||||
delete: os.input(z.object({ id: z.string() })).handler(async ({ input }) => {
|
||||
await recordingsService.deleteById(input.id);
|
||||
return { ok: true };
|
||||
}),
|
||||
summary: os.handler(async () => recordingsService.getSummary()),
|
||||
};
|
||||
|
||||
// ── Recordings (removed) / Voice (removed) ─────────────────────
|
||||
// Music/media playback (media.mjs) was removed with the voice feature.
|
||||
// ── Analysis (search) ──────────────────────────────────────────────
|
||||
const analysisRouter = {
|
||||
search: os
|
||||
@@ -416,11 +318,8 @@ const configRouter = {
|
||||
backlogSyncBatchSize: config.BACKLOG_SYNC_BATCH_SIZE,
|
||||
retentionMessagesDays: config.RETENTION_MESSAGES_DAYS,
|
||||
retentionAttachmentsDays: config.RETENTION_ATTACHMENTS_DAYS,
|
||||
retentionVoiceDays: config.RETENTION_VOICE_DAYS,
|
||||
autoDeleteFlaggedEnabled: config.AUTO_DELETE_FLAGGED_ENABLED,
|
||||
aiAnalysisEnabled: config.AI_ANALYSIS_ENABLED,
|
||||
voiceGuildId: config.VOICE_GUILD_ID || null,
|
||||
voiceChannelId: config.VOICE_CHANNEL_ID || null,
|
||||
logLevel: config.LOG_LEVEL,
|
||||
})),
|
||||
};
|
||||
@@ -438,9 +337,6 @@ export const appRouter = {
|
||||
dashboard: dashboardRouter,
|
||||
messages: messagesRouter,
|
||||
moderation: moderationRouter,
|
||||
media: mediaRouter,
|
||||
voice: voiceRouter,
|
||||
recordings: recordingsRouter,
|
||||
analysis: analysisRouter,
|
||||
chatbot: chatbotRouter,
|
||||
config: configRouter,
|
||||
|
||||
@@ -33,29 +33,6 @@ export const configSchema = z
|
||||
.transform((v) => v.split(",").filter(Boolean))
|
||||
.describe("Thread IDs to exclude from capture"),
|
||||
|
||||
// ── Legacy voice ─────────────────────────────────────────────────────
|
||||
VOICE_GUILD_ID: z.string().min(1).optional(),
|
||||
VOICE_CHANNEL_ID: z.string().min(1).optional(),
|
||||
|
||||
// ── Recording ────────────────────────────────────────────────────────
|
||||
RECORDINGS_DIR: z.string().default("./recordings"),
|
||||
RECORDING_SEGMENT_MS: z.coerce.number().positive().default(5000),
|
||||
|
||||
// ── Decoder ──────────────────────────────────────────────────────────
|
||||
DECODER_ROTATE_MS: z.coerce.number().positive().default(5000),
|
||||
DECODER_COOLDOWN_MS: z.coerce.number().positive().default(30000),
|
||||
|
||||
// ── Audio ────────────────────────────────────────────────────────────
|
||||
AUDIO_STREAM_SILENCE_DURATION_MS: z.coerce
|
||||
.number()
|
||||
.positive()
|
||||
.default(3000),
|
||||
PACKET_FILTER_MIN_SIZE: z.coerce.number().positive().default(8),
|
||||
OPUS_FRAME_SIZE: z.coerce.number().positive().default(960),
|
||||
AUDIO_SAMPLE_RATE: z.coerce.number().positive().default(48000),
|
||||
AUDIO_CHANNELS: z.coerce.number().positive().default(2),
|
||||
AVATAR_SIZE: z.coerce.number().positive().default(64),
|
||||
|
||||
// ── Server ───────────────────────────────────────────────────────────
|
||||
WEBSERVER_PORT: z.coerce.number().positive().default(3001),
|
||||
NODE_ENV: z
|
||||
@@ -92,17 +69,8 @@ export const configSchema = z
|
||||
|
||||
// ── Redis ────────────────────────────────────────────────────────────
|
||||
REDIS_URL: z.string().default("redis://localhost:6379"),
|
||||
// ── Voice PCM WebSocket (direct gateway→backend, bypasses Redis) ────
|
||||
VOICE_PCM_WS_ENABLED: z
|
||||
.string()
|
||||
.optional()
|
||||
.transform((v) => v === "true")
|
||||
.default(true),
|
||||
BACKEND_WS_URL: z.string().default("ws://backend:3000/ws"),
|
||||
BACKEND_WS_TOKEN: z.string().optional().default(""),
|
||||
|
||||
// ── Connection ───────────────────────────────────────────────────────
|
||||
VOICE_CONNECTION_TIMEOUT_MS: z.coerce.number().positive().default(15000),
|
||||
RECONNECT_TIMEOUT_MS: z.coerce.number().positive().default(5000),
|
||||
|
||||
// ── Attachments ─────────────────────────────────────────────────────
|
||||
@@ -203,13 +171,6 @@ export const configSchema = z
|
||||
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
|
||||
.string()
|
||||
.optional()
|
||||
.transform((v) => v === "true")
|
||||
.default(false),
|
||||
|
||||
// ── OpenAI Moderation ───────────────────────────────────────────────
|
||||
OPENAI_MODERATION_API_KEY: z.string().optional(),
|
||||
OPENAI_MODERATION_BASE_URL: z
|
||||
@@ -252,7 +213,6 @@ export const configSchema = z
|
||||
// ── Retention ───────────────────────────────────────────────────────
|
||||
RETENTION_MESSAGES_DAYS: z.coerce.number().int().min(0).default(0),
|
||||
RETENTION_ATTACHMENTS_DAYS: z.coerce.number().int().min(0).default(0),
|
||||
RETENTION_VOICE_DAYS: z.coerce.number().int().min(0).default(0),
|
||||
RETENTION_CLEANUP_INTERVAL_MS: z.coerce
|
||||
.number()
|
||||
.positive()
|
||||
@@ -291,7 +251,6 @@ export const configSchema = z
|
||||
|
||||
export type AppConfig = z.infer<typeof configSchema> & {
|
||||
EFFECTIVE_TEXT_GUILD_ID?: string;
|
||||
EFFECTIVE_VOICE_GUILD_ID?: string;
|
||||
EFFECTIVE_MONITOR_GUILD_IDS: string[];
|
||||
};
|
||||
|
||||
@@ -301,7 +260,6 @@ export function loadConfig(env: NodeJS.ProcessEnv = process.env): AppConfig {
|
||||
return {
|
||||
...parsed,
|
||||
EFFECTIVE_TEXT_GUILD_ID: parsed.TEXT_GUILD_ID ?? parsed.MONITOR_GUILD_ID,
|
||||
EFFECTIVE_VOICE_GUILD_ID: parsed.VOICE_GUILD_ID,
|
||||
EFFECTIVE_MONITOR_GUILD_IDS:
|
||||
parsed.MONITOR_GUILD_IDS.length > 0
|
||||
? parsed.MONITOR_GUILD_IDS
|
||||
|
||||
@@ -24,9 +24,6 @@ export interface BroadcasterClient {
|
||||
messageAnalyzed: (data: unknown) => void;
|
||||
attachmentCreated: (data: unknown) => void;
|
||||
attachmentUploaded: (data: unknown) => void;
|
||||
voiceRecordingStarted: (data: unknown) => void;
|
||||
voiceRecordingStopped: (data: unknown) => void;
|
||||
voiceRecordingUploaded: (data: unknown) => void;
|
||||
analysisQueueStatus: (data: unknown) => void;
|
||||
}
|
||||
|
||||
@@ -102,17 +99,6 @@ export interface AttachmentRecord {
|
||||
uploaded_at: number | null;
|
||||
}
|
||||
|
||||
export interface VoiceSegmentRecord {
|
||||
id: string;
|
||||
user_id: string;
|
||||
session_id: string;
|
||||
guild_id: string;
|
||||
channel_id: string;
|
||||
filename: string;
|
||||
duration_ms: number;
|
||||
created_at: number;
|
||||
}
|
||||
|
||||
export interface DashboardMessage {
|
||||
id: string;
|
||||
channel_id: string;
|
||||
@@ -121,7 +107,7 @@ export interface DashboardMessage {
|
||||
avatar_url: string | null;
|
||||
content: string;
|
||||
created_at: number;
|
||||
type: "text" | "image" | "voice";
|
||||
type: "text" | "image";
|
||||
}
|
||||
|
||||
export interface MessageQuery {
|
||||
@@ -154,23 +140,6 @@ export interface AnalysisResult {
|
||||
evidence?: string[];
|
||||
}
|
||||
|
||||
export interface VoiceRecordingUploadData {
|
||||
id: string;
|
||||
user_id: string;
|
||||
username: string;
|
||||
avatar_url: string | null;
|
||||
guild_id: string | null;
|
||||
channel_id: string | null;
|
||||
channel_name: string | null;
|
||||
filename: string;
|
||||
size_bytes: number;
|
||||
download_url: string;
|
||||
upload_status: string;
|
||||
created_at: number;
|
||||
uploaded_at: number;
|
||||
transcription?: string | null;
|
||||
}
|
||||
|
||||
export interface AnalysisQueueStatus {
|
||||
queuedConversations: number;
|
||||
activeRequests: number;
|
||||
@@ -221,7 +190,6 @@ export interface RetentionPolicy {
|
||||
channel_id: string | null;
|
||||
retention_days: number;
|
||||
apply_to_media: boolean;
|
||||
apply_to_voice: boolean;
|
||||
enabled: boolean;
|
||||
created_at: number;
|
||||
updated_at: number;
|
||||
|
||||
@@ -15,11 +15,6 @@ export const DISCORD_MESSAGE_DELETED = "discord:message:deleted";
|
||||
export const DISCORD_MESSAGE_ANALYZED = "discord:message:analyzed";
|
||||
export const DISCORD_ATTACHMENT_CREATED = "discord:attachment:created";
|
||||
export const DISCORD_ATTACHMENT_UPLOADED = "discord:attachment:uploaded";
|
||||
export const DISCORD_VOICE_STARTED = "discord:voice:started";
|
||||
export const DISCORD_VOICE_STOPPED = "discord:voice:stopped";
|
||||
export const DISCORD_VOICE_UPLOADED = "discord:voice:uploaded";
|
||||
export const DISCORD_VOICE_ACTIVE_USER = "discord:voice:active_user";
|
||||
export const DISCORD_VOICE_PCM = "discord:voice:pcm";
|
||||
export const DISCORD_ANALYSIS_QUEUE_STATUS = "discord:analysis:queue_status";
|
||||
export const DISCORD_REACTION_ADDED = "discord:reaction:added";
|
||||
export const DISCORD_REACTION_REMOVED = "discord:reaction:removed";
|
||||
@@ -37,35 +32,15 @@ export const DISCORD_MODERATION_ACTION = "discord:moderation:action";
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export const BACKEND_COMMAND = "backend:command";
|
||||
export const BACKEND_VOICE_TRANSMIT = "backend:voice:transmit";
|
||||
export const BACKEND_COMMAND_REPLY_PREFIX = "backend:command:reply:";
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Status keys (set by discord-gateway, read by backend via Redis GET)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export const VOICE_STATUS_KEY = "voice:status";
|
||||
export const MEDIA_STATUS_KEY = "media:status";
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Command types (used as the `type` field in CommandMessage envelopes)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export const COMMAND_VOICE_CONNECT = "voice:connect";
|
||||
export const COMMAND_VOICE_DISCONNECT = "voice:disconnect";
|
||||
export const COMMAND_VOICE_DISCONNECT_GUILD = "voice:disconnect:guild";
|
||||
export const COMMAND_VOICE_CHANNELS = "voice:channels";
|
||||
export const COMMAND_VOICE_TRANSMIT_START = "voice:transmit:start";
|
||||
export const COMMAND_VOICE_TRANSMIT_STOP = "voice:transmit:stop";
|
||||
export const COMMAND_GUILDS_LIST = "guilds:list";
|
||||
export const COMMAND_GUILDS_TEXT_CHANNELS = "guilds:text-channels";
|
||||
export const COMMAND_MEDIA_QUEUE = "media:queue";
|
||||
export const COMMAND_MEDIA_SKIP = "media:skip";
|
||||
export const COMMAND_MEDIA_STOP = "media:stop";
|
||||
export const COMMAND_MEDIA_VOLUME = "media:volume";
|
||||
export const COMMAND_MEDIA_LOOP = "media:loop";
|
||||
export const COMMAND_MODERATION_ACTION = "moderation:action";
|
||||
export const DISCORD_VOICE_ANALYZED = "discord:voice:analyzed";
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Event envelope — used by discord-gateway when publishing to Redis
|
||||
@@ -111,13 +86,7 @@ export const DISCORD_CHANNEL_TO_WS_EVENT: Record<string, string> = {
|
||||
[DISCORD_MESSAGE_ANALYZED]: "message_analyzed",
|
||||
[DISCORD_ATTACHMENT_CREATED]: "attachment_created",
|
||||
[DISCORD_ATTACHMENT_UPLOADED]: "attachment_uploaded",
|
||||
[DISCORD_VOICE_STARTED]: "voice_recording_started",
|
||||
[DISCORD_VOICE_STOPPED]: "voice_recording_stopped",
|
||||
[DISCORD_VOICE_UPLOADED]: "voice_recording_uploaded",
|
||||
[DISCORD_ANALYSIS_QUEUE_STATUS]: "analysis_queue_status",
|
||||
[DISCORD_VOICE_ACTIVE_USER]: "voice_active_user",
|
||||
[DISCORD_VOICE_PCM]: "voice_pcm_data",
|
||||
[DISCORD_VOICE_ANALYZED]: "voice_analyzed",
|
||||
[DISCORD_REACTION_ADDED]: "reaction_added",
|
||||
[DISCORD_REACTION_REMOVED]: "reaction_removed",
|
||||
[DISCORD_THREAD_CREATED]: "thread_created",
|
||||
|
||||
@@ -49,7 +49,6 @@ export function clearBroadcastFunctions(): void {
|
||||
function shouldLog(type: string): boolean {
|
||||
if (!_enabled) return false;
|
||||
// Avoid logging high-volume events
|
||||
if (type === "voice_pcm_data") return false;
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
@@ -1,17 +1,8 @@
|
||||
import Redis from "ioredis";
|
||||
import {
|
||||
clearAllSpeakers,
|
||||
recordSpeaker,
|
||||
} from "../modules/voice/live-speaker.js";
|
||||
import { config } from "../shared/config/index.js";
|
||||
import {
|
||||
DISCORD_CHANNEL_TO_WS_EVENT,
|
||||
DISCORD_VOICE_ACTIVE_USER,
|
||||
DISCORD_VOICE_PCM,
|
||||
DISCORD_VOICE_STOPPED,
|
||||
} from "../shared/index.js";
|
||||
import { DISCORD_CHANNEL_TO_WS_EVENT } from "../shared/index.js";
|
||||
import { createChildLogger } from "../shared/logger/index.js";
|
||||
import { broadcastBinary, broadcastEvent } from "./broadcast.js";
|
||||
import { broadcastEvent } from "./broadcast.js";
|
||||
|
||||
const logger = createChildLogger("ws.redis-bridge");
|
||||
|
||||
@@ -49,67 +40,10 @@ function handleSubscriptionMessage(channel: string, message: string): void {
|
||||
// We only want <actual payload>, not the full envelope.
|
||||
const data = envelope.data !== undefined ? envelope.data : envelope;
|
||||
|
||||
// Voice PCM: decode base64 → binary broadcast instead of JSON
|
||||
if (channel === DISCORD_VOICE_PCM) {
|
||||
const pcmPayload = data as { userId?: string; pcm?: string };
|
||||
if (pcmPayload?.pcm && pcmPayload?.userId) {
|
||||
try {
|
||||
const pcmBuffer = Buffer.from(pcmPayload.pcm, "base64");
|
||||
// Prepend userId as 4-byte FNV-1a hash
|
||||
const userIdHash = hashUserId(pcmPayload.userId);
|
||||
const binary = Buffer.alloc(4 + pcmBuffer.length);
|
||||
binary.writeUInt32LE(userIdHash, 0);
|
||||
pcmBuffer.copy(binary, 4);
|
||||
broadcastBinary(binary);
|
||||
return;
|
||||
} catch {
|
||||
// fallback to JSON broadcast on error
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Aggregate live-voice state authoritatively BEFORE broadcasting.
|
||||
// Every browser hears the same `voice_active_user` deltas, so the backend
|
||||
// can maintain the single shared snapshot for late-joining clients.
|
||||
if (channel === DISCORD_VOICE_ACTIVE_USER) {
|
||||
const speaker = data as {
|
||||
userId?: string;
|
||||
username?: string;
|
||||
avatar?: string | null;
|
||||
speaking?: boolean;
|
||||
};
|
||||
if (speaker?.userId) {
|
||||
recordSpeaker({
|
||||
userId: speaker.userId,
|
||||
username: speaker.username,
|
||||
avatar: speaker.avatar,
|
||||
speaking: Boolean(speaker.speaking),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
// When the gateway stops voice recording (disconnects from voice channel),
|
||||
// clear all speakers from the authoritative snapshot so frontends don't
|
||||
// show ghost participants.
|
||||
if (channel === DISCORD_VOICE_STOPPED) {
|
||||
clearAllSpeakers();
|
||||
logger.info("Voice recording stopped — cleared all live speakers");
|
||||
}
|
||||
|
||||
logger.debug({ channel, eventType }, "Broadcasting Redis event");
|
||||
broadcastEvent(eventType, data);
|
||||
}
|
||||
|
||||
/** Simple 32-bit FNV-1a hash for userId → 4-byte identifier */
|
||||
function hashUserId(userId: string): number {
|
||||
let hash = 0x811c9dc5;
|
||||
for (let i = 0; i < userId.length; i++) {
|
||||
hash ^= userId.charCodeAt(i);
|
||||
hash = Math.imul(hash, 0x01000193);
|
||||
}
|
||||
return hash >>> 0;
|
||||
}
|
||||
|
||||
export async function startRedisBridge(): Promise<void> {
|
||||
if (!config.REDIS_URL) {
|
||||
logger.info("Redis not configured, skipping Redis bridge");
|
||||
|
||||
@@ -2,8 +2,6 @@ import type { IncomingMessage, Server } from "node:http";
|
||||
import type { Duplex } from "node:stream";
|
||||
import { WebSocket, WebSocketServer } from "ws";
|
||||
import { messagesService } from "../modules/messages/messages.service.js";
|
||||
import { config } from "../shared/config/index.js";
|
||||
import { BACKEND_COMMAND, BACKEND_VOICE_TRANSMIT } from "../shared/index.js";
|
||||
import { createChildLogger } from "../shared/logger/index.js";
|
||||
import { setBroadcastFunctions } from "./broadcast.js";
|
||||
|
||||
@@ -17,8 +15,6 @@ interface BroadcastEvent {
|
||||
|
||||
interface JsonMessage {
|
||||
type: string;
|
||||
buffer?: string;
|
||||
command?: string;
|
||||
payload?: Record<string, unknown>;
|
||||
}
|
||||
|
||||
@@ -54,36 +50,6 @@ async function sendInitialStates(ws: WebSocket): Promise<void> {
|
||||
} catch (err) {
|
||||
logger.warn({ err }, "Failed to send initial ui_state");
|
||||
}
|
||||
|
||||
// Send initial media state
|
||||
try {
|
||||
const { getStatus } = await import("../modules/media/media.service.js");
|
||||
const mediaState = await getStatus();
|
||||
ws.send(
|
||||
JSON.stringify({
|
||||
type: "media_state",
|
||||
state: mediaState,
|
||||
}),
|
||||
);
|
||||
} catch (err) {
|
||||
logger.warn({ err }, "Failed to send initial media_state");
|
||||
}
|
||||
|
||||
// Send initial live-voice snapshot (shared authoritative state — a browser
|
||||
// joining mid-call sees the same speakers as everyone else, not an empty DB).
|
||||
try {
|
||||
const { getActiveSpeakers } = await import(
|
||||
"../modules/voice/live-speaker.js"
|
||||
);
|
||||
ws.send(
|
||||
JSON.stringify({
|
||||
type: "voice_state",
|
||||
state: { activeSpeakers: getActiveSpeakers() },
|
||||
}),
|
||||
);
|
||||
} catch (err) {
|
||||
logger.warn({ err }, "Failed to send initial voice_state");
|
||||
}
|
||||
}
|
||||
|
||||
export function closeWebSocketServer(): void {
|
||||
@@ -94,9 +60,7 @@ export function closeWebSocketServer(): void {
|
||||
}
|
||||
|
||||
export function createWebSocketServer(server: Server): WebSocketServer {
|
||||
// Separate tracking: gateway sends PCM → forwarded to frontend only
|
||||
const frontendClients = new Set<WebSocket>();
|
||||
const gatewayClients = new Set<WebSocket>();
|
||||
|
||||
const wss = new WebSocketServer({ noServer: true, perMessageDeflate: true });
|
||||
_wss = wss;
|
||||
@@ -115,35 +79,6 @@ export function createWebSocketServer(server: Server): WebSocketServer {
|
||||
// Map-based dispatcher for JSON WebSocket message types
|
||||
const jsonHandlers = new Map<string, MessageHandler>();
|
||||
|
||||
jsonHandlers.set("voice_transmit", async (_ws, message) => {
|
||||
if (!message.buffer) return;
|
||||
const { getCommandPublisher } = await import("../shared/redis/index.js");
|
||||
const publisher = getCommandPublisher();
|
||||
await publisher.publish(
|
||||
BACKEND_VOICE_TRANSMIT,
|
||||
JSON.stringify({ type: "pcm", buffer: message.buffer }),
|
||||
);
|
||||
});
|
||||
|
||||
jsonHandlers.set("voice_command", async (_ws, message) => {
|
||||
if (!message.command) return;
|
||||
const { getCommandPublisher } = await import("../shared/redis/index.js");
|
||||
const publisher = getCommandPublisher();
|
||||
const commandId = `cmd-${Date.now()}-${Math.random().toString(36).slice(2, 9)}`;
|
||||
await publisher.publish(
|
||||
BACKEND_COMMAND,
|
||||
JSON.stringify({
|
||||
id: commandId,
|
||||
type: message.command,
|
||||
payload: message.payload ?? {},
|
||||
replyChannel: `reply:${commandId}`,
|
||||
}),
|
||||
);
|
||||
});
|
||||
|
||||
// Stream historical messages one-by-one over WS (no 50-row batch).
|
||||
// The frontend requests it once per channel switch; the backend emits one
|
||||
// `message_snapshot` frame per message so the UI renders progressively.
|
||||
jsonHandlers.set("stream_messages", async (ws, message) => {
|
||||
if (ws.readyState !== WebSocket.OPEN) return;
|
||||
const payload = (message.payload ?? {}) as {
|
||||
@@ -208,74 +143,15 @@ export function createWebSocketServer(server: Server): WebSocketServer {
|
||||
}
|
||||
});
|
||||
|
||||
wss.on("connection", (ws: WebSocket, req) => {
|
||||
// Parse auth token from query string
|
||||
const rawUrl = req.url ?? "/";
|
||||
let isGateway = false;
|
||||
|
||||
try {
|
||||
const url = new URL(rawUrl, "http://localhost");
|
||||
const token = url.searchParams.get("token");
|
||||
isGateway =
|
||||
token !== null &&
|
||||
config.BACKEND_WS_TOKEN !== "" &&
|
||||
token === config.BACKEND_WS_TOKEN;
|
||||
} catch {
|
||||
// Malformed URL — treat as frontend
|
||||
}
|
||||
|
||||
if (isGateway) {
|
||||
gatewayClients.add(ws);
|
||||
logger.info("Discord gateway WebSocket client authenticated");
|
||||
// Gateway doesn't need initial states
|
||||
} else {
|
||||
frontendClients.add(ws);
|
||||
logger.info(`Frontend client connected (${frontendClients.size} total)`);
|
||||
// Send initial states (user, ui, media) — fire-and-forget
|
||||
sendInitialStates(ws).catch((err) =>
|
||||
logger.error({ err }, "sendInitialStates failed"),
|
||||
);
|
||||
}
|
||||
wss.on("connection", (ws: WebSocket) => {
|
||||
frontendClients.add(ws);
|
||||
logger.info(`Frontend client connected (${frontendClients.size} total)`);
|
||||
// Send initial states (user, ui) — fire-and-forget
|
||||
sendInitialStates(ws).catch((err) =>
|
||||
logger.error({ err }, "sendInitialStates failed"),
|
||||
);
|
||||
|
||||
ws.on("message", (data: Buffer) => {
|
||||
// Gateway PCM forward — broadcast raw binary to frontend clients only
|
||||
if (isGateway && Buffer.isBuffer(data)) {
|
||||
broadcastBinary(data);
|
||||
return;
|
||||
}
|
||||
|
||||
// Handle binary PCM from browser (FE→Discord transmit)
|
||||
// Format: 4-byte magic "PCM\0" + raw PCM Int16 LE
|
||||
if (
|
||||
Buffer.isBuffer(data) &&
|
||||
data.length > 4 &&
|
||||
data[0] === 0x50 && // 'P'
|
||||
data[1] === 0x43 && // 'C'
|
||||
data[2] === 0x4d && // 'M'
|
||||
data[3] === 0x00 // '\0'
|
||||
) {
|
||||
const pcmBuffer = data.subarray(4);
|
||||
const base64 = pcmBuffer.toString("base64");
|
||||
import("../shared/redis/index.js").then(({ getCommandPublisher }) => {
|
||||
const publisher = getCommandPublisher();
|
||||
publisher
|
||||
.publish(
|
||||
BACKEND_VOICE_TRANSMIT,
|
||||
JSON.stringify({
|
||||
type: "pcm",
|
||||
buffer: base64,
|
||||
}),
|
||||
)
|
||||
.catch((err: Error) => {
|
||||
logger.error(
|
||||
{ err },
|
||||
"Failed to publish voice transmit to Redis",
|
||||
);
|
||||
});
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
// Handle JSON messages from browser
|
||||
if (
|
||||
typeof data === "string" ||
|
||||
@@ -296,24 +172,15 @@ export function createWebSocketServer(server: Server): WebSocketServer {
|
||||
});
|
||||
|
||||
ws.on("close", () => {
|
||||
if (isGateway) {
|
||||
gatewayClients.delete(ws);
|
||||
logger.info("Discord gateway WebSocket disconnected");
|
||||
} else {
|
||||
frontendClients.delete(ws);
|
||||
logger.info(
|
||||
`Frontend client disconnected (${frontendClients.size} total)`,
|
||||
);
|
||||
}
|
||||
frontendClients.delete(ws);
|
||||
logger.info(
|
||||
`Frontend client disconnected (${frontendClients.size} total)`,
|
||||
);
|
||||
});
|
||||
|
||||
ws.on("error", (err: Error) => {
|
||||
logger.error({ err }, "WebSocket client error");
|
||||
if (isGateway) {
|
||||
gatewayClients.delete(ws);
|
||||
} else {
|
||||
frontendClients.delete(ws);
|
||||
}
|
||||
frontendClients.delete(ws);
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user