feat: Video streaming via MusicBot with Go WebRTC sidecar

Extend MusicBot to stream video (YouTube, URLs, local files) to TS6
channels via WebRTC P2P. Uses a Go sidecar (Pion WebRTC + FFmpeg) for
media relay and a new stream signaling layer over the TS3 UDP protocol.

- Go sidecar: WebRTC peer management, RTP forwarding, FFmpeg control
- Stream signaling: setupstream, respondjoinstreamrequest, streamsignaling
- YouTube URL resolution via yt-dlp before passing to FFmpeg
- Quality presets (480p/720p/1080p) with configurable bitrate/framerate
- WebUI: Video tab in MusicBots page with live WebRTC preview player
- Viewer list with kick capability, chat commands (!stream, !stopstream, !viewers)
- TS3 UDP command fragmentation for large SDP payloads
- Docker: multi-stage sidecar build, separate container with SIDECAR_URL
- Graceful sidecar error handling (no backend crash on missing binary)

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
Clusterzx
2026-03-10 17:14:58 +01:00
co-authored by Claude Opus 4.6
parent 888f0d3088
commit 591e4ac950
22 changed files with 2421 additions and 1 deletions
+1
View File
@@ -46,4 +46,5 @@ RUN mkdir -p data /data/music && \
echo "--js-runtimes node" > /root/.config/yt-dlp/config
EXPOSE 3001
ENV SIDECAR_URL=http://ts6-sidecar:9800
CMD ["sh", "-c", "npx prisma db push --skip-generate && npx prisma db seed || true && node dist/index.js"]
+15
View File
@@ -0,0 +1,15 @@
# syntax=docker/dockerfile:1.4
FROM golang:1.22-bookworm AS build
WORKDIR /src
COPY packages/sidecar/go.mod packages/sidecar/go.sum ./
RUN go mod download
COPY packages/sidecar/main.go .
RUN CGO_ENABLED=0 go build -ldflags="-s -w" -o /sidecar main.go
FROM debian:bookworm-slim
RUN apt-get update && \
apt-get install -y --no-install-recommends ffmpeg ca-certificates && \
apt-get clean && rm -rf /var/lib/apt/lists/*
COPY --from=build /sidecar /usr/local/bin/sidecar
EXPOSE 9800
CMD ["sidecar"]
+17
View File
@@ -17,6 +17,8 @@ services:
- JWT_REFRESH_EXPIRY=7d
- FRONTEND_URL=${FRONTEND_URL:-http://localhost:3000}
- MUSIC_DIR=/data/music
- SIDECAR_URL=http://ts6-sidecar:9800
- SIDECAR_BINARY_PATH=${SIDECAR_BINARY_PATH:-}
volumes:
- backend-data:/app/packages/backend/data
- music-data:/data/music
@@ -24,6 +26,21 @@ services:
# - ./cookies.txt:/app/cookies.txt
ports:
- "3001:3001"
depends_on:
- sidecar
networks:
- ts6-network
sidecar:
build:
context: .
dockerfile: Dockerfile.sidecar
container_name: ts6-sidecar
restart: unless-stopped
environment:
- SIDECAR_PORT=9800
ports:
- "9800:9800"
networks:
- ts6-network
+12
View File
@@ -149,6 +149,8 @@ model MusicBot {
volume Int @default(50)
identityData String?
autoStart Boolean @default(false)
streamPreset String @default("720p")
sidecarPort Int @default(9800)
createdAt DateTime @default(now())
updatedAt DateTime @updatedAt
@@ -231,3 +233,13 @@ model MusicRequest {
@@unique([serverConfigId, url])
}
model StreamSession {
id Int @id @default(autoincrement())
musicBotId Int
source String
preset String @default("720p")
startedAt DateTime @default(now())
endedAt DateTime?
peakViewers Int @default(0)
}
@@ -409,6 +409,7 @@ musicBotRoutes.get('/:id/state', async (req: Request, res: Response, next) => {
shuffle: bot.queue.shuffle,
repeat: bot.queue.repeat,
isStreaming: bot.isStreaming,
videoStream: bot.videoStreamStatus,
});
} catch (err) { next(err); }
});
@@ -537,3 +538,100 @@ musicBotRoutes.post('/:id/queue/repeat', async (req: Request, res: Response, nex
res.json({ success: true, repeat: bot.queue.repeat });
} catch (err) { next(err); }
});
// === Video Streaming ===
// POST /:id/stream/start — Start video stream
musicBotRoutes.post('/:id/stream/start', async (req: Request, res: Response, next) => {
try {
const manager: VoiceBotManager = req.app.locals.voiceBotManager;
const bot = manager.getBot(parseInt(req.params.id as string));
if (!bot) throw new AppError(404, 'Music bot not found');
const { source, preset } = req.body;
if (!source) throw new AppError(400, 'source is required');
await bot.startVideoStream(source, preset);
res.json({ success: true, status: bot.videoStreamStatus });
} catch (err) { next(err); }
});
// POST /:id/stream/stop — Stop video stream
musicBotRoutes.post('/:id/stream/stop', async (req: Request, res: Response, next) => {
try {
const manager: VoiceBotManager = req.app.locals.voiceBotManager;
const bot = manager.getBot(parseInt(req.params.id as string));
if (!bot) throw new AppError(404, 'Music bot not found');
await bot.stopVideoStream();
res.json({ success: true });
} catch (err) { next(err); }
});
// POST /:id/stream/source — Change video source
musicBotRoutes.post('/:id/stream/source', async (req: Request, res: Response, next) => {
try {
const manager: VoiceBotManager = req.app.locals.voiceBotManager;
const bot = manager.getBot(parseInt(req.params.id as string));
if (!bot) throw new AppError(404, 'Music bot not found');
const { source } = req.body;
if (!source) throw new AppError(400, 'source is required');
await bot.setVideoSource(source);
res.json({ success: true });
} catch (err) { next(err); }
});
// GET /:id/stream/status — Get video stream status
musicBotRoutes.get('/:id/stream/status', async (req: Request, res: Response, next) => {
try {
const manager: VoiceBotManager = req.app.locals.voiceBotManager;
const bot = manager.getBot(parseInt(req.params.id as string));
if (!bot) throw new AppError(404, 'Music bot not found');
res.json(bot.videoStreamStatus);
} catch (err) { next(err); }
});
// DELETE /:id/stream/viewer/:clid — Kick a viewer
musicBotRoutes.delete('/:id/stream/viewer/:clid', async (req: Request, res: Response, next) => {
try {
const manager: VoiceBotManager = req.app.locals.voiceBotManager;
const bot = manager.getBot(parseInt(req.params.id as string));
if (!bot) throw new AppError(404, 'Music bot not found');
await bot.kickVideoViewer(parseInt(req.params.clid as string));
res.json({ success: true });
} catch (err) { next(err); }
});
// POST /:id/stream/webrtc/offer — Get WebRTC offer for preview player
musicBotRoutes.post('/:id/stream/webrtc/offer', async (req: Request, res: Response, next) => {
try {
const manager: VoiceBotManager = req.app.locals.voiceBotManager;
const bot = manager.getBot(parseInt(req.params.id as string));
if (!bot) throw new AppError(404, 'Music bot not found');
const offer = await bot.getWebRtcOffer();
if (!offer) throw new AppError(400, 'No active video stream');
res.json(offer);
} catch (err) { next(err); }
});
// POST /:id/stream/webrtc/answer — Set WebRTC answer from preview player
musicBotRoutes.post('/:id/stream/webrtc/answer', async (req: Request, res: Response, next) => {
try {
const manager: VoiceBotManager = req.app.locals.voiceBotManager;
const bot = manager.getBot(parseInt(req.params.id as string));
if (!bot) throw new AppError(404, 'Music bot not found');
const { sdp } = req.body;
if (!sdp) throw new AppError(400, 'sdp is required');
await bot.setWebRtcAnswer(sdp);
res.json({ success: true });
} catch (err) { next(err); }
});
// POST /:id/stream/webrtc/ice — Add ICE candidate from preview player
musicBotRoutes.post('/:id/stream/webrtc/ice', async (req: Request, res: Response, next) => {
try {
const manager: VoiceBotManager = req.app.locals.voiceBotManager;
const bot = manager.getBot(parseInt(req.params.id as string));
if (!bot) throw new AppError(404, 'Music bot not found');
const { candidate, sdpMid, sdpMLineIndex } = req.body;
await bot.addWebRtcIceCandidate(candidate, sdpMid, sdpMLineIndex ?? 0);
res.json({ success: true });
} catch (err) { next(err); }
});
@@ -10,6 +10,7 @@ const CMD_PREFIX = '!';
const MUSIC_COMMANDS = new Set([
'radio', 'play', 'stop', 'pause', 'skip', 'next', 'prev',
'vol', 'volume', 'np', 'nowplaying', 'queue', 'add',
'stream', 'stopstream', 'viewers',
]);
/**
@@ -98,6 +99,15 @@ export class MusicCommandHandler {
case 'add':
await this.handleQueue(bot, userClid, args);
break;
case 'stream':
await this.handleStream(bot, userClid, args);
break;
case 'stopstream':
await this.handleStopStream(bot, userClid);
break;
case 'viewers':
this.handleViewers(bot, userClid);
break;
}
} catch (err: any) {
console.error(`[MusicCmd] Error handling !${command}: ${err.message}`);
@@ -348,6 +358,69 @@ export class MusicCommandHandler {
this.reply(bot, userClid, `Now playing: ${artist}${np.title}`);
}
// ─── Video Streaming Commands ─────────────────────────────
private async handleStream(bot: VoiceBot, userClid: number, args: string): Promise<void> {
if (!args) {
this.reply(bot, userClid, 'Usage: !stream <url> [preset] — Presets: 480p, 720p, 1080p');
return;
}
const parts = args.split(/\s+/);
const url = parts[0];
const preset = parts[1] || undefined;
if (!url.startsWith('http://') && !url.startsWith('https://')) {
this.reply(bot, userClid, 'Please provide a valid URL.');
return;
}
if (bot.videoStreaming) {
// Change source if already streaming
try {
await bot.setVideoSource(url);
this.reply(bot, userClid, `Stream source changed to: ${url}`);
} catch (err: any) {
this.reply(bot, userClid, `Error: ${err.message}`);
}
return;
}
this.reply(bot, userClid, 'Starting video stream...');
try {
await bot.startVideoStream(url, preset);
this.reply(bot, userClid, `Video stream started: ${url}`);
} catch (err: any) {
this.reply(bot, userClid, `Failed to start stream: ${err.message}`);
}
}
private async handleStopStream(bot: VoiceBot, userClid: number): Promise<void> {
if (!bot.videoStreaming) {
this.reply(bot, userClid, 'No active video stream.');
return;
}
await bot.stopVideoStream();
this.reply(bot, userClid, 'Video stream stopped.');
}
private handleViewers(bot: VoiceBot, userClid: number): void {
const status = bot.videoStreamStatus;
if (!status.streaming) {
this.reply(bot, userClid, 'No active video stream.');
return;
}
if (status.viewers.length === 0) {
this.reply(bot, userClid, 'No viewers connected.');
return;
}
const lines = status.viewers.map((v) => {
const duration = Math.floor((Date.now() - v.joinedAt) / 1000);
return ` clid=${v.clid} (${duration}s)`;
});
this.reply(bot, userClid, `Viewers (${status.viewerCount}):\n${lines.join('\n')}`);
}
private saveMusicRequest(bot: VoiceBot, item: QueueItem): void {
if (!item.sourceUrl || !bot.currentConfig.serverConfigId) return;
this.prisma.musicRequest.upsert({
@@ -0,0 +1,84 @@
/**
* HTTP client for the Go WebRTC sidecar process.
* Manages peers, media sources, and health checks.
*/
export interface SidecarStats {
videoPort: number;
audioPort: number;
peerCount: number;
peers: Record<string, { active: boolean; state: string }>;
source: string;
}
export class SidecarClient {
private baseUrl: string;
constructor(port: number = 9800) {
this.baseUrl = `http://127.0.0.1:${port}`;
}
async waitHealthy(timeoutMs: number = 6000): Promise<void> {
const maxAttempts = Math.ceil(timeoutMs / 200);
for (let i = 0; i < maxAttempts; i++) {
try {
await this.call('GET', '/health');
return;
} catch {
await new Promise((r) => setTimeout(r, 200));
}
}
throw new Error('Sidecar health check timeout');
}
async createPeer(id: string): Promise<{ sdp: string }> {
return this.call('POST', '/peer/create', { id });
}
async setAnswer(id: string, sdp: string): Promise<void> {
await this.call('POST', '/peer/answer', { id, sdp });
}
async addIceCandidate(id: string, candidate: string, sdpMid: string, sdpMLineIndex: number): Promise<void> {
await this.call('POST', '/peer/ice', { id, candidate, sdpMid, sdpMLineIndex });
}
async closePeer(id: string): Promise<void> {
await this.call('POST', '/peer/close', { id });
}
async setSource(source: string): Promise<void> {
await this.call('POST', '/source', { source });
}
async stopSource(): Promise<void> {
await this.call('POST', '/source/stop');
}
async getStats(): Promise<SidecarStats> {
return this.call('GET', '/stats');
}
async getHealth(): Promise<{ status: string; videoPort: number; audioPort: number }> {
return this.call('GET', '/health');
}
private async call(method: string, endpoint: string, body?: any): Promise<any> {
const res = await fetch(`${this.baseUrl}${endpoint}`, {
method,
headers: body !== undefined ? { 'Content-Type': 'application/json' } : {},
body: body !== undefined ? JSON.stringify(body) : undefined,
});
if (!res.ok) {
const text = await res.text();
throw new Error(`Sidecar ${endpoint}: ${res.status} ${text}`);
}
const text = await res.text();
if (!text) return {};
try {
return JSON.parse(text);
} catch {
return {};
}
}
}
@@ -0,0 +1,112 @@
/**
* Manages the Go sidecar child process lifecycle.
* Spawns the WebRTC media relay binary and monitors it.
*/
import { spawn, type ChildProcess } from 'child_process';
import { EventEmitter } from 'events';
export interface SidecarConfig {
binaryPath: string;
port?: number;
ffmpegPath?: string;
stunServers?: string[];
videoCodec?: 'vp8' | 'vp9' | 'h264';
videoBitrate?: string;
videoResolution?: { width: number; height: number };
videoFramerate?: number;
audioBitrate?: string;
}
export class SidecarProcess extends EventEmitter {
private process: ChildProcess | null = null;
private config: SidecarConfig;
constructor(config: SidecarConfig) {
super();
this.config = config;
}
start(): void {
if (this.process) return;
const port = this.config.port ?? 9800;
const env: Record<string, string> = {
...(process.env as Record<string, string>),
SIDECAR_PORT: String(port),
};
if (this.config.ffmpegPath) env.FFMPEG_PATH = this.config.ffmpegPath;
if (this.config.stunServers?.length) env.STUN_SERVERS = this.config.stunServers.join(',');
if (this.config.videoCodec) env.VIDEO_CODEC = this.config.videoCodec;
if (this.config.videoBitrate) env.VIDEO_BITRATE = this.config.videoBitrate;
if (this.config.videoResolution) {
env.VIDEO_WIDTH = String(this.config.videoResolution.width);
env.VIDEO_HEIGHT = String(this.config.videoResolution.height);
}
if (this.config.videoFramerate) env.VIDEO_FRAMERATE = String(this.config.videoFramerate);
if (this.config.audioBitrate) env.AUDIO_BITRATE = this.config.audioBitrate;
console.log(`[Sidecar] Starting: ${this.config.binaryPath} (port ${port})`);
this.process = spawn(this.config.binaryPath, [], {
env,
stdio: ['ignore', 'pipe', 'pipe'],
});
this.process.stdout?.on('data', (data: Buffer) => {
const line = data.toString().trim();
if (line) {
console.log(`[Sidecar] ${line}`);
this.emit('stdout', line);
}
});
this.process.stderr?.on('data', (data: Buffer) => {
const line = data.toString().trim();
if (line) {
console.log(`[Sidecar] ${line}`);
this.emit('stderr', line);
}
});
this.process.on('close', (code) => {
console.log(`[Sidecar] Exited (code=${code})`);
this.process = null;
this.emit('exited', code);
});
this.process.on('error', (err) => {
console.error(`[Sidecar] Process error: ${err.message}`);
this.process = null;
// Don't re-emit as unhandled — just log and notify via 'exited'
this.emit('exited', -1);
});
}
async stop(): Promise<void> {
if (!this.process) return;
const proc = this.process;
this.process = null;
return new Promise<void>((resolve) => {
const timeout = setTimeout(() => {
try { proc.kill('SIGKILL'); } catch { /* ignore */ }
resolve();
}, 3000);
proc.once('close', () => {
clearTimeout(timeout);
resolve();
});
proc.removeAllListeners('close');
try { proc.kill('SIGTERM'); } catch { /* ignore */ }
});
}
isRunning(): boolean {
return this.process !== null;
}
}
@@ -0,0 +1,274 @@
/**
* TS6 Stream Signaling — handles stream protocol commands.
*
* Outgoing: setupstream, respondjoinstreamrequest, streamsignaling, stopstream, removeclientfromstream
* Incoming: notifystreamstarted, notifystreamstopped, notifyjoinstreamrequest,
* notifyrespondjoinstreamrequest, notifystreamsignaling,
* notifystreamclientjoined, notifystreamclientleft, notifystreaminfo
*/
import { EventEmitter } from 'events';
import type { Ts3Client } from '../tslib/client.js';
import { buildCommand } from '../tslib/commands.js';
export interface ActiveStream {
id: string;
clid: number;
name: string;
type: number;
access: number;
mode: number;
bitrate: number;
viewerLimit: number;
audio: boolean;
startedAt: number;
}
export interface SignalingMessage {
type: 'offer' | 'answer' | 'ice_candidate' | 'reconnect' | 'stream_started' | 'stream_stopped' | 'join_response' | 'unknown';
streamId?: string;
clid?: number;
sdp?: string;
sdpMid?: string;
sdpMlineIndex?: number;
candidate?: string;
isReconnect?: boolean;
raw: string;
stream?: ActiveStream;
}
export class StreamSignaling extends EventEmitter {
private client: Ts3Client;
private activeStreams: Map<string, ActiveStream> = new Map();
constructor(client: Ts3Client) {
super();
this.client = client;
this.client.on('command', (parsed: any) => this.handleCommand(parsed));
}
getActiveStreams(): Map<string, ActiveStream> {
return new Map(this.activeStreams);
}
private handleCommand(parsed: any): void {
switch (parsed.name) {
case 'notifystreamstarted':
this.handleStreamStarted(parsed);
break;
case 'notifystreamstopped':
this.handleStreamStopped(parsed);
break;
case 'notifystreamsignaling':
this.handleStreamSignaling(parsed);
break;
case 'notifyjoinstreamrequest':
this.emit('joinStreamRequest', parsed.params);
break;
case 'notifyrespondjoinstreamrequest':
this.handleJoinStreamResponse(parsed);
break;
case 'notifystreamclientjoined':
this.emit('streamClientJoined', parsed.params);
break;
case 'notifystreamclientleft':
this.emit('streamClientLeft', parsed.params);
break;
case 'notifystreaminfo':
this.handleStreamInfo(parsed);
break;
}
}
private handleStreamStarted(parsed: any): void {
const p = parsed.params;
const stream: ActiveStream = {
id: p.id || '',
clid: parseInt(p.clid) || 0,
name: p.name || '',
type: parseInt(p.type) || 0,
access: parseInt(p.access) || 0,
mode: parseInt(p.mode) || 0,
bitrate: parseInt(p.bitrate) || 0,
viewerLimit: parseInt(p.viewer_limit) || 0,
audio: p.audio === '1',
startedAt: Date.now(),
};
this.activeStreams.set(stream.id, stream);
this.emit('streamStarted', stream);
this.emit('signalingMessage', {
type: 'stream_started',
streamId: stream.id,
clid: stream.clid,
raw: JSON.stringify(parsed),
stream,
} as SignalingMessage);
}
private handleStreamStopped(parsed: any): void {
const streamId = parsed.params.id || parsed.params.stream_id || '';
const stream = this.activeStreams.get(streamId);
this.activeStreams.delete(streamId);
this.emit('streamStopped', { streamId, stream });
this.emit('signalingMessage', {
type: 'stream_stopped',
streamId,
clid: parseInt(parsed.params.clid) || stream?.clid,
raw: JSON.stringify(parsed),
stream,
} as SignalingMessage);
}
private handleStreamSignaling(parsed: any): void {
const p = parsed.params;
const dataStr = p.json || p.data || '';
if (!dataStr) return;
try {
const payload = JSON.parse(dataStr);
const cmd = payload.cmd;
const args = payload.args || {};
switch (cmd) {
case 'offer':
case 'reconnectOffer':
this.emit('signalingMessage', {
type: 'offer',
sdp: args.offer || args.sdp,
isReconnect: cmd === 'reconnectOffer',
clid: parseInt(p.clid) || undefined,
streamId: p.id || p.stream_id,
raw: dataStr,
} as SignalingMessage);
break;
case 'answer':
this.emit('signalingMessage', {
type: 'answer',
sdp: args.answer || args.sdp,
clid: parseInt(p.clid) || undefined,
streamId: p.id || p.stream_id,
raw: dataStr,
} as SignalingMessage);
break;
case 'iceCandidate':
this.emit('signalingMessage', {
type: 'ice_candidate',
candidate: args.sdp || args.candidate,
sdpMid: args.mid || args.sdp_mid,
sdpMlineIndex: args.mLine ?? args.sdp_mline_index,
clid: parseInt(p.clid) || undefined,
streamId: p.id || p.stream_id,
raw: dataStr,
} as SignalingMessage);
break;
case 'reconnect':
this.emit('signalingMessage', {
type: 'reconnect',
isReconnect: true,
clid: parseInt(p.clid) || undefined,
streamId: p.id || p.stream_id,
raw: dataStr,
} as SignalingMessage);
break;
}
} catch {
console.warn(`[StreamSignaling] Failed to parse signaling JSON: ${dataStr.substring(0, 200)}`);
}
}
private handleJoinStreamResponse(parsed: any): void {
const p = parsed.params;
const decision = parseInt(p.decision) || 0;
if (decision === 1 && p.offer) {
this.emit('signalingMessage', {
type: 'offer',
sdp: p.offer,
clid: parseInt(p.clid) || undefined,
streamId: p.id || p.stream_id,
raw: JSON.stringify(p),
} as SignalingMessage);
} else {
this.emit('signalingMessage', {
type: 'join_response',
clid: parseInt(p.clid) || undefined,
streamId: p.id || p.stream_id,
raw: JSON.stringify(p),
} as SignalingMessage);
}
}
private handleStreamInfo(parsed: any): void {
const p = parsed.params;
if (!p.id) return;
const stream: ActiveStream = {
id: p.id,
clid: parseInt(p.clid) || 0,
name: p.name || '',
type: parseInt(p.type) || 0,
access: parseInt(p.accessibility) || parseInt(p.access) || 0,
mode: parseInt(p.mode) || 0,
bitrate: parseInt(p.bitrate) || 0,
viewerLimit: parseInt(p.viewer_limit) || 0,
audio: p.audio === '1',
startedAt: Date.now(),
};
this.activeStreams.set(stream.id, stream);
this.emit('streamStarted', stream);
}
// --- Outgoing commands ---
sendSetupStream(params: {
name?: string;
type?: number;
bitrate?: number;
accessibility?: number;
mode?: number;
viewerLimit?: number;
audio?: boolean;
} = {}): void {
this.client.sendCommand(buildCommand('setupstream', {
name: params.name || 'Bot Stream',
type: String(params.type ?? 3),
bitrate: String(params.bitrate ?? 4608),
accessibility: String(params.accessibility ?? 1),
mode: String(params.mode ?? 1),
viewer_limit: String(params.viewerLimit ?? 0),
audio: params.audio === false ? '0' : '1',
}));
}
sendJoinResponse(viewerClid: number, streamId: string, accept: boolean = true, offerSdp?: string): void {
this.client.sendCommand(buildCommand('respondjoinstreamrequest', {
id: streamId,
clid: String(viewerClid),
msg: '',
offer: offerSdp || '',
decision: accept ? '1' : '0',
}));
}
sendSignaling(targetClid: number, data: { cmd: string; args: Record<string, any> }, streamId?: string): void {
this.client.sendCommand(buildCommand('streamsignaling', {
id: streamId || '',
clid: String(targetClid),
json: JSON.stringify(data),
}));
}
sendStreamStop(streamId: string): void {
this.client.sendCommand(buildCommand('stopstream', { id: streamId }));
}
sendRemoveClient(clid: number, streamId: string): void {
this.client.sendCommand(buildCommand('removeclientfromstream', { id: streamId, clid: String(clid) }));
}
registerStreamNotifications(): void {
for (const event of ['channel', 'server', 'textchannel']) {
this.client.sendCommand(buildCommand('servernotifyregister', { event }));
}
}
}
@@ -0,0 +1,36 @@
/**
* Video streaming types and quality presets
*/
export interface VideoStreamPreset {
label: string;
width: number;
height: number;
bitrate: string;
framerate: number;
}
export const STREAM_PRESETS: Record<string, VideoStreamPreset> = {
'480p': { label: '480p', width: 854, height: 480, bitrate: '1000k', framerate: 24 },
'720p': { label: '720p', width: 1280, height: 720, bitrate: '2500k', framerate: 30 },
'1080p': { label: '1080p', width: 1920, height: 1080, bitrate: '4500k', framerate: 30 },
};
export const DEFAULT_PRESET = '720p';
export interface VideoViewerInfo {
clid: number;
joinedAt: number;
iceState: string;
}
export interface VideoStreamStatus {
streaming: boolean;
streamId: string | null;
source: string | null;
preset: string;
startedAt: number | null;
viewerCount: number;
viewers: VideoViewerInfo[];
sidecar: { videoPort: number; audioPort: number } | null;
}
+23 -1
View File
@@ -254,7 +254,29 @@ export class Ts3Client extends EventEmitter {
}
sendCommand(cmd: string): void {
this.sendOutgoing(Buffer.from(cmd, "utf-8"), PacketType.Command);
const data = Buffer.from(cmd, "utf-8");
if (data.length <= MAX_OUT_CONTENT) {
this.sendOutgoing(data, PacketType.Command);
return;
}
// Fragment large commands
const chunks: Buffer[] = [];
for (let offset = 0; offset < data.length; offset += MAX_OUT_CONTENT) {
chunks.push(data.subarray(offset, Math.min(offset + MAX_OUT_CONTENT, data.length)));
}
const cmdName = cmd.split(' ')[0];
console.log(`[TS3Client] Fragmenting ${cmdName}: ${data.length} bytes → ${chunks.length} fragments`);
for (let i = 0; i < chunks.length; i++) {
const isFirst = i === 0;
const isLast = i === chunks.length - 1;
// TS3 fragmentation: FLAG_FRAGMENTED on first and last fragment
const extraFlags = (isFirst || isLast) ? FLAG_FRAGMENTED : 0;
this.sendOutgoing(chunks[i], PacketType.Command, extraFlags);
}
}
// ====== Packet Building & Sending ======
@@ -64,6 +64,9 @@ export class VoiceBotManager extends EventEmitter {
channelPassword: dbBot.channelPassword ?? undefined,
volume: dbBot.volume,
identity,
sidecarBinaryPath: process.env.SIDECAR_BINARY_PATH,
sidecarPort: (dbBot as any).sidecarPort ?? 9800,
streamPreset: (dbBot as any).streamPreset ?? '720p',
};
const bot = this.createBotInstance(config);
@@ -126,6 +129,27 @@ export class VoiceBotManager extends EventEmitter {
}
});
// Video streaming events
bot.on('videoStreamStarted', (data: any) => {
this.broadcast('music:bot:videoStreamStarted', { botId: config.id, ...data });
});
bot.on('videoStreamStopped', () => {
this.broadcast('music:bot:videoStreamStopped', { botId: config.id });
});
bot.on('videoViewerJoined', (viewer: any) => {
this.broadcast('music:bot:videoViewerJoined', { botId: config.id, viewer });
});
bot.on('videoViewerLeft', (clid: number) => {
this.broadcast('music:bot:videoViewerLeft', { botId: config.id, clid });
});
bot.on('videoSourceChanged', (source: string) => {
this.broadcast('music:bot:videoSourceChanged', { botId: config.id, source });
});
bot.on('fatalError', (msg: string) => {
console.error(`[VoiceBotManager] Bot ${config.id}: fatal error — ${msg}. No reconnect.`);
this.clearReconnect(config.id);
@@ -198,6 +222,9 @@ export class VoiceBotManager extends EventEmitter {
channelPassword: dbBot.channelPassword ?? undefined,
volume: dbBot.volume,
identity,
sidecarBinaryPath: process.env.SIDECAR_BINARY_PATH,
sidecarPort: 9800,
streamPreset: '720p',
};
const bot = this.createBotInstance(config);
+349
View File
@@ -3,6 +3,50 @@ import { Ts3Client, type Ts3ClientOptions, generateIdentity, type IdentityData,
import { AudioPipeline, FRAME_MS, BYTES_PER_FRAME } from './audio/pipeline.js';
import { PlayQueue, type QueueItem } from './playlist/queue.js';
import { fetchIcyMetadata } from './audio/icy-metadata.js';
import { StreamSignaling, type ActiveStream, type SignalingMessage } from './streaming/stream-signaling.js';
import { SidecarClient } from './streaming/sidecar-client.js';
import { SidecarProcess, type SidecarConfig } from './streaming/sidecar-process.js';
import { STREAM_PRESETS, DEFAULT_PRESET, type VideoViewerInfo, type VideoStreamStatus } from './streaming/types.js';
import { spawn } from 'child_process';
/** Resolve a YouTube/yt-dlp-compatible URL to a direct stream URL */
function resolveVideoUrl(url: string): Promise<string> {
// Only resolve YouTube and other yt-dlp-supported sites
if (!url.includes('youtube.com/') && !url.includes('youtu.be/') && !url.includes('twitch.tv/')) {
return Promise.resolve(url);
}
return new Promise((resolve, reject) => {
const proc = spawn('yt-dlp', [
'-f', 'best[ext=mp4]/best',
'--no-playlist',
'-g', // print direct URL only
url,
], { shell: false });
let stdout = '';
let stderr = '';
proc.stdout.on('data', (chunk: Buffer) => { stdout += chunk.toString(); });
proc.stderr.on('data', (chunk: Buffer) => { stderr += chunk.toString(); });
proc.on('close', (code) => {
if (code !== 0) {
return reject(new Error(`yt-dlp failed (code ${code}): ${stderr.slice(0, 200)}`));
}
// yt-dlp -g returns the direct URL(s), take the first one
const directUrl = stdout.trim().split('\n')[0];
if (!directUrl) {
return reject(new Error('yt-dlp returned no URL'));
}
console.log(`[VideoResolve] Resolved: ${url.substring(0, 60)}... → direct URL`);
resolve(directUrl);
});
proc.on('error', (err) => {
reject(new Error(`yt-dlp not found: ${err.message}`));
});
});
}
export type VoiceBotStatus = 'stopped' | 'starting' | 'connected' | 'playing' | 'paused' | 'error';
@@ -23,6 +67,9 @@ export interface VoiceBotConfig {
channelPassword?: string;
volume: number; // 0-100
identity?: IdentityData;
sidecarBinaryPath?: string;
sidecarPort?: number;
streamPreset?: string;
}
export class VoiceBot extends EventEmitter {
@@ -68,6 +115,17 @@ export class VoiceBot extends EventEmitter {
// Reconnect: distinguishes manual stop from unexpected disconnect
private _manuallyStopped: boolean = false;
// Video streaming state
private signaling: StreamSignaling | null = null;
private sidecarProc: SidecarProcess | null = null;
private sidecarHttp: SidecarClient | null = null;
private _videoStreaming: boolean = false;
private _activeStreamId: string | null = null;
private _videoSource: string | null = null;
private _videoPreset: string = DEFAULT_PRESET;
private _videoStartedAt: number | null = null;
private _viewers: Map<number, VideoViewerInfo> = new Map();
constructor(config: VoiceBotConfig) {
super();
this.config = config;
@@ -272,6 +330,10 @@ export class VoiceBot extends EventEmitter {
this.resetNickname();
this.stopPlayback();
this._nowPlaying = null;
// Stop video stream if active
if (this._videoStreaming) {
await this.stopVideoStream();
}
this.client.disconnect();
}
@@ -666,4 +728,291 @@ export class VoiceBot extends EventEmitter {
this.streamChunks = [];
this.streamChunksSize = 0;
}
// ─── Video Streaming ────────────────────────────────────────
get videoStreaming(): boolean {
return this._videoStreaming;
}
get videoStreamStatus(): VideoStreamStatus {
return {
streaming: this._videoStreaming,
streamId: this._activeStreamId,
source: this._videoSource,
preset: this._videoPreset,
startedAt: this._videoStartedAt,
viewerCount: this._viewers.size,
viewers: Array.from(this._viewers.values()),
sidecar: null,
};
}
/** Start video streaming to TS6 via WebRTC */
async startVideoStream(source: string, preset?: string): Promise<void> {
if (this._status !== 'connected' && this._status !== 'playing' && this._status !== 'paused') {
throw new Error('Bot is not connected');
}
if (this._videoStreaming) {
throw new Error('Video stream already active');
}
const sidecarBinary = this.config.sidecarBinaryPath || process.env.SIDECAR_BINARY_PATH || 'sidecar';
const sidecarPort = this.config.sidecarPort || 9800;
this._videoPreset = preset || this.config.streamPreset || DEFAULT_PRESET;
const presetConfig = STREAM_PRESETS[this._videoPreset] || STREAM_PRESETS[DEFAULT_PRESET];
// Check if sidecar URL is set (Docker mode — sidecar runs as separate container)
const sidecarUrl = process.env.SIDECAR_URL;
if (sidecarUrl) {
// Docker mode: sidecar is an external service, don't spawn it
const url = new URL(sidecarUrl);
this.sidecarHttp = new SidecarClient(parseInt(url.port) || 9800);
} else {
// Local mode: spawn sidecar binary
const sidecarConfig: SidecarConfig = {
binaryPath: sidecarBinary,
port: sidecarPort,
videoBitrate: presetConfig.bitrate,
videoResolution: { width: presetConfig.width, height: presetConfig.height },
videoFramerate: presetConfig.framerate,
};
this.sidecarProc = new SidecarProcess(sidecarConfig);
this.sidecarProc.on('exited', (code: number | null) => {
console.log(`[VoiceBot ${this.config.id}] Sidecar exited (code=${code})`);
if (this._videoStreaming) {
this._videoStreaming = false;
this._activeStreamId = null;
this._viewers.clear();
this.emit('videoStreamStopped');
this.emit('statusChange', this._status);
}
});
try {
this.sidecarProc.start();
} catch (err: any) {
this.sidecarProc = null;
throw new Error(`Failed to start sidecar: ${err.message}`);
}
this.sidecarHttp = new SidecarClient(sidecarPort);
}
// Wait for sidecar to be healthy
await this.sidecarHttp.waitHealthy();
console.log(`[VoiceBot ${this.config.id}] Sidecar ready`);
// Setup stream signaling on the TS3 client
this.signaling = new StreamSignaling(this.client);
this.setupSignalingListeners();
this.signaling.registerStreamNotifications();
// Wait for server to confirm stream
const streamPromise = new Promise<ActiveStream>((resolve, reject) => {
const timeout = setTimeout(() => reject(new Error('setupstream timeout')), 10000);
const handler = (stream: ActiveStream) => {
if (stream.clid === this.client.getClientId()) {
clearTimeout(timeout);
this.signaling!.removeListener('streamStarted', handler);
resolve(stream);
}
};
this.signaling!.on('streamStarted', handler);
});
// Send setupstream command
this.signaling.sendSetupStream({
name: `${this.config.nickname} Stream`,
type: 3,
bitrate: 4608,
accessibility: 1,
mode: 1,
viewerLimit: 0,
audio: true,
});
const stream = await streamPromise;
this._activeStreamId = stream.id;
this._videoStreaming = true;
this._videoSource = source;
this._videoStartedAt = Date.now();
// Resolve YouTube/streaming URLs via yt-dlp, then start ffmpeg
const resolvedSource = await resolveVideoUrl(source);
await this.sidecarHttp.setSource(resolvedSource);
console.log(`[VoiceBot ${this.config.id}] Video stream started: ${stream.id}, source: ${source}`);
this.emit('videoStreamStarted', { streamId: stream.id, source, preset: this._videoPreset });
this.emit('statusChange', this._status);
}
/** Stop video streaming */
async stopVideoStream(): Promise<void> {
if (!this._videoStreaming) return;
// Stop ffmpeg
try { await this.sidecarHttp?.stopSource(); } catch { /* ignore */ }
// Close all peers
for (const [clid] of this._viewers) {
try { await this.sidecarHttp?.closePeer(String(clid)); } catch { /* ignore */ }
}
this._viewers.clear();
// Stop TS6 stream
if (this.signaling && this._activeStreamId) {
this.signaling.sendStreamStop(this._activeStreamId);
}
// Stop sidecar process (only in local mode)
if (this.sidecarProc) {
await this.sidecarProc.stop();
this.sidecarProc = null;
}
this._activeStreamId = null;
this._videoSource = null;
this._videoStreaming = false;
this._videoStartedAt = null;
this.signaling = null;
console.log(`[VoiceBot ${this.config.id}] Video stream stopped`);
this.emit('videoStreamStopped');
this.emit('statusChange', this._status);
}
/** Change video source while streaming */
async setVideoSource(source: string): Promise<void> {
if (!this._videoStreaming || !this.sidecarHttp) {
throw new Error('No active video stream');
}
this._videoSource = source;
const resolvedSource = await resolveVideoUrl(source);
await this.sidecarHttp.setSource(resolvedSource);
console.log(`[VoiceBot ${this.config.id}] Video source changed: ${source}`);
this.emit('videoSourceChanged', source);
}
/** Kick a viewer from the video stream */
async kickVideoViewer(clid: number): Promise<void> {
if (!this._videoStreaming || !this.signaling || !this._activeStreamId) {
throw new Error('No active video stream');
}
try { await this.sidecarHttp?.closePeer(String(clid)); } catch { /* ignore */ }
this.signaling.sendRemoveClient(clid, this._activeStreamId);
this._viewers.delete(clid);
this.emit('videoViewerLeft', clid);
}
/** Get WebRTC offer for WebUI preview player */
async getWebRtcOffer(): Promise<{ sdp: string } | null> {
if (!this._videoStreaming || !this.sidecarHttp) return null;
return this.sidecarHttp.createPeer('webui-preview');
}
/** Set WebRTC answer from WebUI preview player */
async setWebRtcAnswer(sdp: string): Promise<void> {
if (!this.sidecarHttp) throw new Error('No sidecar');
await this.sidecarHttp.setAnswer('webui-preview', sdp);
}
/** Add ICE candidate from WebUI preview player */
async addWebRtcIceCandidate(candidate: string, sdpMid: string, sdpMLineIndex: number): Promise<void> {
if (!this.sidecarHttp) throw new Error('No sidecar');
await this.sidecarHttp.addIceCandidate('webui-preview', candidate, sdpMid, sdpMLineIndex);
}
private setupSignalingListeners(): void {
if (!this.signaling) return;
this.signaling.on('signalingMessage', (msg: SignalingMessage) => {
this.handleSignalingMessage(msg);
});
this.signaling.on('joinStreamRequest', (params: Record<string, string>) => {
const viewerClid = parseInt(params.clid) || 0;
const streamId = params.id || this._activeStreamId;
if (!streamId || !viewerClid) return;
console.log(`[VoiceBot ${this.config.id}] Viewer join request: clid=${viewerClid}`);
this.handleViewerJoin(viewerClid, streamId);
});
this.signaling.on('streamClientLeft', (params: Record<string, string>) => {
const clid = parseInt(params.clid) || 0;
if (this._viewers.has(clid)) {
console.log(`[VoiceBot ${this.config.id}] Viewer left: clid=${clid}`);
this.sidecarHttp?.closePeer(String(clid)).catch(() => {});
this._viewers.delete(clid);
this.emit('videoViewerLeft', clid);
}
});
}
private async handleSignalingMessage(msg: SignalingMessage): Promise<void> {
if (!this.sidecarHttp) return;
switch (msg.type) {
case 'answer':
if (msg.sdp && msg.clid) {
try {
await this.sidecarHttp.setAnswer(String(msg.clid), msg.sdp);
} catch (err: any) {
console.error(`[VoiceBot ${this.config.id}] setAnswer error (clid=${msg.clid}): ${err.message}`);
}
}
break;
case 'ice_candidate':
if (msg.candidate && msg.clid) {
try {
await this.sidecarHttp.addIceCandidate(
String(msg.clid),
msg.candidate,
msg.sdpMid || '0',
msg.sdpMlineIndex ?? 0
);
} catch (err: any) {
console.error(`[VoiceBot ${this.config.id}] addIceCandidate error (clid=${msg.clid}): ${err.message}`);
}
}
break;
case 'reconnect':
if (msg.clid && this._activeStreamId) {
console.log(`[VoiceBot ${this.config.id}] Reconnect from clid=${msg.clid}`);
try { await this.sidecarHttp.closePeer(String(msg.clid)); } catch { /* ignore */ }
this._viewers.delete(msg.clid);
await this.handleViewerJoin(msg.clid, this._activeStreamId);
}
break;
}
}
private async handleViewerJoin(viewerClid: number, streamId: string): Promise<void> {
if (!this.sidecarHttp || !this.signaling) return;
try {
if (this._viewers.has(viewerClid)) {
try { await this.sidecarHttp.closePeer(String(viewerClid)); } catch { /* ignore */ }
}
const result = await this.sidecarHttp.createPeer(String(viewerClid));
console.log(`[VoiceBot ${this.config.id}] Peer created for clid=${viewerClid}, sdp length=${result.sdp?.length ?? 'undefined'}, starts with: ${result.sdp?.substring(0, 30)}`);
const viewer: VideoViewerInfo = {
clid: viewerClid,
joinedAt: Date.now(),
iceState: 'new',
};
this._viewers.set(viewerClid, viewer);
this.signaling.sendJoinResponse(viewerClid, streamId, true, result.sdp);
console.log(`[VoiceBot ${this.config.id}] Sent respondjoinstreamrequest: clid=${viewerClid}, streamId=${streamId}, decision=1`);
console.log(`[VoiceBot ${this.config.id}] Viewer accepted: clid=${viewerClid} (${this._viewers.size} total)`);
this.emit('videoViewerJoined', viewer);
} catch (err: any) {
console.error(`[VoiceBot ${this.config.id}] handleViewerJoin error (clid=${viewerClid}): ${err.message}`);
this._viewers.delete(viewerClid);
}
}
}
+38
View File
@@ -131,3 +131,41 @@ export interface YouTubeUrlInfo {
type: 'video' | 'playlist';
items: YouTubeSearchResult[];
}
// === Video Streaming Types ===
export type VideoStreamPresetKey = '480p' | '720p' | '1080p';
export interface VideoStreamPreset {
label: string;
width: number;
height: number;
bitrate: string;
framerate: number;
}
export interface VideoStreamStatus {
streaming: boolean;
streamId: string | null;
source: string | null;
preset: string;
startedAt: number | null;
viewerCount: number;
viewers: VideoViewerInfo[];
sidecar: { videoPort: number; audioPort: number } | null;
}
export interface VideoViewerInfo {
clid: number;
joinedAt: number;
iceState: string;
}
export interface StartVideoStreamRequest {
source: string;
preset?: VideoStreamPresetKey;
}
export interface SetVideoSourceRequest {
source: string;
}
+14
View File
@@ -34,6 +34,20 @@ export const musicBotsApi = {
clearQueue: (id: number) => api.delete(`/music-bots/${id}/queue`).then((r) => r.data),
shuffle: (id: number, enabled: boolean) => api.post(`/music-bots/${id}/queue/shuffle`, { enabled }).then((r) => r.data),
repeat: (id: number, mode: string) => api.post(`/music-bots/${id}/queue/repeat`, { mode }).then((r) => r.data),
// Video Streaming
startStream: (id: number, source: string, preset?: string) =>
api.post(`/music-bots/${id}/stream/start`, { source, preset }).then((r) => r.data),
stopStream: (id: number) => api.post(`/music-bots/${id}/stream/stop`).then((r) => r.data),
setStreamSource: (id: number, source: string) =>
api.post(`/music-bots/${id}/stream/source`, { source }).then((r) => r.data),
streamStatus: (id: number) => api.get(`/music-bots/${id}/stream/status`).then((r) => r.data),
kickViewer: (id: number, clid: number) => api.delete(`/music-bots/${id}/stream/viewer/${clid}`).then((r) => r.data),
webrtcOffer: (id: number) => api.post(`/music-bots/${id}/stream/webrtc/offer`).then((r) => r.data),
webrtcAnswer: (id: number, sdp: string) =>
api.post(`/music-bots/${id}/stream/webrtc/answer`, { sdp }).then((r) => r.data),
webrtcIce: (id: number, candidate: string, sdpMid: string, sdpMLineIndex: number) =>
api.post(`/music-bots/${id}/stream/webrtc/ice`, { candidate, sdpMid, sdpMLineIndex }).then((r) => r.data),
};
// === Music Library API ===
@@ -0,0 +1,139 @@
/**
* WebRTC Video Player — connects to the sidecar via the backend proxy
* to display a live preview of the video stream in the WebUI.
*/
import { useRef, useEffect, useCallback, useState } from 'react';
import { musicBotsApi } from '@/api/music.api';
interface VideoPlayerProps {
botId: number;
streaming: boolean;
}
export function VideoPlayer({ botId, streaming }: VideoPlayerProps) {
const videoRef = useRef<HTMLVideoElement>(null);
const pcRef = useRef<RTCPeerConnection | null>(null);
const [connected, setConnected] = useState(false);
const [error, setError] = useState<string | null>(null);
const cleanup = useCallback(() => {
if (pcRef.current) {
pcRef.current.close();
pcRef.current = null;
}
setConnected(false);
}, []);
const connect = useCallback(async () => {
cleanup();
setError(null);
try {
// Get SDP offer from backend (which gets it from sidecar)
const { sdp: offerSdp } = await musicBotsApi.webrtcOffer(botId);
const pc = new RTCPeerConnection({
iceServers: [{ urls: 'stun:stun.l.google.com:19302' }],
});
pcRef.current = pc;
pc.ontrack = (ev) => {
if (videoRef.current && ev.streams[0]) {
videoRef.current.srcObject = ev.streams[0];
}
};
pc.oniceconnectionstatechange = () => {
const state = pc.iceConnectionState;
if (state === 'connected') {
setConnected(true);
} else if (state === 'disconnected' || state === 'failed' || state === 'closed') {
setConnected(false);
}
};
pc.onicecandidate = async (ev) => {
if (ev.candidate) {
try {
await musicBotsApi.webrtcIce(
botId,
ev.candidate.candidate,
ev.candidate.sdpMid || '0',
ev.candidate.sdpMLineIndex ?? 0
);
} catch { /* ignore ICE errors */ }
}
};
// Set remote offer
await pc.setRemoteDescription(new RTCSessionDescription({
type: 'offer',
sdp: offerSdp,
}));
// Create and send answer
const answer = await pc.createAnswer();
await pc.setLocalDescription(answer);
await musicBotsApi.webrtcAnswer(botId, answer.sdp!);
} catch (err: any) {
setError(err.message || 'Failed to connect');
cleanup();
}
}, [botId, cleanup]);
useEffect(() => {
if (streaming) {
// Small delay to ensure sidecar is ready
const timer = setTimeout(connect, 500);
return () => {
clearTimeout(timer);
cleanup();
};
} else {
cleanup();
}
}, [streaming, connect, cleanup]);
if (!streaming) {
return (
<div className="flex items-center justify-center bg-black/50 rounded-lg aspect-video">
<p className="text-muted-foreground text-sm">No active video stream</p>
</div>
);
}
return (
<div className="relative rounded-lg overflow-hidden bg-black aspect-video">
<video
ref={videoRef}
autoPlay
playsInline
muted
className="w-full h-full object-contain"
/>
{!connected && !error && (
<div className="absolute inset-0 flex items-center justify-center bg-black/70">
<p className="text-white text-sm animate-pulse">Connecting to stream...</p>
</div>
)}
{error && (
<div className="absolute inset-0 flex flex-col items-center justify-center bg-black/70 gap-2">
<p className="text-red-400 text-sm">{error}</p>
<button
onClick={connect}
className="text-xs text-blue-400 hover:text-blue-300 underline"
>
Retry
</button>
</div>
)}
{connected && (
<div className="absolute top-2 right-2 flex items-center gap-1.5 bg-black/60 px-2 py-1 rounded text-xs">
<span className="w-2 h-2 rounded-full bg-red-500 animate-pulse" />
<span className="text-white">LIVE</span>
</div>
)}
</div>
);
}
@@ -0,0 +1,246 @@
/**
* Video Stream Tab — embedded in the MusicBots page.
* Controls video streaming: source input, quality preset, start/stop,
* live WebRTC preview, and viewer management.
*/
import { useState } from 'react';
import { VideoPlayer } from './VideoPlayer';
import {
useVideoStreamStatus,
useStartVideoStream,
useStopVideoStream,
useSetStreamSource,
useKickVideoViewer,
} from '@/hooks/use-music-bots';
import { Button } from '@/components/ui/button';
import { Input } from '@/components/ui/input';
import { Label } from '@/components/ui/label';
import { Card, CardContent, CardHeader, CardTitle } from '@/components/ui/card';
import { Badge } from '@/components/ui/badge';
const PRESETS = [
{ value: '480p', label: '480p (854x480, 1 Mbps)' },
{ value: '720p', label: '720p (1280x720, 2.5 Mbps)' },
{ value: '1080p', label: '1080p (1920x1080, 4.5 Mbps)' },
];
interface VideoStreamTabProps {
botId: number;
botStatus: string;
}
export function VideoStreamTab({ botId, botStatus }: VideoStreamTabProps) {
const [sourceUrl, setSourceUrl] = useState('');
const [preset, setPreset] = useState('720p');
const { data: streamStatus } = useVideoStreamStatus(botId);
const startStream = useStartVideoStream();
const stopStream = useStopVideoStream();
const setSource = useSetStreamSource();
const kickViewer = useKickVideoViewer();
const isStreaming = streamStatus?.streaming ?? false;
const isBotConnected = botStatus === 'connected' || botStatus === 'playing' || botStatus === 'paused';
const handleStart = () => {
if (!sourceUrl.trim()) return;
startStream.mutate({ botId, source: sourceUrl.trim(), preset });
};
const handleStop = () => {
stopStream.mutate(botId);
};
const handleChangeSource = () => {
if (!sourceUrl.trim()) return;
setSource.mutate({ botId, source: sourceUrl.trim() });
};
const formatDuration = (ms: number) => {
const s = Math.floor(ms / 1000);
const h = Math.floor(s / 3600);
const m = Math.floor((s % 3600) / 60);
const sec = s % 60;
return h > 0
? `${h}:${String(m).padStart(2, '0')}:${String(sec).padStart(2, '0')}`
: `${m}:${String(sec).padStart(2, '0')}`;
};
return (
<div className="space-y-4">
{/* Stream Controls */}
<Card>
<CardHeader className="pb-3">
<div className="flex items-center justify-between">
<CardTitle className="text-base">Video Stream</CardTitle>
{isStreaming && (
<Badge variant="destructive" className="gap-1">
<span className="w-2 h-2 rounded-full bg-white animate-pulse" />
LIVE
</Badge>
)}
</div>
</CardHeader>
<CardContent className="space-y-4">
{!isBotConnected && (
<p className="text-sm text-muted-foreground">
Bot must be connected to start video streaming.
</p>
)}
{isBotConnected && (
<>
<div className="space-y-2">
<Label>Source URL</Label>
<div className="flex gap-2">
<Input
placeholder="https://youtube.com/watch?v=... or direct video URL"
value={sourceUrl}
onChange={(e) => setSourceUrl(e.target.value)}
disabled={startStream.isPending}
/>
{isStreaming ? (
<Button
onClick={handleChangeSource}
disabled={!sourceUrl.trim() || setSource.isPending}
variant="outline"
className="shrink-0"
>
Switch
</Button>
) : null}
</div>
<p className="text-xs text-muted-foreground">
YouTube, direct video URLs (MP4, HLS), or local file paths
</p>
</div>
{!isStreaming && (
<div className="space-y-2">
<Label>Quality Preset</Label>
<div className="flex gap-2">
{PRESETS.map((p) => (
<Button
key={p.value}
variant={preset === p.value ? 'default' : 'outline'}
size="sm"
onClick={() => setPreset(p.value)}
>
{p.value}
</Button>
))}
</div>
</div>
)}
<div className="flex gap-2">
{!isStreaming ? (
<Button
onClick={handleStart}
disabled={!sourceUrl.trim() || startStream.isPending}
>
{startStream.isPending ? 'Starting...' : 'Start Stream'}
</Button>
) : (
<Button
onClick={handleStop}
variant="destructive"
disabled={stopStream.isPending}
>
{stopStream.isPending ? 'Stopping...' : 'Stop Stream'}
</Button>
)}
</div>
{(startStream.isError || stopStream.isError || setSource.isError) && (
<p className="text-sm text-red-500">
{(startStream.error as any)?.message ||
(stopStream.error as any)?.message ||
(setSource.error as any)?.message}
</p>
)}
</>
)}
</CardContent>
</Card>
{/* Live Preview */}
{isBotConnected && (
<Card>
<CardHeader className="pb-3">
<CardTitle className="text-base">Live Preview</CardTitle>
</CardHeader>
<CardContent>
<VideoPlayer botId={botId} streaming={isStreaming} />
{isStreaming && streamStatus && (
<div className="mt-3 flex flex-wrap gap-3 text-xs text-muted-foreground">
<span>Preset: <strong>{streamStatus.preset}</strong></span>
{streamStatus.source && (
<span className="truncate max-w-xs">
Source: <strong>{streamStatus.source}</strong>
</span>
)}
{streamStatus.startedAt && (
<span>
Uptime: <strong>{formatDuration(Date.now() - streamStatus.startedAt)}</strong>
</span>
)}
</div>
)}
</CardContent>
</Card>
)}
{/* Viewers */}
{isStreaming && streamStatus && (
<Card>
<CardHeader className="pb-3">
<div className="flex items-center justify-between">
<CardTitle className="text-base">
Viewers ({streamStatus.viewerCount})
</CardTitle>
</div>
</CardHeader>
<CardContent>
{streamStatus.viewers.length === 0 ? (
<p className="text-sm text-muted-foreground">No viewers connected</p>
) : (
<div className="space-y-2">
{streamStatus.viewers.map((viewer: any) => {
const duration = Math.floor((Date.now() - viewer.joinedAt) / 1000);
const mins = Math.floor(duration / 60);
const secs = duration % 60;
return (
<div
key={viewer.clid}
className="flex items-center justify-between py-1.5 px-3 rounded bg-muted/50"
>
<div className="flex items-center gap-2">
<span className={`w-2 h-2 rounded-full ${
viewer.iceState === 'connected' ? 'bg-green-500' : 'bg-yellow-500'
}`} />
<span className="text-sm">Client #{viewer.clid}</span>
<span className="text-xs text-muted-foreground">
{mins > 0 ? `${mins}m ${secs}s` : `${secs}s`}
</span>
</div>
<Button
variant="ghost"
size="sm"
className="h-7 text-xs text-red-500 hover:text-red-400"
onClick={() => kickViewer.mutate({ botId, clid: viewer.clid })}
>
Kick
</Button>
</div>
);
})}
</div>
)}
</CardContent>
</Card>
)}
</div>
);
}
@@ -202,3 +202,55 @@ export function useSetRepeat() {
onSuccess: (_, { botId }) => qc.invalidateQueries({ queryKey: ['music-bot-state', botId] }),
});
}
// === Video Streaming Hooks ===
export function useVideoStreamStatus(botId: number | null) {
return useQuery({
queryKey: ['video-stream-status', botId],
queryFn: () => musicBotsApi.streamStatus(botId!),
enabled: !!botId,
refetchInterval: 2000,
});
}
export function useStartVideoStream() {
const qc = useQueryClient();
return useMutation({
mutationFn: ({ botId, source, preset }: { botId: number; source: string; preset?: string }) =>
musicBotsApi.startStream(botId, source, preset),
onSuccess: (_, { botId }) => {
qc.invalidateQueries({ queryKey: ['video-stream-status', botId] });
qc.invalidateQueries({ queryKey: ['music-bot-state', botId] });
},
});
}
export function useStopVideoStream() {
const qc = useQueryClient();
return useMutation({
mutationFn: (botId: number) => musicBotsApi.stopStream(botId),
onSuccess: (_, botId) => {
qc.invalidateQueries({ queryKey: ['video-stream-status', botId] });
qc.invalidateQueries({ queryKey: ['music-bot-state', botId] });
},
});
}
export function useSetStreamSource() {
const qc = useQueryClient();
return useMutation({
mutationFn: ({ botId, source }: { botId: number; source: string }) =>
musicBotsApi.setStreamSource(botId, source),
onSuccess: (_, { botId }) => qc.invalidateQueries({ queryKey: ['video-stream-status', botId] }),
});
}
export function useKickVideoViewer() {
const qc = useQueryClient();
return useMutation({
mutationFn: ({ botId, clid }: { botId: number; clid: number }) =>
musicBotsApi.kickViewer(botId, clid),
onSuccess: (_, { botId }) => qc.invalidateQueries({ queryKey: ['video-stream-status', botId] }),
});
}
+58
View File
@@ -34,7 +34,9 @@ import {
Volume2, VolumeX, Upload, Search, Download, ListMusic, Shuffle,
Repeat, Repeat1, Power, PowerOff, RefreshCw, Pencil, X, Loader2,
Youtube, FileAudio, Link, GripVertical, Music2, Radio, Clock,
Video,
} from 'lucide-react';
import { VideoStreamTab } from '@/components/video/VideoStreamTab';
import { toast } from 'sonner';
import { formatBytes } from '@/lib/utils';
import type { MusicBotSummary, PlaybackState, SongInfo, PlaylistSummary, PlaylistDetail, YouTubeSearchResult, RadioStationInfo, RadioPreset } from '@ts6/common';
@@ -1431,6 +1433,60 @@ function RadioTab() {
);
}
// ─── Video Streaming Tab ─────────────────────────────────────────────────────
function VideoTab() {
const { data } = useMusicBots();
const bots = Array.isArray(data) ? data : [];
const [selectedBotId, setSelectedBotId] = useState<number | null>(null);
// Auto-select first running bot
const runningBots = bots.filter((b: MusicBotSummary) => b.status !== 'stopped' && b.status !== 'error');
useEffect(() => {
if (!selectedBotId && runningBots.length > 0) {
setSelectedBotId(runningBots[0].id);
}
}, [runningBots, selectedBotId]);
const selectedBot = bots.find((b: MusicBotSummary) => b.id === selectedBotId);
return (
<div className="space-y-4">
{bots.length === 0 ? (
<EmptyState icon={Video} title="No bots available" description="Create a music bot first, then use it for video streaming." />
) : (
<>
{/* Bot selector */}
<div className="flex items-center gap-3">
<Label className="shrink-0">Select Bot:</Label>
<Select
value={selectedBotId ? String(selectedBotId) : ''}
onValueChange={(v) => setSelectedBotId(parseInt(v))}
>
<SelectTrigger className="w-64">
<SelectValue placeholder="Choose a bot..." />
</SelectTrigger>
<SelectContent>
{bots.map((b: MusicBotSummary) => (
<SelectItem key={b.id} value={String(b.id)}>
{b.name} — {b.status}
</SelectItem>
))}
</SelectContent>
</Select>
</div>
{selectedBot ? (
<VideoStreamTab botId={selectedBot.id} botStatus={selectedBot.status} />
) : (
<p className="text-sm text-muted-foreground">Select a bot to manage video streaming.</p>
)}
</>
)}
</div>
);
}
// ─── Main Page ───────────────────────────────────────────────────────────────
export default function MusicBots() {
@@ -1446,12 +1502,14 @@ export default function MusicBots() {
<Tabs defaultValue="bots" className="space-y-4">
<TabsList>
<TabsTrigger value="bots"><Music2 className="h-3.5 w-3.5 mr-1.5" /> Bots</TabsTrigger>
<TabsTrigger value="video"><Video className="h-3.5 w-3.5 mr-1.5" /> Video</TabsTrigger>
<TabsTrigger value="library"><FileAudio className="h-3.5 w-3.5 mr-1.5" /> Library</TabsTrigger>
<TabsTrigger value="playlists"><ListMusic className="h-3.5 w-3.5 mr-1.5" /> Playlists</TabsTrigger>
<TabsTrigger value="radio"><Radio className="h-3.5 w-3.5 mr-1.5" /> Radio</TabsTrigger>
</TabsList>
<TabsContent value="bots"><BotsTab /></TabsContent>
<TabsContent value="video"><VideoTab /></TabsContent>
<TabsContent value="library"><LibraryTab /></TabsContent>
<TabsContent value="playlists"><PlaylistsTab /></TabsContent>
<TabsContent value="radio"><RadioTab /></TabsContent>
+30
View File
@@ -0,0 +1,30 @@
module ts6-media-sidecar
go 1.22
require (
github.com/pion/interceptor v0.1.37
github.com/pion/rtp v1.8.9
github.com/pion/webrtc/v4 v4.0.5
)
require (
github.com/google/uuid v1.6.0 // indirect
github.com/pion/datachannel v1.5.9 // indirect
github.com/pion/dtls/v3 v3.0.4 // indirect
github.com/pion/ice/v4 v4.0.3 // indirect
github.com/pion/logging v0.2.2 // indirect
github.com/pion/mdns/v2 v2.0.7 // indirect
github.com/pion/randutil v0.1.0 // indirect
github.com/pion/rtcp v1.2.14 // indirect
github.com/pion/sctp v1.8.34 // indirect
github.com/pion/sdp/v3 v3.0.9 // indirect
github.com/pion/srtp/v3 v3.0.4 // indirect
github.com/pion/stun/v3 v3.0.0 // indirect
github.com/pion/transport/v3 v3.0.7 // indirect
github.com/pion/turn/v4 v4.0.0 // indirect
github.com/wlynxg/anet v0.0.5 // indirect
golang.org/x/crypto v0.29.0 // indirect
golang.org/x/net v0.31.0 // indirect
golang.org/x/sys v0.27.0 // indirect
)
+61
View File
@@ -0,0 +1,61 @@
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/pion/datachannel v1.5.9 h1:LpIWAOYPyDrXtU+BW7X0Yt/vGtYxtXQ8ql7dFfYUVZA=
github.com/pion/datachannel v1.5.9/go.mod h1:kDUuk4CU4Uxp82NH4LQZbISULkX/HtzKa4P7ldf9izE=
github.com/pion/dtls/v3 v3.0.4 h1:44CZekewMzfrn9pmGrj5BNnTMDCFwr+6sLH+cCuLM7U=
github.com/pion/dtls/v3 v3.0.4/go.mod h1:R373CsjxWqNPf6MEkfdy3aSe9niZvL/JaKlGeFphtMg=
github.com/pion/ice/v4 v4.0.3 h1:9s5rI1WKzF5DRqhJ+Id8bls/8PzM7mau0mj1WZb4IXE=
github.com/pion/ice/v4 v4.0.3/go.mod h1:VfHy0beAZ5loDT7BmJ2LtMtC4dbawIkkkejHPRZNB3Y=
github.com/pion/interceptor v0.1.37 h1:aRA8Zpab/wE7/c0O3fh1PqY0AJI3fCSEM5lRWJVorwI=
github.com/pion/interceptor v0.1.37/go.mod h1:JzxbJ4umVTlZAf+/utHzNesY8tmRkM2lVmkS82TTj8Y=
github.com/pion/logging v0.2.2 h1:M9+AIj/+pxNsDfAT64+MAVgJO0rsyLnoJKCqf//DoeY=
github.com/pion/logging v0.2.2/go.mod h1:k0/tDVsRCX2Mb2ZEmTqNa7CWsQPc+YYCB7Q+5pahoms=
github.com/pion/mdns/v2 v2.0.7 h1:c9kM8ewCgjslaAmicYMFQIde2H9/lrZpjBkN8VwoVtM=
github.com/pion/mdns/v2 v2.0.7/go.mod h1:vAdSYNAT0Jy3Ru0zl2YiW3Rm/fJCwIeM0nToenfOJKA=
github.com/pion/randutil v0.1.0 h1:CFG1UdESneORglEsnimhUjf33Rwjubwj6xfiOXBa3mA=
github.com/pion/randutil v0.1.0/go.mod h1:XcJrSMMbbMRhASFVOlj/5hQial/Y8oH/HVo7TBZq+j8=
github.com/pion/rtcp v1.2.14 h1:KCkGV3vJ+4DAJmvP0vaQShsb0xkRfWkO540Gy102KyE=
github.com/pion/rtcp v1.2.14/go.mod h1:sn6qjxvnwyAkkPzPULIbVqSKI5Dv54Rv7VG0kNxh9L4=
github.com/pion/rtp v1.8.9 h1:E2HX740TZKaqdcPmf4pw6ZZuG8u5RlMMt+l3dxeu6Wk=
github.com/pion/rtp v1.8.9/go.mod h1:pBGHaFt/yW7bf1jjWAoUjpSNoDnw98KTMg+jWWvziqU=
github.com/pion/sctp v1.8.34 h1:rCuD3m53i0oGxCSp7FLQKvqVx0Nf5AUAHhMRXTTQjBc=
github.com/pion/sctp v1.8.34/go.mod h1:yWkCClkXlzVW7BXfI2PjrUGBwUI0CjXJBkhLt+sdo4U=
github.com/pion/sdp/v3 v3.0.9 h1:pX++dCHoHUwq43kuwf3PyJfHlwIj4hXA7Vrifiq0IJY=
github.com/pion/sdp/v3 v3.0.9/go.mod h1:B5xmvENq5IXJimIO4zfp6LAe1fD9N+kFv+V/1lOdz8M=
github.com/pion/srtp/v3 v3.0.4 h1:2Z6vDVxzrX3UHEgrUyIGM4rRouoC7v+NiF1IHtp9B5M=
github.com/pion/srtp/v3 v3.0.4/go.mod h1:1Jx3FwDoxpRaTh1oRV8A/6G1BnFL+QI82eK4ms8EEJQ=
github.com/pion/stun/v3 v3.0.0 h1:4h1gwhWLWuZWOJIJR9s2ferRO+W3zA/b6ijOI6mKzUw=
github.com/pion/stun/v3 v3.0.0/go.mod h1:HvCN8txt8mwi4FBvS3EmDghW6aQJ24T+y+1TKjB5jyU=
github.com/pion/transport/v3 v3.0.7 h1:iRbMH05BzSNwhILHoBoAPxoB9xQgOaJk+591KC9P1o0=
github.com/pion/transport/v3 v3.0.7/go.mod h1:YleKiTZ4vqNxVwh77Z0zytYi7rXHl7j6uPLGhhz9rwo=
github.com/pion/turn/v4 v4.0.0 h1:qxplo3Rxa9Yg1xXDxxH8xaqcyGUtbHYw4QSCvmFWvhM=
github.com/pion/turn/v4 v4.0.0/go.mod h1:MuPDkm15nYSklKpN8vWJ9W2M0PlyQZqYt1McGuxG7mA=
github.com/pion/webrtc/v4 v4.0.5 h1:8cVPojcv3cQTwVga2vF1rzCNvkiEimnYdCCG7yF317I=
github.com/pion/webrtc/v4 v4.0.5/go.mod h1:LvP8Np5b/sM0uyJIcUPvJcCvhtjHxJwzh2H2PYzE6cQ=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw=
github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo=
github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA=
github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU=
github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo=
github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY=
github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA=
github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY=
github.com/wlynxg/anet v0.0.5 h1:J3VJGi1gvo0JwZ/P1/Yc/8p63SoW98B5dHkYDmpgvvU=
github.com/wlynxg/anet v0.0.5/go.mod h1:eay5PRQr7fIVAMbTbchTnO9gG65Hg/uYGdc7mguHxoA=
golang.org/x/crypto v0.29.0 h1:L5SG1JTTXupVV3n6sUqMTeWbjAyfPwoda2DLX8J8FrQ=
golang.org/x/crypto v0.29.0/go.mod h1:+F4F4N5hv6v38hfeYwTdx20oUvLLc+QfrE9Ax9HtgRg=
golang.org/x/net v0.31.0 h1:68CPQngjLL0r2AlUKiSxtQFKvzRVbnzLwMUn5SzcLHo=
golang.org/x/net v0.31.0/go.mod h1:P4fl1q7dY2hnZFxEk4pPSkDHF+QqjitcnDjUQyMM+pM=
golang.org/x/sys v0.27.0 h1:wBqf8DvsY9Y/2P8gAfPDEYNuS30J4lPHJxXSb/nJZ+s=
golang.org/x/sys v0.27.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+662
View File
@@ -0,0 +1,662 @@
package main
import (
"encoding/json"
"fmt"
"log"
"net"
"net/http"
"os"
"os/exec"
"os/signal"
"strconv"
"strings"
"sync"
"syscall"
"time"
"github.com/pion/interceptor"
"github.com/pion/interceptor/pkg/intervalpli"
"github.com/pion/rtp"
"github.com/pion/webrtc/v4"
)
var defaultStunServers = []string{
"stun:49.13.204.141:3478",
"stun:176.58.93.154:3478",
"stun:185.40.234.113:3478",
"stun:68.183.90.120:3478",
"stun:45.159.97.233:3478",
"stun:172.105.166.103:3478",
"stun:172.237.28.183:3478",
"stun:208.72.155.133:3478",
"stun:stun.l.google.com:19302",
}
func getStunServers() []string {
if env := os.Getenv("STUN_SERVERS"); env != "" {
return strings.Split(env, ",")
}
return defaultStunServers
}
func envOrDefault(key, def string) string {
if v := os.Getenv(key); v != "" {
return v
}
return def
}
func envIntOrDefault(key string, def int) int {
if v := os.Getenv(key); v != "" {
if n, err := strconv.Atoi(v); err == nil {
return n
}
}
return def
}
func getFfmpegPath() string {
return envOrDefault("FFMPEG_PATH", "ffmpeg")
}
type Peer struct {
ID string
PC *webrtc.PeerConnection
VideoTrack *webrtc.TrackLocalStaticRTP
AudioTrack *webrtc.TrackLocalStaticRTP
Active bool
mu sync.Mutex
}
type Sidecar struct {
peers map[string]*Peer
peersLock sync.RWMutex
videoPort int
audioPort int
videoConn *net.UDPConn
audioConn *net.UDPConn
ffmpeg *exec.Cmd
ffmpegLock sync.Mutex
source string
running bool
videoTsFirst uint32
videoTsGot bool
videoWallFirst int64
audioTsFirst uint32
audioTsGot bool
audioWallFirst int64
tsLock sync.Mutex
}
func NewSidecar() *Sidecar {
return &Sidecar{
peers: make(map[string]*Peer),
}
}
func (s *Sidecar) resetTimestamps() {
s.tsLock.Lock()
s.videoTsGot = false
s.audioTsGot = false
s.tsLock.Unlock()
log.Printf("[SYNC] Timestamp normalization reset")
}
func (s *Sidecar) StartRTP() error {
var err error
s.videoConn, err = net.ListenUDP("udp4", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1), Port: 0})
if err != nil {
return fmt.Errorf("bind video UDP: %w", err)
}
s.videoPort = s.videoConn.LocalAddr().(*net.UDPAddr).Port
s.audioConn, err = net.ListenUDP("udp4", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1), Port: 0})
if err != nil {
return fmt.Errorf("bind audio UDP: %w", err)
}
s.audioPort = s.audioConn.LocalAddr().(*net.UDPAddr).Port
log.Printf("[RTP] Video port: %d, Audio port: %d", s.videoPort, s.audioPort)
s.running = true
go s.readVideoRTP()
go s.readAudioRTP()
return nil
}
func (s *Sidecar) readVideoRTP() {
buf := make([]byte, 1500)
pkt := &rtp.Packet{}
count := 0
for s.running {
n, err := s.videoConn.Read(buf)
if err != nil {
if s.running {
log.Printf("[RTP] Video read error: %v", err)
}
return
}
if err := pkt.Unmarshal(buf[:n]); err != nil {
continue
}
s.tsLock.Lock()
if !s.videoTsGot {
s.videoTsFirst = pkt.Timestamp
s.videoWallFirst = time.Now().UnixNano()
s.videoTsGot = true
log.Printf("[VIDEO] First ts=%d wall=%d", pkt.Timestamp, s.videoWallFirst)
}
first := s.videoTsFirst
var offsetTicks uint32
if s.audioTsGot && s.videoTsGot {
wallDelta := s.videoWallFirst - s.audioWallFirst
offsetTicks = uint32(wallDelta * 90000 / 1_000_000_000)
}
s.tsLock.Unlock()
pkt.Timestamp = (pkt.Timestamp - first) + offsetTicks
count++
if count <= 5 || count%300 == 0 {
log.Printf("[VIDEO] #%d ts=%d (%.3fs) offset=%d marker=%v",
count, pkt.Timestamp, float64(pkt.Timestamp)/90000.0, offsetTicks, pkt.Marker)
}
s.peersLock.RLock()
for _, peer := range s.peers {
if peer.Active && peer.VideoTrack != nil {
peer.VideoTrack.WriteRTP(pkt)
}
}
s.peersLock.RUnlock()
}
}
func (s *Sidecar) readAudioRTP() {
buf := make([]byte, 1500)
pkt := &rtp.Packet{}
count := 0
for s.running {
n, err := s.audioConn.Read(buf)
if err != nil {
if s.running {
log.Printf("[RTP] Audio read error: %v", err)
}
return
}
if err := pkt.Unmarshal(buf[:n]); err != nil {
continue
}
s.tsLock.Lock()
if !s.audioTsGot {
s.audioTsFirst = pkt.Timestamp
s.audioWallFirst = time.Now().UnixNano()
s.audioTsGot = true
log.Printf("[AUDIO] First ts=%d wall=%d", pkt.Timestamp, s.audioWallFirst)
}
first := s.audioTsFirst
var offsetTicks uint32
if s.audioTsGot && s.videoTsGot {
wallDelta := s.audioWallFirst - s.videoWallFirst
offsetTicks = uint32(wallDelta * 48000 / 1_000_000_000)
}
s.tsLock.Unlock()
pkt.Timestamp = (pkt.Timestamp - first) + offsetTicks
count++
if count <= 5 || count%500 == 0 {
log.Printf("[AUDIO] #%d ts=%d (%.3fs) offset=%d",
count, pkt.Timestamp, float64(pkt.Timestamp)/48000.0, offsetTicks)
}
s.peersLock.RLock()
for _, peer := range s.peers {
if peer.Active && peer.AudioTrack != nil {
peer.AudioTrack.WriteRTP(pkt)
}
}
s.peersLock.RUnlock()
}
}
func (s *Sidecar) CreatePeer(id string) (string, error) {
iceServers := []webrtc.ICEServer{}
for _, stun := range getStunServers() {
iceServers = append(iceServers, webrtc.ICEServer{URLs: []string{stun}})
}
m := &webrtc.MediaEngine{}
if err := m.RegisterCodec(webrtc.RTPCodecParameters{
RTPCodecCapability: webrtc.RTPCodecCapability{
MimeType: webrtc.MimeTypeVP8,
ClockRate: 90000,
SDPFmtpLine: "",
},
PayloadType: 96,
}, webrtc.RTPCodecTypeVideo); err != nil {
return "", err
}
if err := m.RegisterCodec(webrtc.RTPCodecParameters{
RTPCodecCapability: webrtc.RTPCodecCapability{
MimeType: webrtc.MimeTypeOpus,
ClockRate: 48000,
Channels: 2,
},
PayloadType: 111,
}, webrtc.RTPCodecTypeAudio); err != nil {
return "", err
}
i := &interceptor.Registry{}
intervalPliFactory, err := intervalpli.NewReceiverInterceptor()
if err != nil {
return "", err
}
i.Add(intervalPliFactory)
if err := webrtc.RegisterDefaultInterceptors(m, i); err != nil {
return "", err
}
api := webrtc.NewAPI(webrtc.WithMediaEngine(m), webrtc.WithInterceptorRegistry(i))
pc, err := api.NewPeerConnection(webrtc.Configuration{
ICEServers: iceServers,
})
if err != nil {
return "", fmt.Errorf("create PeerConnection: %w", err)
}
videoTrack, err := webrtc.NewTrackLocalStaticRTP(
webrtc.RTPCodecCapability{MimeType: webrtc.MimeTypeVP8, ClockRate: 90000},
"video", "ts6-stream",
)
if err != nil {
pc.Close()
return "", err
}
audioTrack, err := webrtc.NewTrackLocalStaticRTP(
webrtc.RTPCodecCapability{MimeType: webrtc.MimeTypeOpus, ClockRate: 48000, Channels: 2},
"audio", "ts6-stream",
)
if err != nil {
pc.Close()
return "", err
}
if _, err = pc.AddTrack(videoTrack); err != nil {
pc.Close()
return "", err
}
if _, err = pc.AddTrack(audioTrack); err != nil {
pc.Close()
return "", err
}
peer := &Peer{
ID: id,
PC: pc,
VideoTrack: videoTrack,
AudioTrack: audioTrack,
Active: false,
}
pc.OnICEConnectionStateChange(func(state webrtc.ICEConnectionState) {
log.Printf("[Peer %s] ICE: %s", id, state.String())
switch state {
case webrtc.ICEConnectionStateConnected:
peer.mu.Lock()
peer.Active = true
peer.mu.Unlock()
case webrtc.ICEConnectionStateDisconnected, webrtc.ICEConnectionStateFailed, webrtc.ICEConnectionStateClosed:
peer.mu.Lock()
peer.Active = false
peer.mu.Unlock()
}
})
s.peersLock.Lock()
if old, exists := s.peers[id]; exists {
old.Active = false
old.PC.Close()
}
s.peers[id] = peer
s.peersLock.Unlock()
offer, err := pc.CreateOffer(nil)
if err != nil {
return "", fmt.Errorf("create offer: %w", err)
}
if err := pc.SetLocalDescription(offer); err != nil {
return "", fmt.Errorf("set local desc: %w", err)
}
gatherComplete := webrtc.GatheringCompletePromise(pc)
<-gatherComplete
return pc.LocalDescription().SDP, nil
}
func (s *Sidecar) SetAnswer(id, sdp string) error {
s.peersLock.RLock()
peer, exists := s.peers[id]
s.peersLock.RUnlock()
if !exists {
return fmt.Errorf("peer %s not found", id)
}
return peer.PC.SetRemoteDescription(webrtc.SessionDescription{
Type: webrtc.SDPTypeAnswer,
SDP: sdp,
})
}
func (s *Sidecar) AddICECandidate(id string, candidate string, sdpMid string, sdpMLineIndex uint16) error {
s.peersLock.RLock()
peer, exists := s.peers[id]
s.peersLock.RUnlock()
if !exists {
return fmt.Errorf("peer %s not found", id)
}
return peer.PC.AddICECandidate(webrtc.ICECandidateInit{
Candidate: candidate,
SDPMid: &sdpMid,
SDPMLineIndex: &sdpMLineIndex,
})
}
func (s *Sidecar) ClosePeer(id string) {
s.peersLock.Lock()
if peer, exists := s.peers[id]; exists {
peer.Active = false
peer.PC.Close()
delete(s.peers, id)
}
s.peersLock.Unlock()
}
func (s *Sidecar) StartFFmpeg(source string) {
s.ffmpegLock.Lock()
defer s.ffmpegLock.Unlock()
s.StopFFmpegLocked()
s.source = source
args := []string{}
if source != "" {
if strings.HasPrefix(source, "http://") || strings.HasPrefix(source, "https://") {
args = append(args, "-reconnect", "1", "-reconnect_streamed", "1", "-reconnect_delay_max", "5")
} else {
args = append(args, "-stream_loop", "-1")
}
args = append(args, "-fflags", "+genpts+discardcorrupt", "-re", "-i", source)
} else {
w := envIntOrDefault("VIDEO_WIDTH", 1280)
h := envIntOrDefault("VIDEO_HEIGHT", 720)
args = append(args, "-re", "-f", "lavfi", "-i", fmt.Sprintf("color=c=black:s=%dx%d:r=1", w, h))
}
w := envIntOrDefault("VIDEO_WIDTH", 1280)
h := envIntOrDefault("VIDEO_HEIGHT", 720)
fps := envOrDefault("VIDEO_FRAMERATE", "30")
vBitrate := envOrDefault("VIDEO_BITRATE", "1500k")
if source != "" {
vf := fmt.Sprintf("scale=%d:%d:force_original_aspect_ratio=decrease,pad=%d:%d:(ow-iw)/2:(oh-ih)/2,format=yuv420p", w, h, w, h)
args = append(args,
"-map", "0:v:0",
"-vf", vf,
"-r", fps,
)
}
args = append(args,
"-pix_fmt", "yuv420p",
"-c:v", "libvpx",
"-cpu-used", "8",
"-deadline", "realtime",
"-lag-in-frames", "0",
"-error-resilient", "1",
"-b:v", vBitrate,
"-maxrate", "2M",
"-bufsize", "100k",
"-keyint_min", "15",
"-g", "15",
"-auto-alt-ref", "0",
"-payload_type", "96",
"-ssrc", "11111111",
"-f", "rtp",
fmt.Sprintf("rtp://127.0.0.1:%d", s.videoPort),
)
if source != "" {
aBitrate := envOrDefault("AUDIO_BITRATE", "128k")
args = append(args,
"-map", "0:a:0?",
"-c:a", "libopus",
"-b:a", aBitrate,
"-ar", "48000",
"-ac", "2",
"-payload_type", "111",
"-ssrc", "22222222",
"-f", "rtp",
fmt.Sprintf("rtp://127.0.0.1:%d", s.audioPort),
)
}
s.resetTimestamps()
log.Printf("[FFmpeg] Starting: source=%s video=:%d audio=:%d", source, s.videoPort, s.audioPort)
cmd := exec.Command(getFfmpegPath(), args...)
cmd.Stdout = nil
cmd.Stderr = os.Stderr
if err := cmd.Start(); err != nil {
log.Printf("[FFmpeg] Start error: %v", err)
return
}
s.ffmpeg = cmd
go func() {
err := cmd.Wait()
log.Printf("[FFmpeg] Exited: %v", err)
}()
}
func (s *Sidecar) StopFFmpegLocked() {
if s.ffmpeg != nil && s.ffmpeg.Process != nil {
s.ffmpeg.Process.Kill()
s.ffmpeg = nil
}
}
func (s *Sidecar) GetStats() map[string]interface{} {
s.peersLock.RLock()
defer s.peersLock.RUnlock()
peers := map[string]interface{}{}
for id, peer := range s.peers {
peers[id] = map[string]interface{}{
"active": peer.Active,
"state": peer.PC.ICEConnectionState().String(),
}
}
return map[string]interface{}{
"videoPort": s.videoPort,
"audioPort": s.audioPort,
"peerCount": len(s.peers),
"peers": peers,
"source": s.source,
}
}
func (s *Sidecar) Stop() {
s.running = false
s.ffmpegLock.Lock()
s.StopFFmpegLocked()
s.ffmpegLock.Unlock()
if s.videoConn != nil {
s.videoConn.Close()
}
if s.audioConn != nil {
s.audioConn.Close()
}
s.peersLock.Lock()
for id, peer := range s.peers {
peer.Active = false
peer.PC.Close()
delete(s.peers, id)
}
s.peersLock.Unlock()
}
func main() {
port := 9800
if p := os.Getenv("SIDECAR_PORT"); p != "" {
if v, err := strconv.Atoi(p); err == nil {
port = v
}
}
sidecar := NewSidecar()
if err := sidecar.StartRTP(); err != nil {
log.Fatalf("Failed to start RTP: %v", err)
}
mux := http.NewServeMux()
mux.HandleFunc("POST /peer/create", func(w http.ResponseWriter, r *http.Request) {
var req struct {
ID string `json:"id"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, err.Error(), 400)
return
}
log.Printf("[API] Creating peer: %s", req.ID)
sdp, err := sidecar.CreatePeer(req.ID)
if err != nil {
log.Printf("[API] CreatePeer error: %v", err)
http.Error(w, err.Error(), 500)
return
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]string{"sdp": sdp})
})
mux.HandleFunc("POST /peer/answer", func(w http.ResponseWriter, r *http.Request) {
var req struct {
ID string `json:"id"`
SDP string `json:"sdp"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, err.Error(), 400)
return
}
log.Printf("[API] Setting answer for peer: %s (%d bytes)", req.ID, len(req.SDP))
if err := sidecar.SetAnswer(req.ID, req.SDP); err != nil {
log.Printf("[API] SetAnswer error: %v", err)
http.Error(w, err.Error(), 500)
return
}
json.NewEncoder(w).Encode(map[string]string{"status": "ok"})
})
mux.HandleFunc("POST /peer/ice", func(w http.ResponseWriter, r *http.Request) {
var req struct {
ID string `json:"id"`
Candidate string `json:"candidate"`
SDPMid string `json:"sdpMid"`
SDPMLineIndex uint16 `json:"sdpMLineIndex"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, err.Error(), 400)
return
}
if err := sidecar.AddICECandidate(req.ID, req.Candidate, req.SDPMid, req.SDPMLineIndex); err != nil {
log.Printf("[API] AddICE error: %v", err)
http.Error(w, err.Error(), 500)
return
}
json.NewEncoder(w).Encode(map[string]string{"status": "ok"})
})
mux.HandleFunc("POST /peer/close", func(w http.ResponseWriter, r *http.Request) {
var req struct {
ID string `json:"id"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, err.Error(), 400)
return
}
sidecar.ClosePeer(req.ID)
json.NewEncoder(w).Encode(map[string]string{"status": "ok"})
})
mux.HandleFunc("POST /source", func(w http.ResponseWriter, r *http.Request) {
var req struct {
Source string `json:"source"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, err.Error(), 400)
return
}
log.Printf("[API] Setting source: %s", req.Source)
sidecar.StartFFmpeg(req.Source)
json.NewEncoder(w).Encode(map[string]string{"status": "ok"})
})
mux.HandleFunc("POST /source/stop", func(w http.ResponseWriter, r *http.Request) {
sidecar.ffmpegLock.Lock()
sidecar.StopFFmpegLocked()
sidecar.ffmpegLock.Unlock()
json.NewEncoder(w).Encode(map[string]string{"status": "ok"})
})
mux.HandleFunc("GET /stats", func(w http.ResponseWriter, r *http.Request) {
json.NewEncoder(w).Encode(sidecar.GetStats())
})
mux.HandleFunc("GET /health", func(w http.ResponseWriter, r *http.Request) {
json.NewEncoder(w).Encode(map[string]interface{}{
"status": "ok",
"videoPort": sidecar.videoPort,
"audioPort": sidecar.audioPort,
})
})
go func() {
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
<-sigCh
log.Println("Shutting down...")
sidecar.Stop()
os.Exit(0)
}()
log.Printf("[Sidecar] HTTP API listening on :%d", port)
if err := http.ListenAndServe(fmt.Sprintf(":%d", port), mux); err != nil {
log.Fatalf("HTTP server error: %v", err)
}
}