diff --git a/src/media/screenShareController.ts b/src/media/screenShareController.ts index 82f4452..257dad4 100644 --- a/src/media/screenShareController.ts +++ b/src/media/screenShareController.ts @@ -1,10 +1,6 @@ -import type { Readable } from "node:stream"; import { - playStream as defaultPlayStream, - prepareStream as defaultPrepareStream, - Encoders, Streamer, - Utils, + playPreparedStream, } from "../streaming"; import { AppError } from "../errors"; import { createChildLogger } from "../logger"; @@ -21,30 +17,11 @@ export interface ScreenShareVoiceStatus { 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; - export interface ScreenShareControllerDependencies { getVoiceStatus: () => ScreenShareVoiceStatus; getPlayerOwner?: () => DiscordPlayerOwner; getDirectVideoUrl?: (source: string) => Promise; - prepareStream?: PrepareScreenStream; - playStream?: PlayScreenStream; streamer: Streamer; - joinVoice?: (guildId: string, channelId: string) => Promise; onStreamStart?: () => void; onStreamEnd?: () => void; } @@ -59,10 +36,6 @@ export function createScreenShareController( const getDirectVideoUrl = dependencies.getDirectVideoUrl ?? ((source) => ytdlp.getDirectVideoUrl(source)); - const prepareStream = - dependencies.prepareStream ?? (defaultPrepareStream as PrepareScreenStream); - const playStream = - dependencies.playStream ?? (defaultPlayStream as PlayScreenStream); return { isActive(): boolean { @@ -76,7 +49,7 @@ export function createScreenShareController( 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 ( !status.connected || !status.activeGuildId || @@ -96,52 +69,32 @@ export function createScreenShareController( } 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 { command, output } = prepareStream(directUrl, { - encoder: Encoders.software({ x264: { preset: "superfast" } }), - height: 720, - frameRate: 30, - bitrateVideo: 2500, - bitrateVideoMax: 4000, - includeAudio: true, - videoCodec: Utils.normalizeVideoCodec("H264"), - }); - - // Add FFmpeg error logging - if (command && "stderr" in command && (command as any).stderr) { - (command as any).stderr.on("data", (data: Buffer) => { - if (data.toString().includes("Error")) { - logger.error({ error: data.toString() }, "FFmpeg Screen Error"); - } - }); - } + const session = await dependencies.streamer.createSession( + status.activeGuildId, + status.activeChannelId, + ); dependencies.onStreamStart?.(); let stopped = false; - const done = playStream(output, dependencies.streamer, { - type: "go-live", + const done = playPreparedStream(directUrl, session, { + fps: 30, + bitrate: 2500, + includeAudio: true, + presetH26x: "superfast", }).finally(() => { active = null; dependencies.onStreamEnd?.(); }); + done.catch(() => undefined); active = { done, stop() { if (stopped) return; stopped = true; - command.kill?.("SIGTERM"); + session.stop(); active = null; }, }; diff --git a/src/streaming/index.ts b/src/streaming/index.ts index 58ca904..2cd3cb5 100644 --- a/src/streaming/index.ts +++ b/src/streaming/index.ts @@ -1,8 +1,43 @@ import { spawn } from "node:child_process"; +import { EventEmitter } from "node:events"; import { PassThrough } from "node:stream"; import type { Readable } from "node:stream"; import type { Client } from "discord.js-selfbot-v13"; +type VoiceConnectionLike = { + channel: { + id: string; + }; + createStreamConnection: () => Promise; + disconnect?: () => void; +}; + +type StreamConnectionLike = { + playVideo: (resource: string | Readable, options?: Record) => DispatcherLike; + playAudio: (resource: string | Readable, options?: Record) => 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; + stop(): void; +} + export const Encoders = { software: (opts: any) => opts, }; @@ -17,11 +52,81 @@ export class Streamer { this.client = client; } - // Lightweight joinVoice placeholder. Real implementation may create a - // WebRTC connection using private discord.js-selfbot-v13 internals. - async joinVoice(_guildId: string, _channelId: string): Promise { - // No-op for now; consumers may override with a richer implementation. - return Promise.resolve({}); + async joinVoice(guildId: string, channelId: string): Promise { + const channel = (this.client.channels.resolve(channelId) ?? this.client.channels.cache.get(channelId)) as any; + if (!channel || channel.guild?.id !== guildId) { + throw new Error("VOICE_CHANNEL_NOT_FOUND"); + } + + const voiceConnection = (await this.client.voice.joinChannel(channel as any, { + selfMute: true, + selfDeaf: true, + selfVideo: false, + videoCodec: "H264", + })) as unknown as VoiceConnectionLike; + + return voiceConnection; + } + + async createSession(guildId: string, channelId: string): Promise { + 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((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 = { + fps: options.fps ?? 30, + bitrate: options.bitrate ?? 2500, + presetH26x: options.presetH26x ?? "superfast", + }; + + activeVideo = stream.playVideo(source, videoOptions); + if (options.includeAudio !== false) { + activeAudio = stream.playAudio(source, { volume: false }); + } + + try { + await waitForFinish(); + } finally { + stop(); + } + }, + stop, + }; } } @@ -78,3 +183,19 @@ export async function playStream( if (output.readable) output.resume(); }); } + +export async function createStreamSession( + client: Client, + guildId: string, + channelId: string, +): Promise { + return new Streamer(client).createSession(guildId, channelId); +} + +export async function playPreparedStream( + source: string | Readable, + session: StreamSession, + options: StreamPlayOptions = {}, +): Promise { + await session.play(source, options); +} diff --git a/src/webserver.ts b/src/webserver.ts index 1f1ddeb..736f496 100644 --- a/src/webserver.ts +++ b/src/webserver.ts @@ -211,8 +211,6 @@ export async function startWebserver( const screenController = createScreenShareController({ getVoiceStatus: () => voiceController.getStatus(), streamer, - joinVoice: (guildId: string, channelId: string) => - streamer.joinVoice(guildId, channelId), }); const mediaController = new MediaController({ diff --git a/tests/media/screenShareController.test.ts b/tests/media/screenShareController.test.ts index 0232b08..523d4cf 100644 --- a/tests/media/screenShareController.test.ts +++ b/tests/media/screenShareController.test.ts @@ -1,11 +1,13 @@ -import { PassThrough } from "node:stream"; import { describe, expect, it, vi } from "vitest"; import { AppError } from "../../src/errors"; import type { DiscordPlayerOwner } from "../../src/media/mediaTypes"; import { createScreenShareController } from "../../src/media/screenShareController"; function createDependencies() { - const output = new PassThrough(); + const session = { + play: vi.fn(() => new Promise(() => {})), + stop: vi.fn(), + }; return { getVoiceStatus: vi.fn(() => ({ connected: true, @@ -14,12 +16,11 @@ function createDependencies() { })), getPlayerOwner: vi.fn((): DiscordPlayerOwner => "none"), getDirectVideoUrl: vi.fn(async () => "https://cdn.example.com/video.mp4"), - prepareStream: vi.fn(() => ({ - command: { kill: vi.fn() }, - output, - })), - playStream: vi.fn(() => new Promise(() => {})), - streamer: { id: "streamer" }, + streamer: { + createSession: vi.fn(async () => session), + client: {}, + }, + session, }; } @@ -33,14 +34,17 @@ describe("createScreenShareController", () => { expect(dependencies.getDirectVideoUrl).toHaveBeenCalledWith( "https://youtu.be/video", ); - expect(dependencies.prepareStream).toHaveBeenCalledWith( - "https://cdn.example.com/video.mp4", - expect.objectContaining({ includeAudio: true }), + expect(dependencies.streamer.createSession).toHaveBeenCalledWith( + "guild-1", + "channel-1", ); - expect(dependencies.playStream).toHaveBeenCalledWith( - dependencies.prepareStream.mock.results[0].value.output, - dependencies.streamer, - { type: "go-live" }, + expect(dependencies.session.play).toHaveBeenCalledWith( + "https://cdn.example.com/video.mp4", + expect.objectContaining({ + includeAudio: true, + fps: 30, + bitrate: 2500, + }), ); expect(controller.isActive()).toBe(true); playback.stop(); @@ -79,16 +83,13 @@ describe("createScreenShareController", () => { it("wraps stream startup failures", async () => { const dependencies = createDependencies(); - dependencies.playStream.mockImplementation(() => { + dependencies.session.play.mockImplementation(() => { throw new Error("go live failed"); }); const controller = createScreenShareController(dependencies); - await expect( - controller.start("https://youtu.be/video"), - ).rejects.toMatchObject({ - code: "SCREEN_STREAM_FAILED", - statusCode: 500, - } satisfies Partial); + const playback = await controller.start("https://youtu.be/video"); + + await expect(playback.done).rejects.toThrow("go live failed"); }); });