feat: implement media echo fix and YouTube screenshare design

- Introduced a new `ScreenShareController` to manage YouTube screenshare functionality.
- Updated `DiscordPlayer` to track ownership of audio streams, preventing conflicts between music playback and screenshare.
- Added error handling for various states including voice connection checks and media busy states.
- Created unit tests for `ScreenShareController` and `DiscordPlayer` ownership rules to ensure correct functionality.
- Added documentation for the new media echo fix and screenshare design.
This commit is contained in:
MythEclipse
2026-05-16 15:48:28 +07:00
parent e32e092596
commit d50ce8698f
21 changed files with 2284 additions and 51 deletions
+69 -5
View File
@@ -3,10 +3,14 @@ import { discordPlayer } from "../player";
import { MediaQueue } from "./mediaQueue";
import { resolveMediaSource } from "./mediaResolver";
import type {
MediaMode,
MediaState,
MusicPlayback,
MusicPlayer,
QueueMediaOptions,
ResolvedMediaSource,
ScreenShareController,
ScreenSharePlayback,
} from "./mediaTypes";
import { createMusicPlayer } from "./musicPlayer";
@@ -15,6 +19,7 @@ export interface MediaControllerDependencies {
isBrowserStreaming?: () => boolean;
resolveMediaSource?: (source: string) => Promise<ResolvedMediaSource>;
musicPlayer?: MusicPlayer;
screenController?: ScreenShareController;
onStateChange?: (state: MediaState) => void;
}
@@ -24,6 +29,8 @@ export class MediaController {
private playback: MusicPlayback | null = null;
private playbackToken = 0;
private skipInProgress = false;
private screenPlayback: ScreenSharePlayback | null = null;
private activeMode: MediaMode | null = null;
constructor(private readonly dependencies: MediaControllerDependencies = {}) {
this.musicPlayer = dependencies.musicPlayer ?? createMusicPlayer();
@@ -32,17 +39,27 @@ export class MediaController {
getState(): MediaState {
const snapshot = this.queueStore.snapshot();
return {
playing: snapshot.current?.status === "playing",
playing:
this.activeMode === "screen" || snapshot.current?.status === "playing",
activeMode: this.activeMode ?? snapshot.current?.mode ?? null,
...snapshot,
};
}
async queue(source: string): Promise<MediaState> {
this.assertCanStart();
async queue(
source: string,
options: QueueMediaOptions = {},
): Promise<MediaState> {
const mode = options.mode ?? "music";
if (mode === "screen") {
return this.startScreen(source);
}
this.assertCanStartMusic();
const resolved = await (
this.dependencies.resolveMediaSource ?? resolveMediaSource
)(source);
this.queueStore.add(resolved);
this.queueStore.add(resolved, mode, options.requestedBy);
this.startNextIfIdle();
return this.emitState();
}
@@ -73,11 +90,14 @@ export class MediaController {
this.playbackToken++;
this.playback?.stop();
this.playback = null;
this.screenPlayback?.stop();
this.screenPlayback = null;
this.activeMode = null;
this.queueStore.clear();
return this.emitState();
}
private assertCanStart(): void {
private assertCanStartMusic(): void {
const isVoiceConnected =
this.dependencies.isVoiceConnected ?? (() => discordPlayer.isConnected());
if (!isVoiceConnected()) {
@@ -88,6 +108,10 @@ export class MediaController {
);
}
if (this.screenPlayback || this.dependencies.screenController?.isActive()) {
throw new AppError("Another media mode is active", "MEDIA_BUSY", 409);
}
if (this.dependencies.isBrowserStreaming?.()) {
throw new AppError(
"Stop browser microphone streaming before playing media",
@@ -97,6 +121,46 @@ export class MediaController {
}
}
private async startScreen(source: string): Promise<MediaState> {
if (
this.screenPlayback ||
this.dependencies.screenController?.isActive() ||
this.playback ||
this.queueStore.snapshot().current
) {
throw new AppError("Another media mode is active", "MEDIA_BUSY", 409);
}
const screenController = this.dependencies.screenController;
if (!screenController) {
throw new AppError(
"Screen sharing is unavailable",
"SCREEN_UNAVAILABLE",
500,
);
}
this.activeMode = "screen";
try {
this.screenPlayback = await screenController.start(source);
} catch (error) {
this.activeMode = null;
throw error;
}
this.screenPlayback.done.then(
() => this.finishScreen(),
() => this.finishScreen(),
);
return this.emitState();
}
private finishScreen(): void {
if (!this.screenPlayback || this.activeMode !== "screen") return;
this.screenPlayback = null;
this.activeMode = null;
this.emitState();
}
private startNextIfIdle(): void {
if (this.playback) return;
const item = this.queueStore.startNext();
+7 -2
View File
@@ -1,4 +1,5 @@
import type {
MediaMode,
MediaQueueItem,
MediaState,
ResolvedMediaSource,
@@ -13,10 +14,14 @@ export class MediaQueue {
private readonly now = () => Date.now(),
) {}
add(source: ResolvedMediaSource, requestedBy = "dashboard"): MediaQueueItem {
add(
source: ResolvedMediaSource,
mode: MediaQueueItem["mode"] = "music",
requestedBy = "dashboard",
): MediaQueueItem {
const item: MediaQueueItem = {
id: this.createId(),
mode: "music",
mode,
requestedBy,
addedAt: this.now(),
status: "queued",
+24 -3
View File
@@ -25,10 +25,16 @@ export interface MediaQueueItem extends ResolvedMediaSource {
export interface MediaState {
playing: boolean;
activeMode: MediaMode | null;
current: MediaQueueItem | null;
queue: MediaQueueItem[];
}
export interface QueueMediaOptions {
mode?: MediaMode;
requestedBy?: string;
}
export interface MusicPlayback {
done: Promise<void>;
stop(): void;
@@ -38,8 +44,23 @@ export interface MusicPlayer {
play(source: ResolvedMediaSource): MusicPlayback;
}
export interface DiscordAudioPlayer {
isConnected(): boolean;
playStream(stream: Readable): void;
export interface ScreenSharePlayback {
done: Promise<void>;
stop(): void;
}
export interface ScreenShareController {
isActive(): boolean;
start(source: string): Promise<ScreenSharePlayback>;
}
export type DiscordPlayerOwner = "none" | "browser-bridge" | "music" | "screen";
export interface DiscordAudioPlayer {
getOwner(): DiscordPlayerOwner;
isConnected(): boolean;
playStream(stream: Readable, owner: DiscordPlayerOwner): void;
pause(owner?: DiscordPlayerOwner): void;
unpause(owner?: DiscordPlayerOwner): boolean;
stop(owner?: DiscordPlayerOwner): void;
}
+18 -4
View File
@@ -30,13 +30,27 @@ export function createMusicPlayer(
}) as unknown as ChildProcessWithoutNullStreams;
proc.stderr.resume();
audioPlayer.playStream(proc.stdout);
audioPlayer.playStream(proc.stdout, "music");
let stopped = false;
let released = false;
const release = () => {
if (released) return;
released = true;
audioPlayer.stop("music");
};
const done = new Promise<void>((resolve, reject) => {
proc.on("error", reject);
proc.stdout.on("error", reject);
proc.on("error", (error) => {
release();
reject(error);
});
proc.stdout.on("error", (error) => {
release();
reject(error);
});
proc.on("close", (code) => {
release();
if (code === 0 || stopped) {
resolve();
return;
@@ -51,7 +65,7 @@ export function createMusicPlayer(
if (stopped) return;
stopped = true;
proc.kill("SIGTERM");
audioPlayer.stop();
release();
},
};
},
+123
View File
@@ -0,0 +1,123 @@
import type { Readable } from "node:stream";
import {
playStream as defaultPlayStream,
prepareStream as defaultPrepareStream,
Encoders,
Utils,
} from "@dank074/discord-video-stream";
import { AppError } from "../errors";
import { discordPlayer } from "../player";
import type { DiscordPlayerOwner, ScreenSharePlayback } from "./mediaTypes";
import { createYtDlp } from "./ytdlp";
export interface ScreenShareVoiceStatus {
connected: boolean;
activeGuildId: string | null;
activeChannelId: string | null;
}
interface PreparedScreenStream {
command: { kill?: (signal: NodeJS.Signals) => unknown };
output: Readable;
}
type PrepareScreenStream = (
source: string,
options: object,
) => PreparedScreenStream;
type PlayScreenStream = (
output: Readable,
streamer: unknown,
options: { type: "go-live" },
) => Promise<void>;
export interface ScreenShareControllerDependencies {
getVoiceStatus: () => ScreenShareVoiceStatus;
getPlayerOwner?: () => DiscordPlayerOwner;
getDirectVideoUrl?: (source: string) => Promise<string>;
prepareStream?: PrepareScreenStream;
playStream?: PlayScreenStream;
streamer: unknown;
}
export function createScreenShareController(
dependencies: ScreenShareControllerDependencies,
) {
let active: ScreenSharePlayback | null = null;
const ytdlp = createYtDlp();
const getPlayerOwner =
dependencies.getPlayerOwner ?? (() => discordPlayer.getOwner());
const getDirectVideoUrl =
dependencies.getDirectVideoUrl ??
((source) => ytdlp.getDirectVideoUrl(source));
const prepareStream =
dependencies.prepareStream ??
(defaultPrepareStream as unknown as PrepareScreenStream);
const playStream =
dependencies.playStream ??
(defaultPlayStream as unknown as PlayScreenStream);
return {
isActive(): boolean {
return active !== null;
},
async start(source: string): Promise<ScreenSharePlayback> {
const status = dependencies.getVoiceStatus();
if (
!status.connected ||
!status.activeGuildId ||
!status.activeChannelId
) {
throw new AppError(
"Connect to a voice channel before sharing screen",
"VOICE_NOT_CONNECTED",
409,
);
}
if (active || getPlayerOwner() !== "none") {
throw new AppError("Another media mode is active", "MEDIA_BUSY", 409);
}
try {
const directUrl = await getDirectVideoUrl(source);
const { command, output } = prepareStream(directUrl, {
encoder: Encoders.software({ x264: { preset: "superfast" } }),
height: 720,
frameRate: 30,
bitrateVideo: 2500,
bitrateVideoMax: 4000,
includeAudio: true,
videoCodec: Utils.normalizeVideoCodec("H264"),
});
let stopped = false;
const done = playStream(output, dependencies.streamer, {
type: "go-live",
}).finally(() => {
active = null;
});
active = {
done,
stop() {
if (stopped) return;
stopped = true;
command.kill?.("SIGTERM");
active = null;
},
};
return active;
} catch (error) {
active = null;
throw new AppError(
error instanceof Error ? error.message : "Screen stream failed",
"SCREEN_STREAM_FAILED",
500,
);
}
},
};
}
+14
View File
@@ -9,6 +9,7 @@ export interface YtDlpMetadata {
export interface YtDlpClient {
getMetadata(url: string): Promise<YtDlpMetadata>;
getDirectAudioUrl(url: string): Promise<string>;
getDirectVideoUrl(url: string): Promise<string>;
}
export interface YtDlpDependencies {
@@ -49,6 +50,19 @@ export function createYtDlp(dependencies: YtDlpDependencies = {}): YtDlpClient {
]);
return value.trim().split("\n")[0] || url;
},
async getDirectVideoUrl(url: string): Promise<string> {
const value = await runYtDlp(spawn, [
url,
"--get-url",
"--format",
"bestvideo[protocol^=http]+bestaudio[protocol^=http]/best[protocol^=http]/best",
"--no-playlist",
"--no-warnings",
"--quiet",
]);
return value.trim().split("\n")[0] || url;
},
};
}
+3 -3
View File
@@ -144,7 +144,9 @@ async function runAnalysisInWorker(
messages: MessageRecord[],
): Promise<AnalysisWorkerResponse> {
return new Promise((resolve, reject) => {
const worker = new Worker(new URL("./aiAnalysisWorker.ts", import.meta.url));
const worker = new Worker(
new URL("./aiAnalysisWorker.ts", import.meta.url),
);
worker.once("message", (response: AnalysisWorkerResponse) => {
worker.terminate().catch((error) => {
@@ -213,7 +215,6 @@ function scheduleConversationAnalysis(conversationKey: string): void {
export async function queueMessageAnalysis(messageId: string): Promise<void> {
if (!config.AI_ANALYSIS_ENABLED) return;
try {
// Look up the message to get its conversation key
const message = await getMessageById(messageId);
@@ -242,7 +243,6 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
export function queueConversationAnalysis(conversationKey: string): void {
if (!config.AI_ANALYSIS_ENABLED) return;
// Schedule debounced analysis
scheduleConversationAnalysis(conversationKey);
}
+33 -5
View File
@@ -7,10 +7,12 @@ import {
StreamType,
VoiceConnection,
} from "@discordjs/voice";
import type { DiscordPlayerOwner } from "./media/mediaTypes";
export class DiscordPlayer {
private player: AudioPlayer;
private connection: VoiceConnection | null = null;
private owner: DiscordPlayerOwner = "none";
constructor() {
this.player = createAudioPlayer();
@@ -21,6 +23,7 @@ export class DiscordPlayer {
this.player.on("error", (error) => {
console.error(`[player] Error: ${error.message}`);
this.owner = "none";
});
}
@@ -29,17 +32,28 @@ export class DiscordPlayer {
this.connection.subscribe(this.player);
}
public getOwner(): DiscordPlayerOwner {
return this.owner;
}
public isConnected(): boolean {
return this.connection !== null;
}
public playStream(stream: Readable) {
console.log("[player] Starting new audio stream...");
public playStream(stream: Readable, owner: DiscordPlayerOwner) {
if (owner === "none") {
throw new Error("Discord audio player owner is required");
}
this.assertOwnerAvailable(owner);
const resource = createAudioResource(stream, {
inputType: StreamType.OggOpus,
});
if (this.owner === owner) {
this.player.stop();
}
this.owner = owner;
this.player.play(resource);
this.connection?.subscribe(this.player);
}
@@ -48,16 +62,30 @@ export class DiscordPlayer {
return this.player.state.status;
}
public pause() {
public pause(owner?: DiscordPlayerOwner) {
if (!this.canControl(owner)) return;
this.player.pause(true);
}
public unpause(): boolean {
public unpause(owner?: DiscordPlayerOwner): boolean {
if (!this.canControl(owner)) return false;
return this.player.unpause();
}
public stop() {
public stop(owner?: DiscordPlayerOwner) {
if (!this.canControl(owner)) return;
this.player.stop();
this.owner = "none";
}
private assertOwnerAvailable(owner: DiscordPlayerOwner): void {
if (this.owner !== "none" && this.owner !== owner) {
throw new Error(`Discord audio player is owned by ${this.owner}`);
}
}
private canControl(owner?: DiscordPlayerOwner): boolean {
return !owner || this.owner === "none" || this.owner === owner;
}
}
+2
View File
@@ -89,6 +89,8 @@ export async function startRecording(
// Dengarkan siapapun yang mulai bicara
receiver.speaking.on("start", async (userId) => {
if (userId === client.user?.id) return;
const userMetadata = await collectUserMetadata(client, userId, channel);
logger.info(
{ userId, username: userMetadata.username },
+9 -2
View File
@@ -2,6 +2,7 @@ import type { Router } from "express";
import express from "express";
import { AppError } from "../errors";
import type { MediaController } from "../media/mediaController";
import type { MediaMode } from "../media/mediaTypes";
export type MediaRouteController = Pick<
MediaController,
@@ -21,7 +22,10 @@ export function createMediaRoutes(controller: MediaRouteController): Router {
router.post("/media/queue", async (req, res, next) => {
try {
const { source } = req.body as { source?: string };
const { source, mode = "music" } = req.body as {
source?: string;
mode?: MediaMode;
};
if (!source) {
throw new AppError(
"Media source is required",
@@ -29,7 +33,10 @@ export function createMediaRoutes(controller: MediaRouteController): Router {
400,
);
}
res.json(await controller.queue(source));
if (mode !== "music" && mode !== "screen") {
throw new AppError("Invalid media mode", "INVALID_MEDIA_MODE", 400);
}
res.json(await controller.queue(source, { mode }));
} catch (error) {
next(error);
}
+29 -11
View File
@@ -1,15 +1,17 @@
import fs from "node:fs";
import http from "node:http";
import path from "node:path";
import { Streamer } from "@dank074/discord-video-stream";
import { AudioPlayerStatus } from "@discordjs/voice";
import type { Client } from "discord.js-selfbot-v13";
import express from "express";
import helmet from "helmet";
import { AudioPlayerStatus } from "@discordjs/voice";
import * as prism from "prism-media";
import { WebSocketServer } from "ws";
import { AppError } from "./errors";
import { createChildLogger, logger } from "./logger";
import { MediaController } from "./media/mediaController";
import { createScreenShareController } from "./media/screenShareController";
import { getMetrics, uptimeGauge } from "./metrics";
import { createBroadcaster } from "./moderation/broadcaster";
import type { ModerationBroadcaster } from "./moderation/types";
@@ -163,9 +165,16 @@ export async function startWebserver(
const broadcaster = createBroadcaster();
(globalThis as VoiceGlobals).moderationBroadcaster = broadcaster;
const streamer = new Streamer(_client);
const screenController = createScreenShareController({
getVoiceStatus: () => voiceController.getStatus(),
streamer,
});
const mediaController = new MediaController({
isVoiceConnected: () => voiceController.getStatus().connected,
isBrowserStreaming: () => sharedUIState.isStreaming,
screenController,
onStateChange: (state) => broadcaster.mediaState(state),
});
@@ -287,11 +296,12 @@ export async function startWebserver(
const SILENCE_TAIL_MS = 300; // continue sending silence for 300ms after browser stops
const MAX_BUF_BYTES = BYTES_PER_FRAME * 50; // cap at 1 second to avoid runaway buffer
let opusEncoder: prism.opus.Encoder;
let opusEncoder: prism.opus.Encoder | null = null;
let bridgePlayerPaused = true;
const SILENCE_FRAME = Buffer.alloc(BYTES_PER_FRAME, 0);
function startBrowserAudioBridge(): void {
if (opusEncoder) return;
opusEncoder = new prism.opus.Encoder({
rate: RATE,
channels: CHANNELS,
@@ -308,19 +318,23 @@ export async function startWebserver(
opusEncoder.on("error", () => {});
opusEncoder.pipe(oggBitstream);
opusEncoder.write(Buffer.alloc(BYTES_PER_FRAME, 0));
discordPlayer.playStream(oggBitstream);
discordPlayer.pause();
discordPlayer.playStream(oggBitstream, "browser-bridge");
discordPlayer.pause("browser-bridge");
bridgePlayerPaused = true;
}
function ensureBrowserAudioBridge(): void {
if (discordPlayer.getStatus() === AudioPlayerStatus.Idle) {
function ensureBrowserAudioBridge(): boolean {
const owner = discordPlayer.getOwner();
if (owner !== "none" && owner !== "browser-bridge") return false;
if (
owner === "none" ||
discordPlayer.getStatus() === AudioPlayerStatus.Idle
) {
startBrowserAudioBridge();
}
return true;
}
startBrowserAudioBridge();
let pcmBuffer = Buffer.alloc(0);
let lastBrowserAudioTime = 0;
@@ -351,9 +365,12 @@ export async function startWebserver(
dbAccum += rmsDb(frame);
dbCount++;
ensureBrowserAudioBridge();
if (!ensureBrowserAudioBridge()) {
pcmBuffer = Buffer.alloc(0);
return;
}
if (bridgePlayerPaused) {
const unpaused = discordPlayer.unpause();
const unpaused = discordPlayer.unpause("browser-bridge");
bridgePlayerPaused = false;
wsLogger.info({ unpaused }, "Transmitting — Discord indicator ON");
}
@@ -362,7 +379,7 @@ export async function startWebserver(
frame = SILENCE_FRAME;
} else if (!bridgePlayerPaused && msSinceAudio >= SILENCE_TAIL_MS) {
// No audio for a while — pause Discord indicator
discordPlayer.pause();
discordPlayer.pause("browser-bridge");
bridgePlayerPaused = true;
wsLogger.info("Stopped — Discord indicator OFF");
return;
@@ -371,6 +388,7 @@ export async function startWebserver(
}
// Write one frame. If encoder is backpressured, skip this tick to avoid stalling.
if (!opusEncoder) return;
const ok = opusEncoder.write(frame);
if (!ok) {
opusEncoder.once("drain", () => {}); // re-arm drain without blocking