Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
150 changes: 128 additions & 22 deletions packages/bot/src/handlers/player/playerFactory.bridge.spec.ts
Original file line number Diff line number Diff line change
@@ -1,14 +1,16 @@
import { beforeEach, describe, expect, it, jest } from '@jest/globals'
import { EventEmitter } from 'events'
import { PassThrough } from 'stream'
import type { Readable } from 'stream'

const playdlSearchMock = jest.fn()
const playdlStreamMock = jest.fn()
const spawnMock = jest.fn()

jest.mock('child_process', () => ({
spawn: (...args: unknown[]) => spawnMock(...args),
}))

// Stub `discord-player` and `@discord-player/extractor` so importing
// playerFactory.ts doesn't pull in their full dependency trees (one of
// which requires the native `file-type` module that isn't present in
// the test environment). We only exercise the bridge fallback chain —
// the Player class and DefaultExtractors aren't touched in these tests.
jest.mock('discord-player', () => ({
Player: class {
extractors = {
Expand Down Expand Up @@ -37,17 +39,36 @@ jest.mock('@lucky/shared/utils', () => ({
debugLog: jest.fn(),
}))

// Import AFTER the mocks so the real module wires against the mock.
// eslint-disable-next-line @typescript-eslint/no-var-requires
import {
createResilientStream,
streamViaYtDlp,
streamViaSoundCloud,
findMatchingSoundCloudResult,
parseDurationString,
} from './playerFactory'

const fakeStream = { on: jest.fn() } as unknown as Readable

function makeSpawnSuccess() {
const stdout = new PassThrough()
const proc = Object.assign(new EventEmitter(), {
stdout,
kill: jest.fn(),
})
setImmediate(() => stdout.emit('data', Buffer.from('audio')))
return proc
}

function makeSpawnError(code = 1) {
const stdout = new PassThrough()
const proc = Object.assign(new EventEmitter(), {
stdout,
kill: jest.fn(),
})
setImmediate(() => proc.emit('close', code))
return proc
}

function makeTrack(overrides: Partial<Record<string, unknown>> = {}) {
return {
title: 'Bohemian Rhapsody',
Expand All @@ -58,7 +79,7 @@ function makeTrack(overrides: Partial<Record<string, unknown>> = {}) {
}
}

describe('parseDurationString (real export)', () => {
describe('parseDurationString', () => {
it('parses m:ss', () => {
expect(parseDurationString('3:45')).toBe(225)
})
Expand All @@ -81,7 +102,7 @@ describe('parseDurationString (real export)', () => {
})
})

describe('findMatchingSoundCloudResult (real export)', () => {
describe('findMatchingSoundCloudResult', () => {
const results = [
{ name: 'Bohemian Rhapsody', url: 'sc://1', durationInSec: 354 },
{ name: 'Unrelated', url: 'sc://2', durationInSec: 180 },
Expand Down Expand Up @@ -166,25 +187,109 @@ describe('streamViaSoundCloud', () => {
})
})

describe('streamViaYtDlp', () => {
beforeEach(() => {
spawnMock.mockReset()
})

it('resolves with stdout when yt-dlp emits data', async () => {
const proc = makeSpawnSuccess()
spawnMock.mockReturnValue(proc)

const stream = await streamViaYtDlp('https://youtube.com/watch?v=test')
expect(stream).toBe(proc.stdout)
expect(spawnMock).toHaveBeenCalledWith(
'yt-dlp',
expect.arrayContaining([
'--no-playlist',
'-o',
'-',
'https://youtube.com/watch?v=test',
]),
expect.objectContaining({ stdio: ['ignore', 'pipe', 'pipe'] }),
)
})

it('rejects for non-https URLs without spawning', async () => {
await expect(
streamViaYtDlp('http://youtube.com/watch?v=test'),
).rejects.toThrow(/only https/i)
expect(spawnMock).not.toHaveBeenCalled()
})

it('rejects for domains not in the allowlist without spawning', async () => {
await expect(
streamViaYtDlp('https://evil.com/audio.mp3'),
).rejects.toThrow(/not in allowlist/i)
expect(spawnMock).not.toHaveBeenCalled()
})

it('accepts all allowed domains', async () => {
const allowedUrls = [
'https://youtube.com/watch?v=test',
'https://www.youtube.com/watch?v=test',
'https://youtu.be/test',
'https://music.youtube.com/watch?v=test',
'https://soundcloud.com/artist/track',
'https://open.spotify.com/track/test',
]
for (const url of allowedUrls) {
const proc = makeSpawnSuccess()
spawnMock.mockReturnValue(proc)
await expect(streamViaYtDlp(url)).resolves.toBeDefined()
spawnMock.mockReset()
}
})

it('rejects when yt-dlp closes with non-zero exit code', async () => {
const proc = makeSpawnError(1)
spawnMock.mockReturnValue(proc)

await expect(
streamViaYtDlp('https://youtube.com/watch?v=test'),
).rejects.toThrow(/exited with code 1/)
})

it('rejects when spawn emits error', async () => {
const stdout = new PassThrough()
const proc = Object.assign(new EventEmitter(), {
stdout,
kill: jest.fn(),
})
spawnMock.mockReturnValue(proc)
setImmediate(() => proc.emit('error', new Error('ENOENT')))

await expect(
streamViaYtDlp('https://youtube.com/watch?v=test'),
).rejects.toThrow('ENOENT')
})
})

describe('createResilientStream', () => {
beforeEach(() => {
spawnMock.mockReset()
playdlSearchMock.mockReset()
playdlStreamMock.mockReset()
})

it('streams directly from source URL on first attempt', async () => {
playdlStreamMock.mockResolvedValueOnce({ stream: fakeStream })
it('streams via yt-dlp from source URL on first attempt', async () => {
const proc = makeSpawnSuccess()
spawnMock.mockReturnValue(proc)

const result = await createResilientStream(makeTrack())
expect(result).toBe(fakeStream)
expect(playdlStreamMock).toHaveBeenCalledWith(
'https://youtube.com/watch?v=fakeBohemian',
expect(result).toBe(proc.stdout)
expect(spawnMock).toHaveBeenCalledWith(
'yt-dlp',
expect.arrayContaining([
'https://youtube.com/watch?v=fakeBohemian',
]),
expect.anything(),
)
expect(playdlSearchMock).not.toHaveBeenCalled()
})

it('falls back to SoundCloud primary search when direct stream fails', async () => {
playdlStreamMock.mockRejectedValueOnce(new Error('403'))
it('falls back to SoundCloud primary search when yt-dlp fails', async () => {
spawnMock.mockReturnValue(makeSpawnError(1))
playdlSearchMock.mockResolvedValueOnce([
{
name: 'Bohemian Rhapsody - Queen',
Expand All @@ -201,7 +306,7 @@ describe('createResilientStream', () => {
})

it('falls back to SoundCloud title-only when primary returns no validated match', async () => {
playdlStreamMock.mockRejectedValueOnce(new Error('403'))
spawnMock.mockReturnValue(makeSpawnError(1))
playdlSearchMock
.mockResolvedValueOnce([
{ name: 'Unrelated', url: 'sc://miss', durationInSec: 180 },
Expand All @@ -221,20 +326,21 @@ describe('createResilientStream', () => {
expect(playdlStreamMock).toHaveBeenCalledWith('sc://secondary')
})

it('streams directly from source URL even for spam uploader channels', async () => {
playdlStreamMock.mockResolvedValueOnce({ stream: fakeStream })
it('streams via yt-dlp even for spam uploader channels', async () => {
const proc = makeSpawnSuccess()
spawnMock.mockReturnValue(proc)

const result = await createResilientStream(
makeTrack({
title: 'GOLDEN - KPOP DEMON HUNTERS - HUNTR/X [Download]',
author: 'Best Songs',
}),
)
expect(result).toBe(fakeStream)
expect(result).toBe(proc.stdout)
expect(playdlSearchMock).not.toHaveBeenCalled()
})

it('throws "Bridge exhausted" when a track has no URL and SoundCloud fails', async () => {
it('throws "Bridge exhausted" when track has no URL and SoundCloud fails', async () => {
playdlSearchMock.mockResolvedValue([])

await expect(
Expand All @@ -243,7 +349,7 @@ describe('createResilientStream', () => {
})

it('throws "Bridge exhausted" when all stages fail', async () => {
playdlStreamMock.mockRejectedValueOnce(new Error('youtube 403'))
spawnMock.mockReturnValue(makeSpawnError(1))
playdlSearchMock.mockResolvedValue([])

await expect(createResilientStream(makeTrack())).rejects.toThrow(
Expand Down
101 changes: 90 additions & 11 deletions packages/bot/src/handlers/player/playerFactory.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
import type { Track } from 'discord-player'
import { DefaultExtractors } from '@discord-player/extractor'
import * as playdl from 'play-dl'
import { spawn } from 'child_process'

Check warning on line 5 in packages/bot/src/handlers/player/playerFactory.ts

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Prefer `node:child_process` over `child_process`.

See more on https://sonarcloud.io/project/issues?id=LucasSantana-Dev_Lucky&issues=AZ14DCiB44PnsT_gIb8q&open=AZ14DCiB44PnsT_gIb8q&pullRequest=520
import type { Readable } from 'stream'
import type { CustomClient } from '../../types'
import { errorLog, infoLog, warnLog, debugLog } from '@lucky/shared/utils'
Expand Down Expand Up @@ -101,14 +102,93 @@
}
}

/**
* Stream audio via yt-dlp subprocess.
*
* Resolves with the stdout Readable once yt-dlp begins writing data.
* Rejects with a timeout error if no data arrives within 15 seconds.
* Exported for testing.
*/
const ALLOWED_YTDLP_DOMAINS = new Set([
'youtube.com',
'www.youtube.com',
'youtu.be',
'music.youtube.com',
'soundcloud.com',
'www.soundcloud.com',
'open.spotify.com',
])

function validateYtDlpUrl(url: string): void {
let parsed: URL
try {
parsed = new URL(url)
} catch {
throw new Error(`yt-dlp: invalid URL`)
}
if (parsed.protocol !== 'https:') {
throw new Error(`yt-dlp: only https URLs are allowed`)
}
if (!ALLOWED_YTDLP_DOMAINS.has(parsed.hostname.toLowerCase())) {
throw new Error(`yt-dlp: domain not in allowlist: ${parsed.hostname}`)
}
}

export function streamViaYtDlp(url: string): Promise<Readable> {
try {
validateYtDlpUrl(url)
} catch (err) {
return Promise.reject(err)
}
return new Promise<Readable>((resolve, reject) => {
const proc = spawn(
'yt-dlp',
[
'--no-playlist',
'-f',
'bestaudio/best',
'-o',
'-',
'--quiet',
'--no-warnings',
'--no-progress',
url,
],
{ stdio: ['ignore', 'pipe', 'pipe'] },
)

const timeout = setTimeout(() => {
proc.kill()
reject(new Error('yt-dlp: timed out waiting for stream start'))
}, 15_000)

proc.stdout!.once('data', () => {

Check warning on line 165 in packages/bot/src/handlers/player/playerFactory.ts

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

This assertion is unnecessary since it does not change the type of the expression.

See more on https://sonarcloud.io/project/issues?id=LucasSantana-Dev_Lucky&issues=AZ14DCiB44PnsT_gIb8s&open=AZ14DCiB44PnsT_gIb8s&pullRequest=520
clearTimeout(timeout)
resolve(proc.stdout!)

Check warning on line 167 in packages/bot/src/handlers/player/playerFactory.ts

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

This assertion is unnecessary since it does not change the type of the expression.

See more on https://sonarcloud.io/project/issues?id=LucasSantana-Dev_Lucky&issues=AZ14DCiB44PnsT_gIb8t&open=AZ14DCiB44PnsT_gIb8t&pullRequest=520
})

proc.once('error', (err) => {
clearTimeout(timeout)
reject(err)
})

proc.once('close', (code) => {
clearTimeout(timeout)
if (code && code !== 0) {
reject(new Error(`yt-dlp exited with code ${code}`))
}
})
Comment on lines +143 to +180

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major

🧩 Analysis chain

🌐 Web query:

Node.js v22 docs: does a Readable stream 'data' listener consume that chunk, and what do the codeandsignalarguments on child_processclose mean when a child exits before producing stdout?

💡 Result:

Yes, adding a 'data' event listener to a Node.js Readable stream consumes the chunk by switching the stream to flowing mode, where data is automatically emitted via 'data' events, relinquishing ownership of chunks to the listener. This occurs when attaching the listener, calling resume, or pipe. The listener receives the chunk as a Buffer or string (if encoding set). For child_process 'close' event (e.g., from spawn), the callback receives (code, signal). Code is the exit code if the child exited normally (non-zero often indicates error), or null if terminated by signal. Signal is the string name of the terminating signal (e.g., 'SIGTERM', 'SIGKILL') or null if no signal. One is always non-null. This applies even if the child exits before producing stdout—no data emitted, but 'close' fires after process ends and stdio streams close. Use 'exit' for process termination (stdio may remain open); 'close' confirms all stdio closed. No stdout means child produced none or exited too quickly (e.g., immediate error).

Citations:


Fix stream and process lifecycle to avoid chunk loss and dangling promises.

The proc.stdout!.once('data') listener at line 140 consumes the first chunk by switching stdout to flowing mode, so that chunk is lost when returning the stream to the caller. Additionally, lines 150–154 clear the timeout without rejecting if yt-dlp exits cleanly (code 0 or signal) before producing output, leaving the promise unsettled and blocking the fallback to SoundCloud.

Use a PassThrough stream and track whether streaming has started: pipe stdout into it, check the started flag in the close handler, and reject if the process closes before any data arrives, regardless of exit code.

Suggested fix
-import type { Readable } from 'stream'
+import { PassThrough } from 'stream'
+import type { Readable } from 'stream'
@@
 return new Promise<Readable>((resolve, reject) => {
+    const stream = new PassThrough()
+    let started = false
     const proc = spawn(
         'yt-dlp',
         [
@@
         { stdio: ['ignore', 'pipe', 'pipe'] },
     )
+    proc.stdout!.pipe(stream)

     const timeout = setTimeout(() => {
         proc.kill()
         reject(new Error('yt-dlp: timed out waiting for stream start'))
     }, 15_000)

     proc.stdout!.once('data', () => {
+        started = true
         clearTimeout(timeout)
-        resolve(proc.stdout!)
+        resolve(stream)
     })

     proc.once('error', (err) => {
         clearTimeout(timeout)
         reject(err)
     })

-    proc.once('close', (code) => {
+    proc.once('close', (code, signal) => {
         clearTimeout(timeout)
+        if (!started) {
+            reject(
+                new Error(
+                    signal
+                        ? `yt-dlp exited before streaming (signal ${signal})`
+                        : `yt-dlp exited before streaming any audio (code ${code ?? 0})`,
+                ),
+            )
+            return
+        }
         if (code && code !== 0) {
             reject(new Error(`yt-dlp exited with code ${code}`))
         }
     })
 })
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
return new Promise<Readable>((resolve, reject) => {
const proc = spawn(
'yt-dlp',
[
'--no-playlist',
'-f',
'bestaudio/best',
'-o',
'-',
'--quiet',
'--no-warnings',
'--no-progress',
url,
],
{ stdio: ['ignore', 'pipe', 'pipe'] },
)
const timeout = setTimeout(() => {
proc.kill()
reject(new Error('yt-dlp: timed out waiting for stream start'))
}, 15_000)
proc.stdout!.once('data', () => {
clearTimeout(timeout)
resolve(proc.stdout!)
})
proc.once('error', (err) => {
clearTimeout(timeout)
reject(err)
})
proc.once('close', (code) => {
clearTimeout(timeout)
if (code && code !== 0) {
reject(new Error(`yt-dlp exited with code ${code}`))
}
})
return new Promise<Readable>((resolve, reject) => {
const stream = new PassThrough()
let started = false
const proc = spawn(
'yt-dlp',
[
'--no-playlist',
'-f',
'bestaudio/best',
'-o',
'-',
'--quiet',
'--no-warnings',
'--no-progress',
url,
],
{ stdio: ['ignore', 'pipe', 'pipe'] },
)
proc.stdout!.pipe(stream)
const timeout = setTimeout(() => {
proc.kill()
reject(new Error('yt-dlp: timed out waiting for stream start'))
}, 15_000)
proc.stdout!.once('data', () => {
started = true
clearTimeout(timeout)
resolve(stream)
})
proc.once('error', (err) => {
clearTimeout(timeout)
reject(err)
})
proc.once('close', (code, signal) => {
clearTimeout(timeout)
if (!started) {
reject(
new Error(
signal
? `yt-dlp exited before streaming (signal ${signal})`
: `yt-dlp exited before streaming any audio (code ${code ?? 0})`,
),
)
return
}
if (code && code !== 0) {
reject(new Error(`yt-dlp exited with code ${code}`))
}
})
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@packages/bot/src/handlers/player/playerFactory.ts` around lines 118 - 155,
Replace the current direct use of proc.stdout and its once('data') listener with
a PassThrough proxy to avoid losing the first chunk: create a PassThrough (e.g.,
const stream = new PassThrough()), pipe proc.stdout into it
(proc.stdout!.pipe(stream)), and resolve the Promise with that PassThrough when
stream.once('data') fires; add a boolean started flag set when data arrives; in
proc.once('close') clear the timeout and if started is false reject the Promise
with a "timed out before stream started" error (regardless of exit code) to
ensure the promise doesn't remain unsettled; keep the existing
proc.once('error') behavior but ensure all paths clear the timeout and either
resolve with the PassThrough or reject so no dangling promises remain (refer to
symbols: proc, timeout, stream/PassThrough, started, proc.stdout!.pipe,
proc.once('close')).

})
}

/**
* Bridge fallback chain (discord-player-youtubei v3 createStream signature).
*
* Priority: YouTube direct → SoundCloud (full query) → SoundCloud (title only).
* Priority: yt-dlp direct → SoundCloud (full query) → SoundCloud (title only).
*
* YouTube direct is attempted first since the track always carries its source
* URL and play-dl can stream it independently of the extractor. SoundCloud
* search is the fallback for tracks blocked by YouTube bot-detection.
* yt-dlp is tried first (play-dl YouTube streaming is broken due to bot
* detection). SoundCloud search is the fallback for tracks unavailable via
* yt-dlp.
*/
export async function createResilientStream(
track: Pick<Track, 'title' | 'author' | 'duration' | 'url'>,
Expand All @@ -130,18 +210,17 @@

if (track.url) {
try {
const direct = await playdl.stream(track.url)
const stream = await streamViaYtDlp(track.url)
infoLog({
message: 'Bridge: streamed directly from source URL',
message: 'Bridge: streamed via yt-dlp',
data: { url: track.url, title: cleanedTitle || track.title },
})
return direct.stream
} catch (directError) {
return stream
} catch (ytdlpError) {
debugLog({
message:
'Bridge: direct stream failed, falling back to SoundCloud',
message: 'Bridge: yt-dlp failed, falling back to SoundCloud',
data: {
error: (directError as Error).message,
error: (ytdlpError as Error).message,
cleanedTitle,
},
})
Expand Down