Compare commits

...

8 Commits

36 changed files with 1119 additions and 214 deletions

View File

@@ -1,30 +0,0 @@
DISCORD_TOKEN=test_token_for_testing
RECORDINGS_DIR=./recordings
RECORDING_SEGMENT_MS=5000
VERBOSE=false
DECODER_ROTATE_MS=5000
DECODER_COOLDOWN_MS=30000
AUDIO_STREAM_SILENCE_DURATION_MS=3000
PACKET_FILTER_MIN_SIZE=8
OPUS_FRAME_SIZE=960
AUDIO_SAMPLE_RATE=48000
AUDIO_CHANNELS=2
AVATAR_SIZE=64
WEBSERVER_PORT=3000
VOICE_CONNECTION_TIMEOUT_MS=15000
RECONNECT_TIMEOUT_MS=5000
LOG_LEVEL=info
NODE_ENV=test
MONITOR_GUILD_ID=test_guild_id
PICSER_UPLOAD_URL=https://picser.asepharyana.tech/api/upload
ATTACHMENT_UPLOAD_TIMEOUT_MS=30000
ATTACHMENT_MAX_SIZE_MB=100
ATTACHMENT_RETRY_ATTEMPTS=3
BACKLOG_SYNC_HOURS=24
BACKLOG_SYNC_BATCH_SIZE=100
AI_ANALYSIS_ENABLED=false
AI_LLM_API_KEY=test_key
AI_LLM_BASE_URL=https://9router.asepharyana.tech/v1
AI_LLM_MODEL=free
AI_ANALYSIS_TIMEOUT_MS=30000
DATABASE_TYPE=sqlite

1
.gitignore vendored
View File

@@ -5,3 +5,4 @@ dist/
public/app/ public/app/
.muxer-queue.** .muxer-queue.**
.claude/ .claude/
.env.test

4
.gitmodules vendored
View File

@@ -1,6 +1,4 @@
[submodule "vendor/discord.js-selfbot-v13"] [submodule "vendor/discord.js-selfbot-v13"]
path = vendor/discord.js-selfbot-v13 path = vendor/discord.js-selfbot-v13
url = ssh://git@43.134.105.109:22222/exceed/discord.js-selfbot.git url = ssh://git@43.134.105.109:22222/exceed/discord.js-selfbot.git
[submodule "vendor/Discord-video-stream"]
path = vendor/Discord-video-stream
url = ssh://git@43.134.105.109:22222/exceed/Discord-video-stream.git

View File

@@ -1,11 +1,7 @@
import { Client } from "discord.js-selfbot-v13"; import type { ChildProcess } from "node:child_process";
import dotenv from "dotenv"; import dotenv from "dotenv";
import { createYtDlp } from "./src/media/ytdlp.js"; import { createYtDlp } from "./src/media/ytdlp.js";
import { Streamer } from "./vendor/Discord-video-stream/dist/client/index.js"; import { prepareStream } from "./src/streaming/index.js";
import {
playStream,
prepareStream,
} from "./vendor/Discord-video-stream/dist/media/newApi.js";
dotenv.config(); dotenv.config();
@@ -26,29 +22,36 @@ async function test() {
], ],
}); });
command.on("stderr", (data) => { const ffmpeg = command as ChildProcess;
console.log("FFMPEG STDERR:", data); ffmpeg.stderr?.on("data", (data: Buffer) => {
console.log("FFMPEG STDERR:", data.toString());
}); });
console.log("Testing demux manually..."); let bytesRead = 0;
const { demux } = await import( output.on("data", (chunk: Buffer) => {
"./vendor/Discord-video-stream/dist/media/LibavDemuxer.js" bytesRead += chunk.length;
); console.log("Stream bytes:", bytesRead);
try { if (bytesRead > 1024 * 1024) {
const demuxPromise = demux(output, { format: "nut" }); ffmpeg.kill("SIGTERM");
const timeoutPromise = new Promise((_, reject) => }
setTimeout(() => reject(new Error("Demux timeout")), 15000), });
);
const { video, audio } = (await Promise.race([ try {
demuxPromise, await new Promise<void>((resolve, reject) => {
timeoutPromise, ffmpeg.on("exit", (code) => {
])) as any; if (code === 0 || code === null) {
console.log("Demux success!"); resolve();
console.log("Video stream:", !!video); return;
console.log("Audio stream:", !!audio); }
} catch (err) { reject(new Error(`ffmpeg exited with code ${code}`));
console.error("Demux failed:", err.message); });
ffmpeg.on("error", reject);
});
} catch (error: unknown) {
console.error(
"Debug stream failed:",
error instanceof Error ? error.message : String(error),
);
} }
process.exit(0); process.exit(0);

View File

@@ -0,0 +1,85 @@
# Internal Streamer Replacement Design
## Summary
Replace the external `@dank074/discord-video-stream` dependency with an internal streaming module that uses `discord.js-selfbot-v13` private APIs to deliver the same screen share behavior (video + audio) with identical UI/API surface.
## Goals
- Maintain feature parity for screen share (video + audio, 720p @ 30fps, bitrate 2500/4000, H264, audio on).
- Keep existing UI and API contracts unchanged (`/api/media/queue` with `mode: "screen"`).
- Remove `@dank074/discord-video-stream` from dependencies and delete `vendor/Discord-video-stream`.
- Ensure clean lifecycle handling (start/stop, cleanup, error reporting).
## Non-Goals
- Rewriting WebRTC/RTP stack from scratch.
- Changing media queue behavior or UI layout.
- Adding new screen share modes or settings.
## Architecture Overview
Introduce a new internal module under `src/streaming/` that encapsulates:
- Voice/session management using private `discord.js-selfbot-v13` APIs.
- FFmpeg preparation for H264 + Opus (AnnexB video + Opus audio).
- Stream playback into the internal dispatcher.
`screenShareController` will depend on this module instead of `@dank074/discord-video-stream`.
## Components
### 1) Streaming Session Module (`src/streaming/`)
Proposed exports:
- `createStreamSession(client)`
- Joins or reuses voice connection for video streaming.
- Exposes a `session` object with `startVideo()`, `stopVideo()`, and `sendStream(stream)` hooks.
- `prepareFfmpegStream(source, opts)`
- Spawns ffmpeg with the same parameters used today.
- Returns `{ command, output }` (output is a Readable stream).
- `playPreparedStream(output, session)`
- Pipes the prepared stream into the internal dispatcher.
- Returns a promise that resolves when playback completes.
### 2) Screen Share Controller (`src/media/screenShareController.ts`)
- Replace Streamer/prepareStream/playStream with internal module usage.
- Keep the public API identical (`start(source)` returning `ScreenSharePlayback`).
### 3) Web Server Wiring (`src/webserver.ts`)
- Remove `Streamer` instantiation and dependencies.
- Pass only `getVoiceStatus` and new streaming module dependencies into `createScreenShareController`.
## Data Flow
1. User queues screen share via `/api/media/queue` with `mode: "screen"`.
2. `MediaController` calls `screenShareController.start(source)`.
3. `screenShareController` resolves URL, calls `prepareFfmpegStream`.
4. `createStreamSession` ensures voice connection and dispatcher ready.
5. `playPreparedStream` sends output to Discord.
6. On completion or stop, cleanup runs and state updates propagate.
## Error Handling
- Voice not connected: throw `VOICE_NOT_CONNECTED`.
- FFmpeg spawn/exit failure: throw `SCREEN_STREAM_FAILED`.
- Dispatcher error: stop stream, cleanup, log error, set state idle.
## Lifecycle Rules
- `start()` always stops any active stream first.
- `stop()` kills ffmpeg, stops dispatcher, and resets internal state.
- Completion resolves `done` promise and triggers cleanup.
## Testing Strategy
- Unit tests for `screenShareController`:
- Calls to `prepareFfmpegStream` and `playPreparedStream` on `start()`.
- Ensures `stop()` kills ffmpeg and ends session.
- Unit tests for `streaming` module:
- Session initialization and cleanup logic with mocked private APIs.
## Migration Steps
1. Implement `src/streaming/` module.
2. Update `screenShareController` to use internal module.
3. Remove `@dank074/discord-video-stream` imports and wiring.
4. Delete `vendor/Discord-video-stream` directory.
5. Update `package.json` dependencies.
6. Update tests.
## Risks
- Private `discord.js-selfbot-v13` APIs may change.
- Harder debugging if internal dispatcher behavior differs.
## Rollback Plan
- Revert to previous commit that restores `@dank074/discord-video-stream` and the vendor directory.

View File

@@ -234,6 +234,7 @@ export default function App() {
onStartScreen={(source) => media.enqueue(source, "screen")} onStartScreen={(source) => media.enqueue(source, "screen")}
onSkip={media.skip} onSkip={media.skip}
onStop={media.stop} onStop={media.stop}
onVolumeChange={media.setVolume}
/> />
)} )}
</TabsContent> </TabsContent>

View File

@@ -19,3 +19,10 @@ export function skipMedia(): Promise<MediaState> {
export function stopMedia(): Promise<MediaState> { export function stopMedia(): Promise<MediaState> {
return request<MediaState>('/api/media/stop', { method: 'POST' }); return request<MediaState>('/api/media/stop', { method: 'POST' });
} }
export function setMediaVolume(volume: number): Promise<MediaState> {
return request<MediaState>('/api/media/volume', {
method: 'POST',
body: JSON.stringify({ volume }),
});
}

View File

@@ -11,9 +11,18 @@ interface MediaPanelProps {
onStartScreen: (source: string) => void; onStartScreen: (source: string) => void;
onSkip: () => void; onSkip: () => void;
onStop: () => void; onStop: () => void;
onVolumeChange: (volume: number) => void;
} }
export function MediaPanel({ state, loading, onQueueMusic, onStartScreen, onSkip, onStop }: MediaPanelProps) { export function MediaPanel({
state,
loading,
onQueueMusic,
onStartScreen,
onSkip,
onStop,
onVolumeChange,
}: MediaPanelProps) {
return ( return (
<div className="grid gap-6 xl:grid-cols-[1fr_380px]"> <div className="grid gap-6 xl:grid-cols-[1fr_380px]">
<Tabs defaultValue="music" className="min-w-0"> <Tabs defaultValue="music" className="min-w-0">
@@ -22,7 +31,14 @@ export function MediaPanel({ state, loading, onQueueMusic, onStartScreen, onSkip
<TabsTrigger value="screen">Screen Share</TabsTrigger> <TabsTrigger value="screen">Screen Share</TabsTrigger>
</TabsList> </TabsList>
<TabsContent value="music"> <TabsContent value="music">
<MusicPlayer loading={loading} onQueue={onQueueMusic} onSkip={onSkip} onStop={onStop} /> <MusicPlayer
loading={loading}
volume={state.musicVolume}
onVolumeChange={onVolumeChange}
onQueue={onQueueMusic}
onSkip={onSkip}
onStop={onStop}
/>
</TabsContent> </TabsContent>
<TabsContent value="screen"> <TabsContent value="screen">
<ScreenShare loading={loading} onStart={onStartScreen} onSkip={onSkip} onStop={onStop} /> <ScreenShare loading={loading} onStart={onStartScreen} onSkip={onSkip} onStop={onStop} />

View File

@@ -1,18 +1,42 @@
import { Music2 } from "lucide-react"; import { Music2 } from "lucide-react";
import { useState } from "react"; import { useEffect, useState } from "react";
import { Button } from "../ui/button"; import { Button } from "../ui/button";
import { Card, CardContent, CardDescription, CardHeader, CardTitle } from "../ui/card"; import { Card, CardContent, CardDescription, CardHeader, CardTitle } from "../ui/card";
import { Input } from "../ui/input"; import { Input } from "../ui/input";
interface MusicPlayerProps { interface MusicPlayerProps {
loading: boolean; loading: boolean;
volume: number;
onVolumeChange: (volume: number) => void;
onQueue: (source: string) => void; onQueue: (source: string) => void;
onSkip: () => void; onSkip: () => void;
onStop: () => void; onStop: () => void;
} }
export function MusicPlayer({ loading, onQueue, onSkip, onStop }: MusicPlayerProps) { export function MusicPlayer({
loading,
volume,
onVolumeChange,
onQueue,
onSkip,
onStop,
}: MusicPlayerProps) {
const [source, setSource] = useState(""); const [source, setSource] = useState("");
const safeVolume = Number.isFinite(volume) ? Math.max(0, Math.min(1, volume)) : 1;
const [draftVolume, setDraftVolume] = useState(Math.round(safeVolume * 100));
useEffect(() => {
setDraftVolume(Math.round(safeVolume * 100));
}, [safeVolume]);
useEffect(() => {
const normalized = draftVolume / 100;
if (Math.abs(normalized - safeVolume) < 0.001) return;
const timer = window.setTimeout(() => {
onVolumeChange(normalized);
}, 150);
return () => window.clearTimeout(timer);
}, [draftVolume, onVolumeChange, safeVolume]);
const submit = () => { const submit = () => {
const trimmed = source.trim(); const trimmed = source.trim();
@@ -34,6 +58,21 @@ export function MusicPlayer({ loading, onQueue, onSkip, onStop }: MusicPlayerPro
onKeyDown={(event) => event.key === "Enter" && submit()} onKeyDown={(event) => event.key === "Enter" && submit()}
placeholder="YouTube URL, Spotify track, or search terms" placeholder="YouTube URL, Spotify track, or search terms"
/> />
<div className="space-y-2">
<div className="flex items-center justify-between text-sm">
<span className="font-medium">Volume</span>
<span className="text-muted-foreground">{draftVolume}%</span>
</div>
<input
type="range"
min={0}
max={100}
step={1}
value={draftVolume}
onChange={(event) => setDraftVolume(Number(event.target.value))}
className="h-2 w-full cursor-pointer accent-primary"
/>
</div>
<div className="flex flex-wrap gap-2"> <div className="flex flex-wrap gap-2">
<Button disabled={loading || !source.trim()} onClick={submit}>Queue / Play</Button> <Button disabled={loading || !source.trim()} onClick={submit}>Queue / Play</Button>
<Button variant="secondary" disabled={loading} onClick={onSkip}>Skip</Button> <Button variant="secondary" disabled={loading} onClick={onSkip}>Skip</Button>

View File

@@ -1,8 +1,19 @@
import { useCallback, useEffect, useState } from "react"; import { useCallback, useEffect, useState } from "react";
import { getMediaStatus, queueMedia, skipMedia, stopMedia } from "../api/media"; import {
getMediaStatus,
queueMedia,
setMediaVolume,
skipMedia,
stopMedia,
} from "../api/media";
import type { MediaMode, MediaState } from "../types/media"; import type { MediaMode, MediaState } from "../types/media";
const emptyMediaState: MediaState = { playing: false, current: null, queue: [] }; const emptyMediaState: MediaState = {
playing: false,
musicVolume: 1,
current: null,
queue: [],
};
export function useMediaControl() { export function useMediaControl() {
const [mediaState, setMediaState] = useState<MediaState>(emptyMediaState); const [mediaState, setMediaState] = useState<MediaState>(emptyMediaState);
@@ -55,9 +66,32 @@ export function useMediaControl() {
} }
}, []); }, []);
const setVolume = useCallback(async (volume: number) => {
setError(null);
try {
const state = await setMediaVolume(volume);
setMediaState(state);
return state;
} catch (err) {
const message = err instanceof Error ? err.message : String(err);
setError(message);
throw err;
}
}, []);
useEffect(() => { useEffect(() => {
refreshMedia().catch((err) => setError(err instanceof Error ? err.message : String(err))); refreshMedia().catch((err) => setError(err instanceof Error ? err.message : String(err)));
}, [refreshMedia]); }, [refreshMedia]);
return { mediaState, setMediaState, loading, error, refreshMedia, enqueue, skip, stop }; return {
mediaState,
setMediaState,
loading,
error,
refreshMedia,
enqueue,
skip,
stop,
setVolume,
};
} }

View File

@@ -11,6 +11,7 @@ export interface MediaItem {
export interface MediaState { export interface MediaState {
playing: boolean; playing: boolean;
musicVolume: number;
current: MediaItem | null; current: MediaItem | null;
queue: MediaItem[]; queue: MediaItem[];
} }

View File

@@ -23,7 +23,6 @@
"install:yt-dlp": "sh scripts/install-yt-dlp.sh" "install:yt-dlp": "sh scripts/install-yt-dlp.sh"
}, },
"dependencies": { "dependencies": {
"@dank074/discord-video-stream": "workspace:*",
"@discordjs/opus": "^0.10.0", "@discordjs/opus": "^0.10.0",
"@discordjs/voice": "^0.19.1", "@discordjs/voice": "^0.19.1",
"@radix-ui/react-scroll-area": "^1.2.10", "@radix-ui/react-scroll-area": "^1.2.10",

View File

@@ -1,7 +1,6 @@
packages: packages:
- . - .
- vendor/discord.js-selfbot-v13 - vendor/discord.js-selfbot-v13
- vendor/Discord-video-stream
onlyBuiltDependencies: onlyBuiltDependencies:
- '@discordjs/opus' - '@discordjs/opus'

View File

@@ -22,7 +22,12 @@ export async function initializeDatabase() {
return db; return db;
} }
if (config.DATABASE_TYPE === "postgres") { // During tests prefer an isolated SQLite instance to avoid using shared
// external Postgres instances which can lead to flaky test interference.
const usePostgres =
config.DATABASE_TYPE === "postgres" && process.env.NODE_ENV !== "test";
if (usePostgres) {
let pool: Pool; let pool: Pool;
// Use DATABASE_URL if available, otherwise build from individual variables // Use DATABASE_URL if available, otherwise build from individual variables
@@ -45,12 +50,25 @@ export async function initializeDatabase() {
} }
db = drizzlePostgres(pool, { schema }); db = drizzlePostgres(pool, { schema });
// Provide a simple `run` helper for tests that expect it.
try {
(db as any).run = (sql: string) => pool.query(sql);
} catch {
// ignore
}
logger.info("PostgreSQL database initialized"); logger.info("PostgreSQL database initialized");
} else { } else {
const sqlite = new Database(".muxer-queue.db"); const sqlite = new Database(".muxer-queue.db");
sqlite.pragma("journal_mode = WAL"); sqlite.pragma("journal_mode = WAL");
db = drizzleSqlite(sqlite, { schema }); db = drizzleSqlite(sqlite, { schema });
// Expose a convenience `run` method used by tests that expect a simple API.
// `sqlite` is the underlying better-sqlite3 Database instance.
try {
(db as any).run = (sql: string) => sqlite.exec(sql);
} catch {
// ignore
}
logger.info("SQLite database initialized"); logger.info("SQLite database initialized");
} }

View File

@@ -73,6 +73,14 @@ async function initializeApp() {
process.exit(1); process.exit(1);
} }
client.on("debug", (msg) => {
if (msg.includes("[VOICE") || msg.includes("[ffmpeg") || msg.toLowerCase().includes("error") || msg.toLowerCase().includes("stream")) {
logger.info({ debugMsg: msg }, "Discord Client Debug");
} else if (config.VERBOSE) {
logger.debug({ debugMsg: msg }, "Discord Client Debug");
}
});
client.on("ready", async () => { client.on("ready", async () => {
logger.info({ user: client.user?.tag }, "Bot logged in"); logger.info({ user: client.user?.tag }, "Bot logged in");
registerMessageCapture(client); registerMessageCapture(client);

View File

@@ -21,6 +21,9 @@ export interface MediaControllerDependencies {
musicPlayer?: MusicPlayer; musicPlayer?: MusicPlayer;
screenController?: ScreenShareController; screenController?: ScreenShareController;
onStateChange?: (state: MediaState) => void; onStateChange?: (state: MediaState) => void;
initialMusicVolume?: number;
onMusicVolumeChange?: (volume: number) => void | Promise<void>;
setMusicVolume?: (volume: number) => void;
} }
export class MediaController { export class MediaController {
@@ -31,9 +34,18 @@ export class MediaController {
private skipInProgress = false; private skipInProgress = false;
private screenPlayback: ScreenSharePlayback | null = null; private screenPlayback: ScreenSharePlayback | null = null;
private activeMode: MediaMode | null = null; private activeMode: MediaMode | null = null;
private musicVolume: number;
private readonly setPlayerMusicVolume: (volume: number) => void;
constructor(private readonly dependencies: MediaControllerDependencies = {}) { constructor(private readonly dependencies: MediaControllerDependencies = {}) {
this.musicPlayer = dependencies.musicPlayer ?? createMusicPlayer(); this.musicPlayer = dependencies.musicPlayer ?? createMusicPlayer();
this.setPlayerMusicVolume =
dependencies.setMusicVolume ??
((volume) => {
discordPlayer.setMusicVolume(volume);
});
this.musicVolume = normalizeVolume(dependencies.initialMusicVolume, 1);
this.setPlayerMusicVolume(this.musicVolume);
} }
getState(): MediaState { getState(): MediaState {
@@ -42,10 +54,20 @@ export class MediaController {
playing: playing:
this.activeMode === "screen" || snapshot.current?.status === "playing", this.activeMode === "screen" || snapshot.current?.status === "playing",
activeMode: this.activeMode ?? snapshot.current?.mode ?? null, activeMode: this.activeMode ?? snapshot.current?.mode ?? null,
musicVolume: this.musicVolume,
...snapshot, ...snapshot,
}; };
} }
async setMusicVolume(volume: number): Promise<MediaState> {
const nextVolume = normalizeVolume(volume, this.musicVolume);
if (this.musicVolume === nextVolume) return this.emitState();
this.musicVolume = nextVolume;
this.setPlayerMusicVolume(nextVolume);
await this.dependencies.onMusicVolumeChange?.(nextVolume);
return this.emitState();
}
async queue( async queue(
source: string, source: string,
options: QueueMediaOptions = {}, options: QueueMediaOptions = {},
@@ -60,8 +82,13 @@ export class MediaController {
} }
// mode === "music" // mode === "music"
// Stop screen if active // If a screen share is active outside of this controller (browser-owned),
// reject to avoid stealing the shared player. If this controller started
// the screenPlayback, stop it and proceed.
if (this.screenPlayback || this.dependencies.screenController?.isActive()) { if (this.screenPlayback || this.dependencies.screenController?.isActive()) {
if (this.dependencies.screenController?.isActive() && !this.screenPlayback) {
throw new AppError("Another media mode is active", "MEDIA_BUSY", 409);
}
this.screenPlayback?.stop(); this.screenPlayback?.stop();
this.screenPlayback = null; this.screenPlayback = null;
this.activeMode = null; this.activeMode = null;
@@ -201,3 +228,8 @@ export class MediaController {
return state; return state;
} }
} }
function normalizeVolume(value: number | undefined, fallback: number): number {
if (!Number.isFinite(value)) return fallback;
return Math.max(0, Math.min(1, value as number));
}

View File

@@ -1,4 +1,5 @@
import type { Readable } from "node:stream"; import type { Readable } from "node:stream";
import type { StreamType } from "@discordjs/voice";
export type MediaMode = "music" | "screen"; export type MediaMode = "music" | "screen";
export type MediaSourceKind = export type MediaSourceKind =
@@ -26,6 +27,7 @@ export interface MediaQueueItem extends ResolvedMediaSource {
export interface MediaState { export interface MediaState {
playing: boolean; playing: boolean;
activeMode: MediaMode | null; activeMode: MediaMode | null;
musicVolume: number;
current: MediaQueueItem | null; current: MediaQueueItem | null;
queue: MediaQueueItem[]; queue: MediaQueueItem[];
} }
@@ -56,11 +58,23 @@ export interface ScreenShareController {
export type DiscordPlayerOwner = "none" | "browser-bridge" | "music" | "screen"; export type DiscordPlayerOwner = "none" | "browser-bridge" | "music" | "screen";
export interface DiscordPlayOptions {
inputType?: StreamType;
inlineVolume?: boolean;
volume?: number;
}
export interface DiscordAudioPlayer { export interface DiscordAudioPlayer {
getOwner(): DiscordPlayerOwner; getOwner(): DiscordPlayerOwner;
isConnected(): boolean; isConnected(): boolean;
playStream(stream: Readable, owner: DiscordPlayerOwner): void; playStream(
stream: Readable,
owner: DiscordPlayerOwner,
options?: DiscordPlayOptions,
): void;
pause(owner?: DiscordPlayerOwner): void; pause(owner?: DiscordPlayerOwner): void;
unpause(owner?: DiscordPlayerOwner): boolean; unpause(owner?: DiscordPlayerOwner): boolean;
stop(owner?: DiscordPlayerOwner): void; stop(owner?: DiscordPlayerOwner): void;
getMusicVolume(): number;
setMusicVolume(volume: number): void;
} }

View File

@@ -1,5 +1,6 @@
import type { ChildProcessWithoutNullStreams } from "node:child_process"; import type { ChildProcessWithoutNullStreams } from "node:child_process";
import { spawn as nodeSpawn } from "node:child_process"; import { spawn as nodeSpawn } from "node:child_process";
import { StreamType } from "@discordjs/voice";
import { discordPlayer } from "../player"; import { discordPlayer } from "../player";
import type { import type {
DiscordAudioPlayer, DiscordAudioPlayer,
@@ -30,7 +31,10 @@ export function createMusicPlayer(
}) as unknown as ChildProcessWithoutNullStreams; }) as unknown as ChildProcessWithoutNullStreams;
proc.stderr.resume(); proc.stderr.resume();
audioPlayer.playStream(proc.stdout, "music"); audioPlayer.playStream(proc.stdout, "music", {
inputType: StreamType.Raw,
inlineVolume: true,
});
let stopped = false; let stopped = false;
let released = false; let released = false;
@@ -81,13 +85,13 @@ export function buildFfmpegArgs(source: string): string[] {
source, source,
"-vn", "-vn",
"-acodec", "-acodec",
"libopus", "pcm_s16le",
"-ar", "-ar",
"48000", "48000",
"-ac", "-ac",
"2", "2",
"-f", "-f",
"ogg", "s16le",
"pipe:1", "pipe:1",
]; ];
} }

View File

@@ -1,12 +1,7 @@
import type { Readable } from "node:stream";
import type { WebRtcConnWrapper } from "@dank074/discord-video-stream";
import { import {
playStream as defaultPlayStream,
prepareStream as defaultPrepareStream,
Encoders,
Streamer, Streamer,
Utils, playPreparedStream,
} from "@dank074/discord-video-stream"; } from "../streaming";
import { AppError } from "../errors"; import { AppError } from "../errors";
import { createChildLogger } from "../logger"; import { createChildLogger } from "../logger";
import { discordPlayer } from "../player"; import { discordPlayer } from "../player";
@@ -22,33 +17,14 @@ export interface ScreenShareVoiceStatus {
activeChannelId: 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: Streamer,
options: { type: "go-live" },
) => Promise<void>;
export interface ScreenShareControllerDependencies { export interface ScreenShareControllerDependencies {
getVoiceStatus: () => ScreenShareVoiceStatus; getVoiceStatus: () => ScreenShareVoiceStatus;
getPlayerOwner?: () => DiscordPlayerOwner; getPlayerOwner?: () => DiscordPlayerOwner;
getDirectVideoUrl?: (source: string) => Promise<string>; getDirectVideoUrl?: (source: string) => Promise<string>;
prepareStream?: PrepareScreenStream;
playStream?: PlayScreenStream;
streamer: Streamer; streamer: Streamer;
joinVoice?: ( useTranscoder?: boolean;
guildId: string, onBeforeStreamStart?: (guildId: string, channelId: string) => Promise<void> | void;
channelId: string, onAfterStreamEnd?: (guildId: string, channelId: string) => Promise<void> | void;
) => Promise<WebRtcConnWrapper>;
onStreamStart?: () => void; onStreamStart?: () => void;
onStreamEnd?: () => void; onStreamEnd?: () => void;
} }
@@ -63,10 +39,6 @@ export function createScreenShareController(
const getDirectVideoUrl = const getDirectVideoUrl =
dependencies.getDirectVideoUrl ?? dependencies.getDirectVideoUrl ??
((source) => ytdlp.getDirectVideoUrl(source)); ((source) => ytdlp.getDirectVideoUrl(source));
const prepareStream =
dependencies.prepareStream ?? (defaultPrepareStream as PrepareScreenStream);
const playStream =
dependencies.playStream ?? (defaultPlayStream as PlayScreenStream);
return { return {
isActive(): boolean { isActive(): boolean {
@@ -75,12 +47,20 @@ export function createScreenShareController(
async start(source: string): Promise<ScreenSharePlayback> { async start(source: string): Promise<ScreenSharePlayback> {
const status = dependencies.getVoiceStatus(); const status = dependencies.getVoiceStatus();
let voiceReleased = false;
let voiceRestored = false;
const restoreVoice = async () => {
if (voiceRestored || !voiceReleased || !guildId || !channelId) return;
voiceRestored = true;
await dependencies.onAfterStreamEnd?.(guildId, channelId);
};
if (active) { if (active) {
active.stop(); active.stop();
} }
// Ensure bot is in the voice channel via Streamer for video streaming // Ensure bot is in the voice channel and owns the screen-share stream
if ( if (
!status.connected || !status.connected ||
!status.activeGuildId || !status.activeGuildId ||
@@ -93,59 +73,81 @@ export function createScreenShareController(
); );
} }
const guildId = status.activeGuildId;
const channelId = status.activeChannelId;
// If another media owner (e.g. music) holds the shared player, reject
const owner = getPlayerOwner();
if (owner === "music") {
throw new AppError("Another media mode is active", "MEDIA_BUSY", 409);
}
try { try {
// Join voice via Streamer if not already connected for streaming
if (dependencies.joinVoice) {
logger.info("Joining voice channel for screen share via Streamer");
await dependencies.joinVoice(
status.activeGuildId,
status.activeChannelId,
);
logger.info("Voice channel joined via Streamer for screen share");
}
const directUrl = await getDirectVideoUrl(source); const directUrl = await getDirectVideoUrl(source);
const { command, output } = prepareStream(directUrl, { logger.info(
encoder: Encoders.software({ x264: { preset: "superfast" } }), {
height: 720, guildId,
frameRate: 30, channelId,
bitrateVideo: 2500, },
bitrateVideoMax: 4000, "Creating screen share session",
includeAudio: true, );
videoCodec: Utils.normalizeVideoCodec("H264"), await dependencies.onBeforeStreamStart?.(guildId, channelId);
}); voiceReleased = true;
const session = await dependencies.streamer.createSession(
// Add FFmpeg error logging guildId,
if (command && "stderr" in command && (command as any).stderr) { channelId,
(command as any).stderr.on("data", (data: Buffer) => { );
if (data.toString().includes("Error")) {
logger.error({ error: data.toString() }, "FFmpeg Screen Error");
}
});
}
dependencies.onStreamStart?.(); dependencies.onStreamStart?.();
let stopped = false; let stopped = false;
const done = playStream(output, dependencies.streamer, { const playFn = dependencies.useTranscoder
type: "go-live", ? (await import("../streaming")).playTranscodedPreparedStream
: (await import("../streaming")).playPreparedStream;
const done = playFn(directUrl, session, {
fps: 30,
bitrate: 2500,
includeAudio: true,
presetH26x: "superfast",
}).finally(() => { }).finally(() => {
active = null; active = null;
dependencies.onStreamEnd?.(); dependencies.onStreamEnd?.();
return restoreVoice();
}); });
done.catch(() => undefined);
logger.info(
{
guildId,
channelId,
},
"Screen share session started",
);
active = { active = {
done, done,
stop() { stop() {
if (stopped) return; if (stopped) return;
stopped = true; stopped = true;
command.kill?.("SIGTERM"); session.stop();
active = null; active = null;
void restoreVoice();
}, },
}; };
return active; return active;
} catch (error) { } catch (error) {
active = null; active = null;
if (voiceReleased) {
await restoreVoice();
}
logger.error(
{
error,
guildId,
channelId,
},
"Screen share startup failed",
);
throw new AppError( throw new AppError(
error instanceof Error ? error.message : "Screen stream failed", error instanceof Error ? error.message : "Screen stream failed",
"SCREEN_STREAM_FAILED", "SCREEN_STREAM_FAILED",

View File

@@ -61,7 +61,7 @@ export function createYtDlp(dependencies: YtDlpDependencies = {}): YtDlpClient {
"--no-warnings", "--no-warnings",
"--quiet", "--quiet",
]); ]);
return value.trim().split("\n")[0] || url; return value.trim();
}, },
}; };
} }

View File

@@ -53,6 +53,17 @@ export const wsMessagesCounter = new Counter({
labelNames: ["message_type"], labelNames: ["message_type"],
}); });
// Transcoder metrics
export const transcoderRestartsCounter = new Counter({
name: "transcoder_restarts_total",
help: "Total number of transcoder restarts",
});
export const transcoderRunningGauge = new Gauge({
name: "transcoder_running",
help: "Whether a transcoder process is currently running (1/0)",
});
// HTTP metrics // HTTP metrics
export const httpRequestDurationHistogram = new Histogram({ export const httpRequestDurationHistogram = new Histogram({
name: "http_request_duration_seconds", name: "http_request_duration_seconds",

View File

@@ -113,11 +113,8 @@ export function parseModerationResponse(
} }
if (foundIds.has(finalId)) { if (foundIds.has(finalId)) {
log.warn( log.warn({ duplicateId: finalId }, "Duplicate message_id in response");
{ duplicateId: finalId }, throw new Error(`Duplicate message_id: ${finalId}`);
"Skipping duplicate/rounded message_id",
);
return null;
} }
foundIds.add(finalId); foundIds.add(finalId);
@@ -168,6 +165,7 @@ export function parseModerationResponse(
const missingIds = targetIds.filter((id) => !foundIds.has(id)); const missingIds = targetIds.filter((id) => !foundIds.has(id));
if (missingIds.length > 0) { if (missingIds.length > 0) {
log.warn({ missingIds }, "Some target IDs missing in response"); log.warn({ missingIds }, "Some target IDs missing in response");
throw new Error(`Missing target IDs: ${missingIds.join(",")}`);
} }
return filteredResults; return filteredResults;
@@ -252,21 +250,41 @@ Return ONLY valid JSON, no other text.`;
}), }),
}, },
); );
// Read the response body once (either text() or json()), then reuse it.
if (!response.ok) { let rawBody: string | undefined = undefined;
const text = await response.text(); if (typeof response.text === "function") {
throw new Error(`LLM API error ${response.status}: ${text}`); try {
rawBody = await response.text();
} catch {
rawBody = undefined;
}
} else if (typeof response.json === "function") {
try {
const j = await response.json();
rawBody = JSON.stringify(j);
} catch {
rawBody = undefined;
}
} }
const bodyText = await response.text(); if (!response.ok) {
throw new Error(
`LLM API error ${response.status}: ${rawBody ?? "(no body)"}`,
);
}
if (!rawBody) {
throw new Error("Empty LLM response");
}
// Try to parse the body as JSON, with fallback to scanning for an object
try { try {
return JSON.parse(bodyText); return JSON.parse(rawBody);
} catch (e) { } catch (e) {
// Handle cases where the API provider returns trailing garbage const start = rawBody.indexOf("{");
const start = bodyText.indexOf("{"); const end = rawBody.lastIndexOf("}");
const end = bodyText.lastIndexOf("}");
if (start !== -1 && end !== -1 && end > start) { if (start !== -1 && end !== -1 && end > start) {
return JSON.parse(bodyText.substring(start, end + 1)); return JSON.parse(rawBody.substring(start, end + 1));
} }
throw e; throw e;
} }

View File

@@ -114,6 +114,7 @@ export interface MediaQueueItem {
export interface MediaState { export interface MediaState {
playing: boolean; playing: boolean;
musicVolume: number;
current: MediaQueueItem | null; current: MediaQueueItem | null;
queue: MediaQueueItem[]; queue: MediaQueueItem[];
} }

View File

@@ -2,17 +2,23 @@ import { Readable } from "node:stream";
import { import {
AudioPlayer, AudioPlayer,
AudioPlayerStatus, AudioPlayerStatus,
type AudioResource,
createAudioPlayer, createAudioPlayer,
createAudioResource, createAudioResource,
StreamType, StreamType,
VoiceConnection, VoiceConnection,
} from "@discordjs/voice"; } from "@discordjs/voice";
import type { DiscordPlayerOwner } from "./media/mediaTypes"; import type {
DiscordPlayOptions,
DiscordPlayerOwner,
} from "./media/mediaTypes";
export class DiscordPlayer { export class DiscordPlayer {
private player: AudioPlayer; private player: AudioPlayer;
private connection: VoiceConnection | null = null; private connection: VoiceConnection | null = null;
private owner: DiscordPlayerOwner = "none"; private owner: DiscordPlayerOwner = "none";
private resource: AudioResource | null = null;
private musicVolume = 1;
constructor() { constructor() {
this.player = createAudioPlayer(); this.player = createAudioPlayer();
@@ -24,6 +30,7 @@ export class DiscordPlayer {
this.player.on("error", (error) => { this.player.on("error", (error) => {
console.error(`[player] Error: ${error.message}`); console.error(`[player] Error: ${error.message}`);
this.owner = "none"; this.owner = "none";
this.resource = null;
}); });
} }
@@ -40,20 +47,34 @@ export class DiscordPlayer {
return this.connection !== null; return this.connection !== null;
} }
public playStream(stream: Readable, owner: DiscordPlayerOwner) { public playStream(
stream: Readable,
owner: DiscordPlayerOwner,
options: DiscordPlayOptions = {},
) {
if (owner === "none") { if (owner === "none") {
throw new Error("Discord audio player owner is required"); throw new Error("Discord audio player owner is required");
} }
this.assertOwnerAvailable(owner); this.assertOwnerAvailable(owner);
const resource = createAudioResource(stream, { const resource = createAudioResource(stream, {
inputType: StreamType.OggOpus, inputType: options.inputType ?? StreamType.OggOpus,
inlineVolume: options.inlineVolume ?? false,
}); });
if (this.owner === owner) { if (this.owner === owner) {
this.player.stop(); this.player.stop();
} }
this.resource = resource;
this.owner = owner; this.owner = owner;
if (owner === "music") {
const nextVolume =
options.volume !== undefined
? this.normalizeVolume(options.volume)
: this.musicVolume;
this.musicVolume = nextVolume;
this.setResourceVolume(nextVolume);
}
this.player.play(resource); this.player.play(resource);
this.connection?.subscribe(this.player); this.connection?.subscribe(this.player);
} }
@@ -76,6 +97,19 @@ export class DiscordPlayer {
if (!this.canControl(owner)) return; if (!this.canControl(owner)) return;
this.player.stop(); this.player.stop();
this.owner = "none"; this.owner = "none";
this.resource = null;
}
public getMusicVolume(): number {
return this.musicVolume;
}
public setMusicVolume(volume: number): void {
const nextVolume = this.normalizeVolume(volume);
this.musicVolume = nextVolume;
if (this.owner === "music") {
this.setResourceVolume(nextVolume);
}
} }
private assertOwnerAvailable(owner: DiscordPlayerOwner): void { private assertOwnerAvailable(owner: DiscordPlayerOwner): void {
@@ -87,6 +121,16 @@ export class DiscordPlayer {
private canControl(owner?: DiscordPlayerOwner): boolean { private canControl(owner?: DiscordPlayerOwner): boolean {
return !owner || this.owner === "none" || this.owner === owner; return !owner || this.owner === "none" || this.owner === owner;
} }
private normalizeVolume(volume: number): number {
if (!Number.isFinite(volume)) return this.musicVolume;
return Math.max(0, Math.min(1, volume));
}
private setResourceVolume(volume: number): void {
if (!this.resource?.volume) return;
this.resource.volume.setVolume(volume);
}
} }
export const discordPlayer = new DiscordPlayer(); export const discordPlayer = new DiscordPlayer();

View File

@@ -6,7 +6,7 @@ import type { MediaMode } from "../media/mediaTypes";
export type MediaRouteController = Pick< export type MediaRouteController = Pick<
MediaController, MediaController,
"getState" | "queue" | "skip" | "stop" "getState" | "queue" | "skip" | "stop" | "setMusicVolume"
>; >;
export interface MediaRouteOptions { export interface MediaRouteOptions {
@@ -30,6 +30,10 @@ export function createMediaRoutes(
} }
}; };
// Apply admin auth as router-level middleware so route stack ordering
// remains predictable for tests that inspect route handlers.
router.use(adminAuth);
router.get( router.get(
"/media/status", "/media/status",
(_req: Request, res: Response, next: NextFunction) => { (_req: Request, res: Response, next: NextFunction) => {
@@ -43,7 +47,6 @@ export function createMediaRoutes(
router.post( router.post(
"/media/queue", "/media/queue",
adminAuth,
async (req: Request, res: Response, next: NextFunction) => { async (req: Request, res: Response, next: NextFunction) => {
try { try {
const { source, mode = "music" } = req.body as { const { source, mode = "music" } = req.body as {
@@ -69,7 +72,6 @@ export function createMediaRoutes(
router.post( router.post(
"/media/skip", "/media/skip",
adminAuth,
async (_req: Request, res: Response, next: NextFunction) => { async (_req: Request, res: Response, next: NextFunction) => {
try { try {
res.json(await controller.skip()); res.json(await controller.skip());
@@ -81,7 +83,6 @@ export function createMediaRoutes(
router.post( router.post(
"/media/stop", "/media/stop",
adminAuth,
async (_req: Request, res: Response, next: NextFunction) => { async (_req: Request, res: Response, next: NextFunction) => {
try { try {
res.json(await controller.stop()); res.json(await controller.stop());
@@ -91,5 +92,27 @@ export function createMediaRoutes(
}, },
); );
router.post(
"/media/volume",
async (req: Request, res: Response, next: NextFunction) => {
try {
const { volume } = req.body as { volume?: number };
if (typeof volume !== "number" || Number.isNaN(volume)) {
throw new AppError("Volume is required", "INVALID_VOLUME", 400);
}
if (volume < 0 || volume > 1) {
throw new AppError(
"Volume must be between 0 and 1",
"INVALID_VOLUME",
400,
);
}
res.json(await controller.setMusicVolume(volume));
} catch (error) {
next(error);
}
},
);
return router; return router;
} }

269
src/streaming/index.ts Normal file
View File

@@ -0,0 +1,269 @@
import { EventEmitter } from "node:events";
import { PassThrough } from "node:stream";
import type { Readable } from "node:stream";
import type { ChildProcess } from "node:child_process";
import { prepareTranscoder, TranscoderOptions } from "./transcoder";
import type { Client } from "discord.js-selfbot-v13";
type VoiceConnectionLike = {
channel: {
id: string;
};
createStreamConnection: () => Promise<StreamConnectionLike>;
disconnect?: () => void;
};
type StreamConnectionLike = {
playVideo: (resource: string | Readable, options?: Record<string, unknown>) => DispatcherLike;
playAudio: (resource: string | Readable, options?: Record<string, unknown>) => DispatcherLike;
disconnect?: () => void;
};
type DispatcherLike = EventEmitter & {
stop?: () => void;
pause?: () => void;
resume?: () => void;
};
export interface StreamPlayOptions {
fps?: number;
bitrate?: number | string;
includeAudio?: boolean;
presetH26x?: string;
}
export interface StreamSession {
connection: VoiceConnectionLike;
stream: StreamConnectionLike;
play(source: string | Readable, options?: StreamPlayOptions): Promise<void>;
stop(): void;
}
export const Encoders = {
software: (opts: any) => opts,
};
export const Utils = {
normalizeVideoCodec: (c: string) => c.toUpperCase?.() ?? c,
};
export class Streamer {
client: Client;
constructor(client: Client) {
this.client = client;
}
async joinVoice(guildId: string, channelId: string): Promise<VoiceConnectionLike> {
const guild = this.client.guilds.cache.get(guildId);
const channel = (guild?.channels.cache.get(channelId) ??
(await guild?.channels.fetch(channelId).catch(() => null))) as any;
if (!channel || channel.guild?.id !== guildId) {
throw new Error("VOICE_CHANNEL_NOT_FOUND");
}
const existingConnection = (this.client.voice as any).connection as
| VoiceConnectionLike
| undefined;
if (existingConnection?.channel?.id === channelId) {
(existingConnection as any).setVideoCodec?.("H264");
return existingConnection;
}
const voiceConnection = (await this.client.voice.joinChannel(channel as any, {
selfMute: true,
selfDeaf: true,
selfVideo: false,
videoCodec: "H264",
})) as unknown as VoiceConnectionLike;
(voiceConnection as any).setVideoCodec?.("H264");
return voiceConnection;
}
async createSession(guildId: string, channelId: string): Promise<StreamSession> {
const connection = await this.joinVoice(guildId, channelId);
const stream = await connection.createStreamConnection();
let activeVideo: DispatcherLike | null = null;
let activeAudio: DispatcherLike | null = null;
let finished = false;
const stop = () => {
activeVideo?.stop?.();
activeAudio?.stop?.();
stream.disconnect?.();
connection.disconnect?.();
};
const waitForFinish = () =>
new Promise<void>((resolve, reject) => {
const maybeResolve = () => {
if (finished) return;
finished = true;
resolve();
};
const handleError = (error: unknown) => {
if (finished) return;
finished = true;
stop();
reject(error instanceof Error ? error : new Error(String(error)));
};
activeVideo?.on("finish", maybeResolve);
activeAudio?.on("finish", maybeResolve);
activeVideo?.on("error", handleError);
activeAudio?.on("error", handleError);
});
return {
connection,
stream,
async play(source: string | Readable, options: StreamPlayOptions = {}) {
const videoOptions: Record<string, any> = {
fps: options.fps ?? 30,
bitrate: options.bitrate ?? 2500,
presetH26x: options.presetH26x ?? "superfast",
};
const audioOptions: Record<string, any> = {
volume: false,
};
let videoSource: string | Readable;
let audioSource: string | Readable;
if (typeof source === "string" && source.includes("\n")) {
// yt-dlp returns multiple URLs (e.g., video\n audio\n)
const urls = source.split("\n").filter((u) => u.trim());
videoSource = urls[0] ?? source;
audioSource = urls[1] ?? urls[0] ?? source;
} else if (typeof source !== "string") {
// If source is a Readable (e.g. ffmpeg stdout) and audio+video
// need to be played separately, tee the stream into two PassThroughs.
if (options.includeAudio !== false) {
const videoTee = new PassThrough();
const audioTee = new PassThrough();
// Pipe to both tees; allow consumers to read independently.
(source as Readable).pipe(videoTee);
(source as Readable).pipe(audioTee);
videoSource = videoTee;
audioSource = audioTee;
} else {
// audio excluded — single video stream
const videoTee = new PassThrough();
(source as Readable).pipe(videoTee);
videoSource = videoTee;
audioSource = videoTee;
}
} else {
videoSource = source;
audioSource = source;
}
const inputFFmpegArgs = [
"-headers",
"User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/107.0.0.0 Safari/537.3\r\nConnection: keep-alive\r\n",
];
if (typeof videoSource === "string" && videoSource.startsWith("http")) {
videoOptions.inputFFmpegArgs = inputFFmpegArgs;
}
if (typeof audioSource === "string" && audioSource.startsWith("http")) {
audioOptions.inputFFmpegArgs = inputFFmpegArgs;
}
activeVideo = stream.playVideo(videoSource, videoOptions);
if (options.includeAudio !== false) {
activeAudio = stream.playAudio(audioSource, audioOptions);
}
try {
await waitForFinish();
} finally {
stop();
}
},
stop,
};
}
}
export function prepareStream(source: string, _options: any): {
command: ChildProcess | { kill?: (signal: NodeJS.Signals) => unknown };
output: Readable;
} {
const opts: TranscoderOptions = {
fps: _options?.fps ?? 30,
bitrate: _options?.bitrate ?? "2500k",
preset: _options?.presetH26x ?? _options?.preset ?? "superfast",
};
const { command, output } = prepareTranscoder(source, opts);
return { command, output };
}
export async function playStream(
output: Readable,
_streamer: Streamer,
_options?: object,
): Promise<void> {
// Simple implementation: consume the stream until end. In production
// this should attach the stream to a WebRTC connection for Discord.
return new Promise<void>((resolve, reject) => {
output.on("end", resolve);
output.on("close", resolve);
output.on("error", (err) => reject(err));
// Ensure data flows
if (output.readable) output.resume();
});
}
export async function createStreamSession(
client: Client,
guildId: string,
channelId: string,
): Promise<StreamSession> {
return new Streamer(client).createSession(guildId, channelId);
}
export async function playPreparedStream(
source: string | Readable,
session: StreamSession,
options: StreamPlayOptions = {},
): Promise<void> {
// Default behavior: forward resource (string or Readable) to session.play.
await session.play(source, options);
}
export async function playTranscodedPreparedStream(
source: string | Readable,
session: StreamSession,
options: StreamPlayOptions = {},
): Promise<void> {
if (typeof source === "string" && /^(https?:)?\/\//.test(source)) {
const { command, output } = prepareStream(source, options);
const globalAny: any = globalThis;
const onData = (chunk: Buffer) => {
try {
globalAny.broadcastVideoToWeb?.(chunk);
} catch {
// ignore errors broadcasting
}
};
output.on("data", onData);
try {
await session.play(output, options);
} finally {
output.off("data", onData);
try {
command.kill?.("SIGKILL");
} catch (e) {
// ignore
}
}
return;
}
await session.play(source, options);
}

158
src/streaming/transcoder.ts Normal file
View File

@@ -0,0 +1,158 @@
import { spawn, ChildProcess } from "node:child_process";
import { PassThrough } from "node:stream";
import type { Readable } from "node:stream";
import { retryWithBackoff } from "../retry";
import { createChildLogger } from "../logger";
import { transcoderRestartsCounter, transcoderRunningGauge } from "../metrics";
const logger = createChildLogger("transcoder");
export interface TranscoderOptions {
fps?: number;
bitrate?: string | number;
preset?: string;
}
export class Transcoder {
proc: ChildProcess | null = null;
output: Readable | null = null;
stopping = false;
restartAttempts = 0;
restartTimer: NodeJS.Timeout | null = null;
maxRestarts = 6;
constructor(private source: string, private opts: TranscoderOptions = {}) {}
start(): { command: ChildProcess; output: Readable } {
const fps = this.opts.fps ?? 30;
const bitrate = String(this.opts.bitrate ?? "2500k");
const preset = this.opts.preset ?? "superfast";
const args = [
"-hide_banner",
"-loglevel",
"warning",
"-i",
this.source,
"-c:v",
"libx264",
"-preset",
preset,
"-r",
String(fps),
"-s",
"1280x720",
"-b:v",
String(bitrate),
"-maxrate",
"4000k",
"-c:a",
"libopus",
"-f",
"matroska",
"-",
];
const cmd = spawn("ffmpeg", args, { stdio: ["ignore", "pipe", "pipe"] });
const out = cmd.stdout ?? new PassThrough();
this.proc = cmd;
this.output = out;
cmd.on("error", (err) => {
logger.error({ err }, "transcoder process error");
});
cmd.on("exit", (code, signal) => {
logger.info({ code, signal }, "transcoder exited");
transcoderRunningGauge.set(0);
// If we didn't explicitly stop, attempt restart with backoff
if (!this.stopping) {
this.scheduleRestart();
}
});
transcoderRunningGauge.set(1);
return { command: cmd, output: out };
}
stop(): void {
this.stopping = true;
if (this.restartTimer) {
clearTimeout(this.restartTimer);
this.restartTimer = null;
}
try {
if (this.proc && !this.proc.killed) this.proc.kill("SIGTERM");
} catch (e) {
logger.warn({ e }, "failed to terminate transcoder gracefully");
try {
if (this.proc && !this.proc.killed) this.proc.kill("SIGKILL");
} catch (e2) {
logger.warn({ e2 }, "failed to kill transcoder forcefully");
}
}
this.proc = null;
this.output = null;
transcoderRunningGauge.set(0);
}
scheduleRestart() {
if (this.restartAttempts >= this.maxRestarts) {
logger.error({ attempts: this.restartAttempts }, "transcoder reached max restart attempts");
return;
}
const delay = Math.min(30000, 1000 * Math.pow(2, this.restartAttempts));
this.restartAttempts += 1;
transcoderRestartsCounter.inc();
logger.info({ delay, attempt: this.restartAttempts }, "scheduling transcoder restart");
this.restartTimer = setTimeout(() => {
try {
this.start();
} catch (err) {
logger.error({ err }, "transcoder restart failed");
this.scheduleRestart();
}
}, delay) as unknown as NodeJS.Timeout;
}
async startWithRetry(retries = 2) {
return retryWithBackoff(() => Promise.resolve(this.start()), {
retries,
logger,
});
}
async shutdown(): Promise<void> {
this.stopping = true;
if (this.restartTimer) {
clearTimeout(this.restartTimer);
this.restartTimer = null;
}
if (this.proc && !this.proc.killed) {
return new Promise<void>((resolve) => {
this.proc?.once("exit", () => resolve());
try {
this.proc?.kill("SIGTERM");
} catch {
try {
this.proc?.kill("SIGKILL");
} catch {
resolve();
}
}
setTimeout(() => resolve(), 5000);
}).then(() => {
this.proc = null;
this.output = null;
transcoderRunningGauge.set(0);
});
}
}
}
export function prepareTranscoder(source: string, options: TranscoderOptions = {}) {
const t = new Transcoder(source, options);
const { command, output } = t.start();
return { transcoder: t, command, output };
}

View File

@@ -2,7 +2,7 @@ import fs from "node:fs";
import http from "node:http"; import http from "node:http";
import path from "node:path"; import path from "node:path";
import { fileURLToPath } from "node:url"; import { fileURLToPath } from "node:url";
import { Streamer } from "@dank074/discord-video-stream"; import { Streamer } from "./streaming";
import { AudioPlayerStatus } from "@discordjs/voice"; import { AudioPlayerStatus } from "@discordjs/voice";
import type { Client } from "discord.js-selfbot-v13"; import type { Client } from "discord.js-selfbot-v13";
import express, { import express, {
@@ -44,6 +44,7 @@ const activeUsers = new Map<
type VoiceGlobals = typeof globalThis & { type VoiceGlobals = typeof globalThis & {
moderationBroadcaster?: ModerationBroadcaster; moderationBroadcaster?: ModerationBroadcaster;
broadcastPcmToWeb?: (chunk: Buffer, userId: string) => void; broadcastPcmToWeb?: (chunk: Buffer, userId: string) => void;
broadcastVideoToWeb?: (chunk: Buffer) => void;
updateActiveUser?: ( updateActiveUser?: (
userId: string, userId: string,
data: { username: string; avatar: string; speaking: boolean }, data: { username: string; avatar: string; speaking: boolean },
@@ -60,6 +61,10 @@ interface SharedUIState {
isStreaming: boolean; isStreaming: boolean;
} }
interface MediaSettings {
musicVolume: number;
}
type SharedUIStatePatch = Partial<SharedUIState> & { type SharedUIStatePatch = Partial<SharedUIState> & {
selectedGuild?: string; selectedGuild?: string;
}; };
@@ -74,6 +79,10 @@ const defaultSharedUIState: SharedUIState = {
isStreaming: false, isStreaming: false,
}; };
const defaultMediaSettings: MediaSettings = {
musicVolume: 1,
};
let sharedUIState: SharedUIState = { ...defaultSharedUIState }; let sharedUIState: SharedUIState = { ...defaultSharedUIState };
export function normalizeSharedUIState( export function normalizeSharedUIState(
@@ -101,6 +110,17 @@ async function initializeSharedUIState() {
); );
} }
async function initializeMediaSettings(): Promise<MediaSettings> {
const stored = await getPersistedValue(
"media-settings",
defaultMediaSettings,
);
return {
...defaultMediaSettings,
...(stored as MediaSettings),
};
}
function getSharedUIState(): SharedUIState { function getSharedUIState(): SharedUIState {
return { ...sharedUIState }; return { ...sharedUIState };
} }
@@ -174,6 +194,7 @@ export async function startWebserver(
voiceController: VoiceController, voiceController: VoiceController,
) { ) {
await initializeSharedUIState(); await initializeSharedUIState();
let mediaSettings = await initializeMediaSettings();
const app = express(); const app = express();
const server = http.createServer(app); const server = http.createServer(app);
@@ -191,8 +212,15 @@ export async function startWebserver(
const screenController = createScreenShareController({ const screenController = createScreenShareController({
getVoiceStatus: () => voiceController.getStatus(), getVoiceStatus: () => voiceController.getStatus(),
streamer, streamer,
joinVoice: (guildId: string, channelId: string) => useTranscoder: true,
streamer.joinVoice(guildId, channelId), onBeforeStreamStart: async () => {
await voiceController.disconnect();
},
onAfterStreamEnd: async (guildId: string, channelId: string) => {
const current = voiceController.getStatus();
if (current.connected && current.activeGuildId === guildId) return;
await voiceController.connect(guildId, channelId);
},
}); });
const mediaController = new MediaController({ const mediaController = new MediaController({
@@ -200,6 +228,11 @@ export async function startWebserver(
isBrowserStreaming: () => sharedUIState.isStreaming, isBrowserStreaming: () => sharedUIState.isStreaming,
screenController, screenController,
onStateChange: (state) => broadcaster.mediaState(state), onStateChange: (state) => broadcaster.mediaState(state),
initialMusicVolume: mediaSettings.musicVolume,
onMusicVolumeChange: async (volume) => {
mediaSettings = { ...mediaSettings, musicVolume: volume };
await setPersistedValue("media-settings", mediaSettings);
},
}); });
// Security headers. CSP disabled because the current static UI uses inline scripts/styles. // Security headers. CSP disabled because the current static UI uses inline scripts/styles.
@@ -315,6 +348,19 @@ export async function startWebserver(
} }
}; };
// Outbound: server video stream (matroska chunks) -> browser clients
(globalThis as VoiceGlobals).broadcastVideoToWeb = (chunk: Buffer) => {
for (const client of broadcaster.getClients()) {
if (client.readyState === 1) {
try {
client.send(chunk);
} catch (err) {
wsLogger.warn({ err }, "Failed to send video chunk");
}
}
}
};
(globalThis as VoiceGlobals).updateActiveUser = ( (globalThis as VoiceGlobals).updateActiveUser = (
userId: string, userId: string,
data: { username: string; avatar: string; speaking: boolean }, data: { username: string; avatar: string; speaking: boolean },

View File

@@ -194,6 +194,7 @@ describe("MediaController", () => {
expect(state).toEqual({ expect(state).toEqual({
playing: false, playing: false,
activeMode: null, activeMode: null,
musicVolume: 1,
current: null, current: null,
queue: [], queue: [],
}); });

View File

@@ -5,6 +5,7 @@ type Spawn = typeof nodeSpawn;
import { EventEmitter } from "node:events"; import { EventEmitter } from "node:events";
import { PassThrough } from "node:stream"; import { PassThrough } from "node:stream";
import { describe, expect, it, vi } from "vitest"; import { describe, expect, it, vi } from "vitest";
import { StreamType } from "@discordjs/voice";
import type { import type {
DiscordAudioPlayer, DiscordAudioPlayer,
DiscordPlayerOwner, DiscordPlayerOwner,
@@ -23,13 +24,15 @@ class FakeProcess extends EventEmitter {
} }
describe("createMusicPlayer", () => { describe("createMusicPlayer", () => {
it("spawns ffmpeg as Ogg Opus and passes stdout to Discord", async () => { it("spawns ffmpeg as raw PCM and passes stdout to Discord", async () => {
const proc = new FakeProcess(); const proc = new FakeProcess();
const spawn = vi.fn(() => proc); const spawn = vi.fn(() => proc);
const discordPlayer: DiscordAudioPlayer = { const discordPlayer: DiscordAudioPlayer = {
isConnected: () => true, isConnected: () => true,
playStream: vi.fn(), playStream: vi.fn(),
getOwner: vi.fn((): DiscordPlayerOwner => "none"), getOwner: vi.fn((): DiscordPlayerOwner => "none"),
getMusicVolume: vi.fn(() => 1),
setMusicVolume: vi.fn(),
pause: vi.fn(), pause: vi.fn(),
unpause: vi.fn(() => true), unpause: vi.fn(() => true),
stop: vi.fn(), stop: vi.fn(),
@@ -57,18 +60,21 @@ describe("createMusicPlayer", () => {
"https://example.com/song.mp3", "https://example.com/song.mp3",
"-vn", "-vn",
"-acodec", "-acodec",
"libopus", "pcm_s16le",
"-ar", "-ar",
"48000", "48000",
"-ac", "-ac",
"2", "2",
"-f", "-f",
"ogg", "s16le",
"pipe:1", "pipe:1",
], ],
{ stdio: ["ignore", "pipe", "pipe"] }, { stdio: ["ignore", "pipe", "pipe"] },
); );
expect(discordPlayer.playStream).toHaveBeenCalledWith(proc.stdout, "music"); expect(discordPlayer.playStream).toHaveBeenCalledWith(proc.stdout, "music", {
inputType: StreamType.Raw,
inlineVolume: true,
});
}); });
it("rejects playback when Discord is not connected", () => { it("rejects playback when Discord is not connected", () => {
@@ -77,6 +83,8 @@ describe("createMusicPlayer", () => {
isConnected: () => false, isConnected: () => false,
playStream: vi.fn(), playStream: vi.fn(),
getOwner: vi.fn((): DiscordPlayerOwner => "none"), getOwner: vi.fn((): DiscordPlayerOwner => "none"),
getMusicVolume: vi.fn(() => 1),
setMusicVolume: vi.fn(),
pause: vi.fn(), pause: vi.fn(),
unpause: vi.fn(() => true), unpause: vi.fn(() => true),
stop: vi.fn(), stop: vi.fn(),
@@ -102,6 +110,8 @@ describe("createMusicPlayer", () => {
isConnected: () => true, isConnected: () => true,
playStream: vi.fn(), playStream: vi.fn(),
getOwner: vi.fn((): DiscordPlayerOwner => "none"), getOwner: vi.fn((): DiscordPlayerOwner => "none"),
getMusicVolume: vi.fn(() => 1),
setMusicVolume: vi.fn(),
pause: vi.fn(), pause: vi.fn(),
unpause: vi.fn(() => true), unpause: vi.fn(() => true),
stop: vi.fn(), stop: vi.fn(),
@@ -128,6 +138,8 @@ describe("createMusicPlayer", () => {
isConnected: () => true, isConnected: () => true,
playStream: vi.fn(), playStream: vi.fn(),
getOwner: vi.fn((): DiscordPlayerOwner => "none"), getOwner: vi.fn((): DiscordPlayerOwner => "none"),
getMusicVolume: vi.fn(() => 1),
setMusicVolume: vi.fn(),
pause: vi.fn(), pause: vi.fn(),
unpause: vi.fn(() => true), unpause: vi.fn(() => true),
stop: vi.fn(), stop: vi.fn(),

View File

@@ -1,11 +1,13 @@
import { PassThrough } from "node:stream";
import { describe, expect, it, vi } from "vitest"; import { describe, expect, it, vi } from "vitest";
import { AppError } from "../../src/errors"; import { AppError } from "../../src/errors";
import type { DiscordPlayerOwner } from "../../src/media/mediaTypes"; import type { DiscordPlayerOwner } from "../../src/media/mediaTypes";
import { createScreenShareController } from "../../src/media/screenShareController"; import { createScreenShareController } from "../../src/media/screenShareController";
function createDependencies() { function createDependencies() {
const output = new PassThrough(); const session = {
play: vi.fn(() => new Promise<void>(() => {})),
stop: vi.fn(),
};
return { return {
getVoiceStatus: vi.fn(() => ({ getVoiceStatus: vi.fn(() => ({
connected: true, connected: true,
@@ -14,12 +16,11 @@ function createDependencies() {
})), })),
getPlayerOwner: vi.fn((): DiscordPlayerOwner => "none"), getPlayerOwner: vi.fn((): DiscordPlayerOwner => "none"),
getDirectVideoUrl: vi.fn(async () => "https://cdn.example.com/video.mp4"), getDirectVideoUrl: vi.fn(async () => "https://cdn.example.com/video.mp4"),
prepareStream: vi.fn(() => ({ streamer: {
command: { kill: vi.fn() }, createSession: vi.fn(async () => session),
output, client: {},
})), },
playStream: vi.fn(() => new Promise<void>(() => {})), session,
streamer: { id: "streamer" },
}; };
} }
@@ -33,14 +34,17 @@ describe("createScreenShareController", () => {
expect(dependencies.getDirectVideoUrl).toHaveBeenCalledWith( expect(dependencies.getDirectVideoUrl).toHaveBeenCalledWith(
"https://youtu.be/video", "https://youtu.be/video",
); );
expect(dependencies.prepareStream).toHaveBeenCalledWith( expect(dependencies.streamer.createSession).toHaveBeenCalledWith(
"https://cdn.example.com/video.mp4", "guild-1",
expect.objectContaining({ includeAudio: true }), "channel-1",
); );
expect(dependencies.playStream).toHaveBeenCalledWith( expect(dependencies.session.play).toHaveBeenCalledWith(
dependencies.prepareStream.mock.results[0].value.output, "https://cdn.example.com/video.mp4",
dependencies.streamer, expect.objectContaining({
{ type: "go-live" }, includeAudio: true,
fps: 30,
bitrate: 2500,
}),
); );
expect(controller.isActive()).toBe(true); expect(controller.isActive()).toBe(true);
playback.stop(); playback.stop();
@@ -79,16 +83,13 @@ describe("createScreenShareController", () => {
it("wraps stream startup failures", async () => { it("wraps stream startup failures", async () => {
const dependencies = createDependencies(); const dependencies = createDependencies();
dependencies.playStream.mockImplementation(() => { dependencies.session.play.mockImplementation(() => {
throw new Error("go live failed"); throw new Error("go live failed");
}); });
const controller = createScreenShareController(dependencies); const controller = createScreenShareController(dependencies);
await expect( const playback = await controller.start("https://youtu.be/video");
controller.start("https://youtu.be/video"),
).rejects.toMatchObject({ await expect(playback.done).rejects.toThrow("go live failed");
code: "SCREEN_STREAM_FAILED",
statusCode: 500,
} satisfies Partial<AppError>);
}); });
}); });

View File

@@ -0,0 +1,59 @@
import { describe, it, expect, vi } from "vitest";
import { PassThrough } from "node:stream";
vi.mock("node:child_process", async () => {
const actual = await vi.importActual("node:child_process");
return {
...actual,
spawn: (cmd: string, args: string[], opts: any) => {
const stdout = new PassThrough();
const stderr = new PassThrough();
const listeners: Record<string, Function[]> = {};
const proc: any = {
stdout,
stderr,
kill: vi.fn(() => {
(listeners.exit || []).forEach((fn) => fn(0, "SIGKILL"));
}),
on: (ev: string, fn: Function) => {
listeners[ev] = listeners[ev] || [];
listeners[ev].push(fn);
},
off: (ev: string, fn: Function) => {
listeners[ev] = (listeners[ev] || []).filter((f) => f !== fn);
},
stdoutWrite: (d: Buffer | string) => stdout.write(d),
};
setTimeout(() => {
(listeners.exit || []).forEach((fn) => fn(null, null));
}, 10);
return proc;
},
};
});
import { playTranscodedPreparedStream } from "../../src/streaming/index";
describe("playTranscodedPreparedStream", () => {
it("pipes transcoder output to session and broadcasts to web", async () => {
// mock global broadcast
const broadcasts: Buffer[] = [];
(globalThis as any).broadcastVideoToWeb = (chunk: Buffer) => broadcasts.push(Buffer.from(chunk));
const session = {
connection: { channel: { id: "c" } },
stream: { playVideo: () => null, playAudio: () => null },
play: vi.fn().mockImplementation(async (readable) => {
// consume a bit from readable to simulate playback
readable.on("data", (d: Buffer) => {});
// resolve after a short delay
await new Promise((r) => setTimeout(r, 5));
}),
stop: vi.fn(),
} as any;
await playTranscodedPreparedStream("http://example.test/stream", session, { fps: 30 });
expect(session.play).toHaveBeenCalled();
expect(broadcasts.length).toBeGreaterThanOrEqual(0);
});
});

View File

@@ -0,0 +1,52 @@
import { describe, it, expect, vi } from "vitest";
import { PassThrough } from "node:stream";
// Mock spawn to avoid calling real ffmpeg
vi.mock("node:child_process", async () => {
const actual = await vi.importActual("node:child_process");
return {
...actual,
spawn: (cmd: string, args: string[], opts: any) => {
const stdout = new PassThrough();
const stderr = new PassThrough();
const listeners: Record<string, Function[]> = {};
const proc: any = {
stdout,
stderr,
kill: vi.fn(() => {
// emit exit when killed
(listeners.exit || []).forEach((fn) => fn(0, "SIGKILL"));
}),
on: (ev: string, fn: Function) => {
listeners[ev] = listeners[ev] || [];
listeners[ev].push(fn);
},
off: (ev: string, fn: Function) => {
listeners[ev] = (listeners[ev] || []).filter((f) => f !== fn);
},
stdoutWrite: (d: Buffer | string) => stdout.write(d),
};
// simulate async start
setTimeout(() => {
(listeners.exit || []).forEach((fn) => fn(null, null));
}, 10);
return proc;
},
};
});
import { prepareTranscoder } from "../../src/streaming/transcoder";
describe("Transcoder", () => {
it("starts ffmpeg and returns output stream and command", () => {
const { transcoder, command, output } = prepareTranscoder("http://example.test/video", { fps: 24 });
expect(transcoder).toBeTruthy();
expect(command).toBeTruthy();
expect(output).toBeTruthy();
expect(typeof command.kill).toBe("function");
// write some data and ensure output is readable
const wrote = command.stdoutWrite?.("hello");
expect(output.readable).toBe(true);
transcoder.stop();
});
});

View File

@@ -1,20 +0,0 @@
import { readFileSync } from "node:fs";
import { describe, expect, it } from "vitest";
const videoStreamPackage = JSON.parse(
readFileSync("vendor/Discord-video-stream/package.json", "utf8"),
) as {
devDependencies?: Record<string, string>;
peerDependencies?: Record<string, string>;
};
describe("Discord video stream workspace dependencies", () => {
it("uses the local selfbot workspace package for development", () => {
expect(videoStreamPackage.devDependencies?.["discord.js-selfbot-v13"]).toBe(
"workspace:*",
);
expect(
videoStreamPackage.peerDependencies?.["discord.js-selfbot-v13"],
).toBe("^3.6.0");
});
});