feat: redesign UI and overhaul audio processing for Discord Gateway v4

This commit is contained in:
baharsah
2026-05-13 02:30:09 +07:00
parent 44ac346c21
commit ad7dcde47c
7 changed files with 396 additions and 297 deletions
+2
View File
@@ -1,4 +1,6 @@
import "./mock-crc";
import "libsodium-wrappers";
import "@snazzah/davey";
import { Client } from "discord.js-selfbot-v13";
import { startRecording } from "./recorder";
import { config } from "./config";
+2 -6
View File
@@ -33,13 +33,9 @@ export class DiscordPlayer {
public playStream(stream: Readable) {
console.log("[player] Starting new audio stream...");
// Use WebmDemuxer to extract Opus packets from browser stream
const demuxer = new prism.opus.WebmDemuxer();
demuxer.on('error', err => console.error("[player] Demuxer error:", err));
const resource = createAudioResource(stream.pipe(demuxer), {
inputType: StreamType.Opus,
const resource = createAudioResource(stream, {
inputType: StreamType.OggOpus,
});
this.player.play(resource);
+55 -39
View File
@@ -63,20 +63,22 @@ export async function startRecording(client: Client, channel: VoiceChannel): Pro
// Dengarkan siapapun yang mulai bicara
receiver.speaking.on("start", async (userId) => {
if (config.verbose) {
// console.log(`[recorder-debug] Speaking 'start' event triggered for userId: ${userId}. Subscriptions has? ${receiver.subscriptions.has(userId)}`);
// Coba ambil data user dari cache atau fetch dari API
const user = client.users.cache.get(userId) || await client.users.fetch(userId).catch(() => null);
const username = user?.username ?? "Unknown User";
const avatar = user?.displayAvatarURL({ format: 'png', size: 64 }) ?? "https://cdn.discordapp.com/embed/avatars/0.png";
// Tampilkan format "nama user [voice activity]"
console.log(`${username} [voice activity]`);
// Notify webserver
if ((global as any).updateActiveUser) {
(global as any).updateActiveUser(userId, { username, avatar, speaking: true });
}
// Jangan record kalau sudah ada stream aktif untuk user ini
if (receiver.subscriptions.has(userId)) return;
// Coba ambil data user dari cache atau fetch dari API
const user = client.users.cache.get(userId) || await client.users.fetch(userId).catch(() => null);
const username = user?.username ?? "Unknown User";
// Tampilkan format "nama user [voice activity]"
console.log(`${username} [voice activity]`);
const timestamp = Date.now();
const userDir = path.join(recordingsDir, userId);
if (!fs.existsSync(userDir)) {
@@ -88,38 +90,60 @@ export async function startRecording(client: Client, channel: VoiceChannel): Pro
const audioStream = receiver.subscribe(userId, {
end: {
behavior: EndBehaviorType.AfterSilence,
duration: 1000, // Stop 1 detik setelah user diam
duration: 3000, // 3 seconds — avoids FFmpeg restart overhead between utterances
},
});
try {
const packetFilter = new PacketFilter(10);
// --- OGG file recording (unchanged) ---
const packetFilterForOgg = new PacketFilter(8);
const oggStream = new prism.opus.OggLogicalBitstream({
opusHead: new prism.opus.OpusHead({
channelCount: 2,
sampleRate: 48000,
}),
pageSizeControl: {
maxPackets: 10,
},
crc: true, // Use our mock node-crc
opusHead: new prism.opus.OpusHead({ channelCount: 2, sampleRate: 48000 }),
pageSizeControl: { maxPackets: 10 },
crc: true,
});
const out = fs.createWriteStream(filename);
audioStream.pipe(packetFilterForOgg).pipe(oggStream).pipe(out);
// Pipe: audioStream -> packetFilter -> oggStream -> out
audioStream.pipe(packetFilter).pipe(oggStream).pipe(out);
// --- Web broadcast: pure JS Opus → PCM, no FFmpeg ---
// Create a fresh decoder for each user session
const opusDecoder = new prism.opus.Decoder({ frameSize: 960, channels: 2, rate: 48000 });
// Forward raw Opus packets to the web shared Ogg stream
packetFilter.on('data', (chunk) => {
if ((global as any).broadcastOpusToWeb) {
(global as any).broadcastOpusToWeb(chunk);
// CRITICAL: Swallow decode errors (DAVE/bad packets) without crashing
opusDecoder.on('error', () => {});
// Downsample 48kHz stereo → 24kHz mono (take left channel, every 2nd sample)
opusDecoder.on('data', (pcm: Buffer) => {
if (!(global as any).broadcastPcmToWeb) return;
// Input: 48kHz stereo s16le → 4 bytes per sample-pair
// Output: 24kHz mono s16le → 2 bytes per sample
const outBuf = Buffer.alloc(pcm.length / 4);
for (let i = 0; i < outBuf.length / 2; i++) {
outBuf.writeInt16LE(pcm.readInt16LE(i * 8), i * 2);
}
(global as any).broadcastPcmToWeb(outBuf, userId);
});
// Feed Opus packets one-by-one; catch per-packet decode errors
let packetCount = 0;
audioStream.on('data', (chunk: Buffer) => {
packetCount++;
if (packetCount <= 5) {
console.log(`[recorder] Pkt #${packetCount} from ${userId}: ${chunk.length}b | 0x${chunk.slice(0,4).toString('hex')}`);
}
if (chunk.length < 8) return; // skip tiny control packets
try {
opusDecoder.write(chunk);
} catch (_) {} // per-packet isolation — don't let one bad packet stop the stream
});
audioStream.on('end', () => {
opusDecoder.end();
if ((global as any).updateActiveUser) {
(global as any).updateActiveUser(userId, { username, avatar, speaking: false });
}
});
if (config.verbose) {
console.log(`[recorder] Recording user ${userId}${filename}`);
}
out.on('finish', async () => {
if (config.verbose) {
@@ -145,17 +169,9 @@ export async function startRecording(client: Client, channel: VoiceChannel): Pro
audioStream.on('error', (err) => {
console.error(`[recorder] Audio Stream error ${userId}:`, err.message);
});
audioStream.on('data', (chunk) => {
if (config.verbose) {
console.log(`[recorder-debug] Received audio packet from ${userId}, size: ${chunk.length} bytes`);
}
packetFilterForOgg.on('error', (err) => {
console.error(`[recorder] PacketFilter(ogg) error ${userId}:`, err.message);
});
packetFilter.on('error', (err) => {
console.error(`[recorder] Packet Filter error ${userId}:`, err.message);
});
out.on('error', (err) => {
console.error(`[recorder] File write error ${userId}:`, err.message);
});
+103 -68
View File
@@ -1,92 +1,127 @@
import express from "express";
import { WebSocketServer } from "ws";
import http from "http";
import { WebSocketServer } from "ws";
import path from "path";
import { PassThrough } from "stream";
import { discordPlayer } from "./player";
import prism from "prism-media";
import { discordPlayer } from "./player";
const activeUsers = new Map<string, { username: string, avatar: string, speaking: boolean }>();
let wsClients = new Set<any>();
// --- Upsampling: 24kHz mono s16le → 48kHz stereo s16le (pure JS, no FFmpeg) ---
// Each input sample is duplicated into 2 stereo pairs to double the sample rate.
function upsample24kMonoTo48kStereo(mono24k: Buffer): Buffer {
const out = Buffer.alloc(mono24k.length * 4); // 2x rate * 2ch = 4x bytes
for (let i = 0; i < mono24k.length / 2; i++) {
const s = mono24k.readInt16LE(i * 2);
out.writeInt16LE(s, i * 8); // t=0 L
out.writeInt16LE(s, i * 8 + 2); // t=0 R
out.writeInt16LE(s, i * 8 + 4); // t=1 L (duplicate for 2x rate)
out.writeInt16LE(s, i * 8 + 6); // t=1 R
}
return out;
}
export function startWebserver(port: number = 3000) {
const app = express();
const server = http.createServer(app);
const wss = new WebSocketServer({ server });
const listeners = new Set<express.Response>();
let headerChunks: Buffer[] = [];
// Create a single, continuous Ogg stream for all web listeners
const oggStream = new prism.opus.OggLogicalBitstream({
opusHead: new prism.opus.OpusHead({
channelCount: 2,
sampleRate: 48000,
}),
pageSizeControl: {
maxPackets: 10,
},
});
// Forward Ogg pages to all connected web listeners
oggStream.on("data", (chunk) => {
// Cache the first 2 chunks (headers)
if (headerChunks.length < 2) {
headerChunks.push(chunk);
}
listeners.forEach(res => res.write(chunk));
});
// Prime the stream with a silent packet to generate headers immediately
// Silent Opus packet (1 frame, 20ms)
const silentPacket = Buffer.from([0xf8, 0xff, 0xfe]);
oggStream.write(silentPacket);
const wsPort = port + 1;
const wss = new WebSocketServer({ port: wsPort, host: "0.0.0.0" });
console.log(`[webserver] WebSocket server listening on ws://0.0.0.0:${wsPort}`);
app.use(express.static(path.join(__dirname, "../public")));
// Endpoint for receiving (listening) audio from Discord
app.get("/listen", (req, res) => {
res.setHeader("Content-Type", "audio/ogg");
res.setHeader("Transfer-Encoding", "chunked");
res.setHeader("Connection", "keep-alive");
// Send cached headers immediately so the browser recognizes the stream
headerChunks.forEach(chunk => res.write(chunk));
listeners.add(res);
console.log(`[webserver] New listener connected. Total: ${listeners.size}`);
req.on("close", () => {
listeners.delete(res);
console.log(`[webserver] Listener disconnected. Total: ${listeners.size}`);
// --- Inbound: Discord PCM → tagged chunks → browser (set in recorder.ts) ---
(global as any).broadcastPcmToWeb = (chunk: Buffer, userId: string) => {
let hash = 0;
for (let i = 0; i < userId.length; i++) {
hash = ((hash << 5) - hash) + userId.charCodeAt(i);
hash |= 0;
}
const header = Buffer.alloc(4);
header.writeInt32LE(hash, 0);
const packet = Buffer.concat([header, chunk]);
wsClients.forEach(client => {
if (client.readyState === 1) client.send(packet);
});
});
// Function to broadcast raw Opus packets from Discord to the shared Ogg stream
(global as any).broadcastOpusToWeb = (chunk: Buffer) => {
oggStream.write(chunk);
};
(global as any).updateActiveUser = (userId: string, data: { username: string, avatar: string, speaking: boolean }) => {
activeUsers.set(userId, data);
broadcastUserState();
};
function broadcastUserState() {
const payload = JSON.stringify({
type: "user_state",
users: Array.from(activeUsers.entries()).map(([id, data]) => ({ id, ...data }))
});
wsClients.forEach(client => {
if (client.readyState === 1) client.send(payload);
});
}
// --- Outbound: browser PCM (24kHz mono) → Opus → Discord, NO FFmpeg ---
const RATE = 48000;
const CHANNELS = 2;
const FRAME_SIZE = 960; // 20ms @ 48kHz
const BYTES_PER_FRAME = FRAME_SIZE * CHANNELS * 2; // 3840 bytes
const opusEncoder = new prism.opus.Encoder({ rate: RATE, channels: CHANNELS, frameSize: FRAME_SIZE });
const oggBitstream = new prism.opus.OggLogicalBitstream({
opusHead: new prism.opus.OpusHead({ channelCount: CHANNELS, sampleRate: RATE }),
pageSizeControl: { maxPackets: 10 },
crc: true,
});
opusEncoder.on('error', () => {});
opusEncoder.pipe(oggBitstream);
// Prime the encoder immediately so OGG headers are emitted before player reads
opusEncoder.write(Buffer.alloc(BYTES_PER_FRAME, 0));
discordPlayer.playStream(oggBitstream);
let pcmBuffer = Buffer.alloc(0);
let lastBrowserAudioTime = 0;
const SILENCE_FRAME = Buffer.alloc(BYTES_PER_FRAME, 0);
// Keep encoder alive with silence when browser isn't sending
setInterval(() => {
if (Date.now() - lastBrowserAudioTime > 40) {
opusEncoder.write(SILENCE_FRAME);
}
}, 20);
wss.on("connection", (ws) => {
console.log("[webserver] New WebSocket connection");
console.log("[webserver] New WebSocket connection on port " + wsPort);
wsClients.add(ws);
const audioStream = new PassThrough();
discordPlayer.playStream(audioStream);
ws.send(JSON.stringify({
type: "user_state",
users: Array.from(activeUsers.entries()).map(([id, data]) => ({ id, ...data }))
}));
ws.on("message", (data: Buffer) => {
// console.log(`[webserver] Received chunk: ${data.length} bytes`);
audioStream.write(data);
ws.on("message", (data: any) => {
if (!Buffer.isBuffer(data)) return;
lastBrowserAudioTime = Date.now();
// Upsample browser 24kHz mono → 48kHz stereo
const upsampled = upsample24kMonoTo48kStereo(data);
pcmBuffer = Buffer.concat([pcmBuffer, upsampled]);
// Encode complete Opus frames
while (pcmBuffer.length >= BYTES_PER_FRAME) {
const frame = pcmBuffer.slice(0, BYTES_PER_FRAME);
pcmBuffer = pcmBuffer.slice(BYTES_PER_FRAME);
opusEncoder.write(frame);
}
});
ws.on("close", () => {
console.log("[webserver] WebSocket connection closed");
audioStream.end();
});
ws.on("error", (err) => {
console.error("[webserver] WebSocket error:", err);
audioStream.end();
});
ws.on("close", () => { wsClients.delete(ws); });
ws.on("error", () => { wsClients.delete(ws); });
});
server.listen(port, () => {
console.log(`[webserver] Server listening on http://localhost:${port}`);
server.listen(port, "0.0.0.0", () => {
console.log(`[webserver] Web interface listening on http://0.0.0.0:${port}`);
});
}