fix(security): complete Electron hardening and review follow-ups

* fix(security): harden Electron IPC against MITM, SSRF, path and injection risks

S1 TLS: validate certs by default on playlist/EPG fetches (opt-out via IPTVNATOR_ALLOW_INSECURE_TLS); new util/secure-https.ts.
S2: write-file IPC restricted to save-dialog-authorized paths.
S3: XTREAM_PROBE_URL guarded by assertRemoteUrlAllowed + maxRedirects:0; new events/url-safety.ts (+19 tests).
S4: EPG titles rendered via interpolation, not [innerHTML].
S5: downloads reveal/play limited to recorded download paths.
S6: Stalker cmd encoded (slash-preserving) to block query injection.
EPG-worker and Stalker fetches reject file://-style/credentialed URLs; LAN/self-hosted targets remain allowed.

* perf(player): lazy-load web video players via @defer

Wrap Video.js/HTML5/ArtPlayer in @defer (on immediate) so video.js, hls.js,
artplayer and mpegts.js split into a deferred chunk loaded on first playback
instead of eagerly on the player route. Embedded MPV (native) stays eager.
Spec uses DeferBlockBehavior.Playthrough.

* fix(player): remove leaked HTML video listeners on destroy

volumechange used a mismatched removeEventListener reference, while
loadedmetadata and timeupdate were never removed at all. Bind all three to
stable handler fields used for both add and remove, and add a teardown
regression test asserting each listener is detached on destroy.

* refactor(dashboard): extract pure navigation helpers from DashboardDataService

Move the 8 stateless link/navigation-state/type-kind helpers into a new
dashboard-navigation.util.ts so the routing logic is independently testable and
the 1260-line god-service shrinks. DashboardDataService keeps the public methods
as thin delegators (facade) so the public API and the single consumer
(workspace-dashboard-rails) are unchanged. First slice of the DashboardDataService
decomposition; verified by the existing service spec (33/33) and the app typecheck.

* fix(review): address PR feedback (IPv6 link-local, write-path cap, @defer placeholder)

- url-safety: broaden IPv6 link-local detection to the full fe80::/10 range
  (fe80:: through febf::), not just the fe80:: prefix (+ regression tests).
- playlist.events: cap authorizedWritePaths (evict oldest past 32) so a save
  dialog opened without a following write cannot accumulate entries until restart.
- web-player-view: add a @placeholder to each @defer (on immediate) player block
  to avoid the one-frame blank/layout-shift before the chunk resolves.

* fix(security): close Electron network and download gaps

* test(downloads): cover cancellation and restart cleanup

* fix(downloads): address Greptile review gaps

* test(security): reproduce remaining Greptile findings

* fix(security): close remaining Greptile findings

* test(downloads): reproduce early database queue stall

* fix(downloads): release queue after setup failures

* test(downloads): reproduce completion queue stall

* fix(downloads): release queue after completion failures
This commit is contained in:
Salem authored and GitHub committed 2026-06-12 15:24:29 +02:00
1 parent 7775553a6b
commit 2c032cd3c8
44 files changed
+3646 -1124

No files matched your search

+9
View File
@@ -304,6 +304,15 @@ Useful narrower flags:
- `IPTVNATOR_TRACE_RENDERER_CONSOLE=1` mirrors renderer console messages into
the Electron terminal output
Security-sensitive network compatibility flags are opt-in:
- `IPTVNATOR_ALLOW_PRIVATE_NETWORK_URLS=1` permits EPG URLs that resolve to
localhost, LAN, or other private addresses. Leave this unset for playlists
you do not fully trust.
- `IPTVNATOR_ALLOW_INSECURE_TLS=1` disables certificate validation for remote
playlist imports and refreshes. Use it only for a trusted provider with a
self-signed or otherwise invalid certificate.
If the local Nx daemon gets into a bad state before rerunning Electron, reset it:
```
@@ -0,0 +1,76 @@
import { resolve } from 'node:path';
import { DownloadDirectoryAuthorizer } from './download-directory-authorization';
describe('DownloadDirectoryAuthorizer', () => {
it('returns a previously authorized native-dialog selection', async () => {
const authorizer = new DownloadDirectoryAuthorizer({
getDefaultDirectory: () => '/downloads',
loadSelectedDirectory: async () => '/media/iptv',
saveSelectedDirectory: jest.fn(),
});
await expect(authorizer.getPreferredDirectory()).resolves.toBe(
resolve('/media/iptv')
);
await expect(authorizer.requireAuthorized('/media/iptv')).resolves.toBe(
resolve('/media/iptv')
);
});
it('rejects a renderer-supplied directory that was never authorized', async () => {
const authorizer = new DownloadDirectoryAuthorizer({
getDefaultDirectory: () => '/downloads',
loadSelectedDirectory: async () => null,
saveSelectedDirectory: jest.fn(),
});
await expect(
authorizer.requireAuthorized('/tmp/renderer-controlled')
).rejects.toThrow(/not authorized/i);
});
it('persists and authorizes a directory selected by the native dialog', async () => {
const saveSelectedDirectory = jest.fn().mockResolvedValue(undefined);
const authorizer = new DownloadDirectoryAuthorizer({
getDefaultDirectory: () => '/downloads',
loadSelectedDirectory: async () => null,
saveSelectedDirectory,
});
await expect(
authorizer.authorizeSelectedDirectory('/media/iptv')
).resolves.toBe(resolve('/media/iptv'));
expect(saveSelectedDirectory).toHaveBeenCalledWith(
resolve('/media/iptv')
);
await expect(authorizer.requireAuthorized('/media/iptv')).resolves.toBe(
resolve('/media/iptv')
);
});
it('matches authorized paths case-insensitively on Windows', async () => {
const authorizer = new DownloadDirectoryAuthorizer({
getDefaultDirectory: () => 'C:\\Downloads',
loadSelectedDirectory: async () => 'C:\\Media\\IPTV',
saveSelectedDirectory: jest.fn(),
platform: 'win32',
});
await expect(
authorizer.requireAuthorized('c:\\media\\iptv')
).resolves.toBe(resolve('c:\\media\\iptv'));
});
it('keeps authorized path matching case-sensitive on POSIX', async () => {
const authorizer = new DownloadDirectoryAuthorizer({
getDefaultDirectory: () => '/downloads',
loadSelectedDirectory: async () => '/media/IPTV',
saveSelectedDirectory: jest.fn(),
platform: 'linux',
});
await expect(
authorizer.requireAuthorized('/media/iptv')
).rejects.toThrow(/not authorized/i);
});
});
@@ -0,0 +1,80 @@
import { resolve } from 'node:path';
export interface DownloadDirectoryAuthorizerOptions {
getDefaultDirectory: () => string;
loadSelectedDirectory: () => Promise<string | null>;
saveSelectedDirectory: (directory: string) => Promise<void>;
platform?: NodeJS.Platform;
}
function normalizeDirectory(directory: string): string {
const value = String(directory ?? '').trim();
if (!value) {
throw new Error('Download directory is required');
}
return resolve(value);
}
function directoryKey(directory: string, platform: NodeJS.Platform): string {
const normalized = normalizeDirectory(directory);
return platform === 'win32' ? normalized.toLowerCase() : normalized;
}
/**
* Keeps download-directory authorization in the Electron main process.
* Only the OS default or a directory persisted after a native dialog may be
* used by renderer-triggered download IPC calls.
*/
export class DownloadDirectoryAuthorizer {
private selectedDirectory: string | null | undefined;
constructor(private readonly options: DownloadDirectoryAuthorizerOptions) {}
private async loadSelectedDirectory(): Promise<string | null> {
if (this.selectedDirectory !== undefined) {
return this.selectedDirectory;
}
const storedDirectory = await this.options.loadSelectedDirectory();
this.selectedDirectory = storedDirectory
? normalizeDirectory(storedDirectory)
: null;
return this.selectedDirectory;
}
async getPreferredDirectory(): Promise<string> {
return (
(await this.loadSelectedDirectory()) ??
normalizeDirectory(this.options.getDefaultDirectory())
);
}
async authorizeSelectedDirectory(directory: string): Promise<string> {
const normalized = normalizeDirectory(directory);
await this.options.saveSelectedDirectory(normalized);
this.selectedDirectory = normalized;
return normalized;
}
async requireAuthorized(directory: string): Promise<string> {
const normalized = normalizeDirectory(directory);
const defaultDirectory = normalizeDirectory(
this.options.getDefaultDirectory()
);
const selectedDirectory = await this.loadSelectedDirectory();
const platform = this.options.platform ?? process.platform;
const requestedKey = directoryKey(normalized, platform);
if (
requestedKey !== directoryKey(defaultDirectory, platform) &&
(!selectedDirectory ||
requestedKey !== directoryKey(selectedDirectory, platform))
) {
throw new Error(
'Download directory was not authorized by a native folder dialog'
);
}
return normalized;
}
}
@@ -0,0 +1,90 @@
import { join } from 'node:path';
import {
removePartialDownload,
reserveAvailableDownloadFile,
} from './download-file-path';
describe('reserveAvailableDownloadFile', () => {
it('atomically reserves the requested filename when it is unused', () => {
const reserveFile = jest.fn();
expect(
reserveAvailableDownloadFile('/downloads', 'movie.mp4', reserveFile)
).toEqual({
filename: 'movie.mp4',
path: join('/downloads', 'movie.mp4'),
});
expect(reserveFile).toHaveBeenCalledWith(
join('/downloads', 'movie.mp4')
);
});
it('retries with a numbered filename after exclusive-create collisions', () => {
const reserveFile = jest.fn((filePath: string) => {
if (!filePath.endsWith('movie (2).mp4')) {
const error = new Error(
'already exists'
) as NodeJS.ErrnoException;
error.code = 'EEXIST';
throw error;
}
});
expect(
reserveAvailableDownloadFile('/downloads', 'movie.mp4', reserveFile)
).toEqual({
filename: 'movie (2).mp4',
path: join('/downloads', 'movie (2).mp4'),
});
expect(reserveFile).toHaveBeenCalledTimes(3);
});
it('does not hide non-collision filesystem errors', () => {
const error = new Error('permission denied') as NodeJS.ErrnoException;
error.code = 'EACCES';
expect(() =>
reserveAvailableDownloadFile('/downloads', 'movie.mp4', () => {
throw error;
})
).toThrow(error);
});
});
describe('removePartialDownload', () => {
it('removes only the actual partial save path', () => {
const requestedPath = join('/downloads', 'movie.mp4');
const partialPath = join('/downloads', 'movie (1).mp4');
const removeFile = jest.fn();
expect(
removePartialDownload(
{ getSavePath: () => partialPath },
(filePath) => filePath === partialPath,
removeFile
)
).toBe(true);
expect(removeFile).toHaveBeenCalledWith(partialPath);
expect(removeFile).not.toHaveBeenCalledWith(requestedPath);
});
it('does not remove a missing or unavailable partial path', () => {
const removeFile = jest.fn();
expect(
removePartialDownload(
{ getSavePath: () => '' },
() => true,
removeFile
)
).toBe(false);
expect(
removePartialDownload(
{ getSavePath: () => '/downloads/missing.mp4' },
() => false,
removeFile
)
).toBe(false);
expect(removeFile).not.toHaveBeenCalled();
});
});
@@ -0,0 +1,57 @@
import { closeSync, existsSync, openSync, unlinkSync } from 'node:fs';
import { extname, join } from 'node:path';
export interface ReservedDownloadFile {
filename: string;
path: string;
}
function createExclusiveFile(filePath: string): void {
const descriptor = openSync(filePath, 'wx');
closeSync(descriptor);
}
export function reserveAvailableDownloadFile(
directory: string,
requestedFilename: string,
reserveFile: (filePath: string) => void = createExclusiveFile
): ReservedDownloadFile {
const extension = extname(requestedFilename);
const stem = extension
? requestedFilename.slice(0, -extension.length)
: requestedFilename;
let candidate = requestedFilename;
let suffix = 1;
for (;;) {
const filePath = join(directory, candidate);
try {
reserveFile(filePath);
return { filename: candidate, path: filePath };
} catch (error) {
if ((error as NodeJS.ErrnoException).code !== 'EEXIST') {
throw error;
}
candidate = `${stem} (${suffix})${extension}`;
suffix += 1;
}
}
}
export interface DownloadSavePathItem {
getSavePath(): string;
}
export function removePartialDownload(
downloadItem: DownloadSavePathItem | undefined,
pathExists: (filePath: string) => boolean = existsSync,
removeFile: (filePath: string) => void = unlinkSync
): boolean {
const savePath = downloadItem?.getSavePath();
if (!savePath || !pathExists(savePath)) {
return false;
}
removeFile(savePath);
return true;
}
@@ -0,0 +1,48 @@
import { cleanupStaleDownloadFiles } from './stale-download-files';
describe('cleanupStaleDownloadFiles', () => {
it('removes every persisted partial path and ignores empty paths', () => {
const removeFile = jest.fn();
cleanupStaleDownloadFiles(
[
{ filePath: '/downloads/one.mp4' },
{ filePath: null },
{ filePath: '/downloads/two.mp4' },
],
removeFile
);
expect(removeFile).toHaveBeenCalledTimes(2);
expect(removeFile).toHaveBeenNthCalledWith(1, '/downloads/one.mp4');
expect(removeFile).toHaveBeenNthCalledWith(2, '/downloads/two.mp4');
});
it('continues cleaning other stale files after one removal fails', () => {
const removeFile = jest.fn((filePath: string) => {
if (filePath.endsWith('one.mp4')) {
throw new Error('locked');
}
});
const consoleError = jest
.spyOn(console, 'error')
.mockImplementation(() => undefined);
cleanupStaleDownloadFiles(
[
{ filePath: '/downloads/one.mp4' },
{ filePath: '/downloads/two.mp4' },
],
removeFile
);
expect(removeFile).toHaveBeenCalledTimes(2);
expect(consoleError).toHaveBeenCalledWith(
'[Downloads] Failed to delete stale partial file:',
'/downloads/one.mp4',
expect.any(Error)
);
consoleError.mockRestore();
});
});
@@ -0,0 +1,45 @@
import { inArray, sql } from 'drizzle-orm';
import { getDatabase } from '../../database/connection';
import * as schema from '../../database/schema';
import { cleanupStaleDownloadFiles } from './stale-download-files';
export async function resetStaleDownloads(): Promise<void> {
try {
const db = await getDatabase();
const staleDownloads = await db
.select({
filePath: schema.downloads.filePath,
id: schema.downloads.id,
})
.from(schema.downloads)
.where(inArray(schema.downloads.status, ['queued', 'downloading']));
const ownedReservations = staleDownloads.filter(
(item) => item.filePath
);
cleanupStaleDownloadFiles(ownedReservations);
await db
.update(schema.downloads)
.set({
errorMessage: 'Download interrupted by application restart',
status: 'failed',
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(inArray(schema.downloads.status, ['queued', 'downloading']));
if (ownedReservations.length > 0) {
await db
.update(schema.downloads)
.set({ filePath: null })
.where(
inArray(
schema.downloads.id,
ownedReservations.map((item) => item.id)
)
);
}
console.log('[Downloads] Reset stale downloads');
} catch (error) {
console.error('[Downloads] Error resetting stale downloads:', error);
}
}
@@ -0,0 +1,206 @@
import { and, eq, sql } from 'drizzle-orm';
import { extname } from 'node:path';
import { getDatabase } from '../../database/connection';
import * as schema from '../../database/schema';
import { assertRemoteUrlAllowed } from '../url-safety';
import { DownloadDirectoryAuthorizer } from './download-directory-authorization';
import { enqueueDownload } from './download-runtime';
export interface StartDownloadRequest {
playlistId: string;
xtreamId: number;
contentType: 'vod' | 'episode';
title: string;
url: string;
posterUrl?: string;
downloadFolder: string;
headers?: { userAgent?: string; referer?: string; origin?: string };
seriesXtreamId?: number;
seasonNumber?: number;
episodeNumber?: number;
playlistName?: string;
playlistType?: 'xtream' | 'stalker' | 'm3u-file' | 'm3u-text' | 'm3u-url';
serverUrl?: string;
portalUrl?: string;
macAddress?: string;
}
function getExtensionFromUrl(url: string): string {
try {
return extname(new URL(url).pathname) || '.mp4';
} catch {
return '.mp4';
}
}
function sanitizeFilename(name: string): string {
return name.replace(/[<>:"/\\|?*]/g, '_').trim();
}
function createFileName(title: string, url: string): string {
return sanitizeFilename(title) + getExtensionFromUrl(url);
}
function createHeaders(
headers: StartDownloadRequest['headers']
): Record<string, string> | undefined {
return headers
? {
'User-Agent': headers.userAgent || '',
Origin: headers.origin || '',
Referer: headers.referer || '',
}
: undefined;
}
export async function startDownloadRequest(
data: StartDownloadRequest,
authorizer: DownloadDirectoryAuthorizer
): Promise<{ success: boolean; error?: string; id?: number }> {
console.log('[Downloads] Enqueue download:', data.title);
const directory = await authorizer.requireAuthorized(data.downloadFolder);
await assertRemoteUrlAllowed(data.url, { allowPrivateNetworks: true });
const db = await getDatabase();
if (!data.playlistId) {
throw new Error('playlistId is required for downloads');
}
const existingPlaylist = await db
.select()
.from(schema.playlists)
.where(eq(schema.playlists.id, data.playlistId))
.limit(1);
if (existingPlaylist.length === 0) {
console.log(
'[Downloads] Creating playlist entry for:',
data.playlistId
);
await db.insert(schema.playlists).values({
id: data.playlistId,
macAddress: data.macAddress,
name: data.playlistName || 'Unknown Playlist',
serverUrl: data.serverUrl,
type: data.playlistType || 'stalker',
url: data.portalUrl,
});
}
const existing = await db
.select()
.from(schema.downloads)
.where(
and(
eq(schema.downloads.playlistId, data.playlistId),
eq(schema.downloads.xtreamId, data.xtreamId),
eq(schema.downloads.contentType, data.contentType)
)
)
.limit(1);
const fileName = createFileName(data.title, data.url);
if (existing.length > 0) {
const item = existing[0];
if (!['completed', 'failed', 'canceled'].includes(item.status)) {
return {
error: 'Download already in progress',
id: item.id,
success: false,
};
}
await db
.update(schema.downloads)
.set({
bytesDownloaded: 0,
errorMessage: null,
fileName,
filePath: null,
status: 'queued',
totalBytes: null,
updatedAt: sql`CURRENT_TIMESTAMP`,
url: data.url,
})
.where(eq(schema.downloads.id, item.id));
enqueueDownload({
directory,
fileName,
headers: createHeaders(data.headers),
id: item.id,
url: data.url,
});
return { id: item.id, success: true };
}
const result = await db.insert(schema.downloads).values({
contentType: data.contentType,
episodeNumber: data.episodeNumber,
fileName,
playlistId: data.playlistId,
posterUrl: data.posterUrl,
seasonNumber: data.seasonNumber,
seriesXtreamId: data.seriesXtreamId,
status: 'queued',
title: data.title,
url: data.url,
xtreamId: data.xtreamId,
});
const insertedId = Number(result.lastInsertRowid);
enqueueDownload({
directory,
fileName,
headers: createHeaders(data.headers),
id: insertedId,
url: data.url,
});
return { id: insertedId, success: true };
}
export async function retryDownloadRequest(
downloadId: number,
downloadFolder: string,
authorizer: DownloadDirectoryAuthorizer
): Promise<{ success: boolean; error?: string }> {
console.log('[Downloads] Retry download:', downloadId);
const directory = await authorizer.requireAuthorized(downloadFolder);
const db = await getDatabase();
const existing = await db
.select()
.from(schema.downloads)
.where(eq(schema.downloads.id, downloadId))
.limit(1);
if (existing.length === 0) {
return { error: 'Download not found', success: false };
}
const item = existing[0];
await assertRemoteUrlAllowed(item.url, { allowPrivateNetworks: true });
if (!['failed', 'canceled'].includes(item.status)) {
return {
error: 'Can only retry failed or canceled downloads',
success: false,
};
}
const fileName = createFileName(item.title, item.url);
await db
.update(schema.downloads)
.set({
bytesDownloaded: 0,
errorMessage: null,
fileName,
filePath: null,
status: 'queued',
totalBytes: null,
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, downloadId));
enqueueDownload({
directory,
fileName,
id: item.id,
url: item.url,
});
return { success: true };
}
@@ -0,0 +1,326 @@
import type { DownloadItem } from 'electron';
import {
attachDownloadItem,
requestDownloadCancellation,
type DownloadTask,
} from './download-task';
function createTask(): DownloadTask {
return {
directory: '/downloads',
fileName: 'movie.mp4',
id: 42,
url: 'https://example.test/movie.mp4',
};
}
async function waitForCallCount(
mock: jest.Mock,
expectedCallCount: number
): Promise<void> {
for (let attempt = 0; attempt < 20; attempt++) {
if (mock.mock.calls.length === expectedCallCount) {
return;
}
await new Promise<void>((resolve) => setImmediate(resolve));
}
expect(mock).toHaveBeenCalledTimes(expectedCallCount);
}
describe('download runtime cancellation', () => {
it('cancels the item when an earlier cancellation request reaches onStarted', () => {
const task = createTask();
const cancel = jest.fn();
requestDownloadCancellation(task);
attachDownloadItem(task, { cancel } as unknown as DownloadItem);
expect(task.cancelRequested).toBe(true);
expect(cancel).toHaveBeenCalledTimes(1);
});
it('cancels an already-started item immediately', () => {
const cancel = jest.fn();
const task = {
...createTask(),
downloadItem: { cancel } as unknown as DownloadItem,
};
requestDownloadCancellation(task);
expect(task.cancelRequested).toBe(true);
expect(cancel).toHaveBeenCalledTimes(1);
});
it('continues the queue when cancellation persistence fails', async () => {
jest.resetModules();
class TestCancelError extends Error {}
let cancellationCallback: Promise<void> | undefined;
const download = jest
.fn()
.mockImplementationOnce(
async (
_window: unknown,
_url: string,
options: {
onCancel: (item: DownloadItem) => Promise<void>;
}
) => {
cancellationCallback = options.onCancel({
getSavePath: () => '/downloads/movie.mp4',
} as unknown as DownloadItem);
throw new TestCancelError();
}
)
.mockImplementationOnce(
async (
_window: unknown,
_url: string,
options: {
onCompleted: (file: {
fileSize: number;
filename: string;
path: string;
}) => Promise<void>;
}
) => {
await options.onCompleted({
fileSize: 1,
filename: 'second.mp4',
path: '/downloads/second.mp4',
});
}
);
const where = jest
.fn()
.mockResolvedValueOnce(undefined)
.mockResolvedValueOnce(undefined)
.mockRejectedValueOnce(new Error('database is busy'))
.mockResolvedValue(undefined);
const db = {
update: jest.fn(() => ({
set: jest.fn(() => ({ where })),
})),
};
jest.doMock('electron-dl', () => ({
CancelError: TestCancelError,
download,
}));
jest.doMock('../../database/connection', () => ({
getDatabase: jest.fn().mockResolvedValue(db),
}));
jest.doMock('./download-file-path', () => ({
removePartialDownload: jest.fn(),
reserveAvailableDownloadFile: jest.fn(
(directory: string, filename: string) => ({
filename,
path: `${directory}/${filename}`,
})
),
}));
const consoleError = jest
.spyOn(console, 'error')
.mockImplementation(() => undefined);
const runtime = await import('./download-runtime');
runtime.setMainWindow({
isDestroyed: () => false,
webContents: { send: jest.fn() },
} as never);
runtime.enqueueDownload(createTask());
runtime.enqueueDownload({
...createTask(),
fileName: 'second.mp4',
id: 43,
});
await waitForCallCount(download, 2);
await expect(cancellationCallback).resolves.toBeUndefined();
consoleError.mockRestore();
});
it('continues the queue when completion persistence fails', async () => {
jest.resetModules();
class TestCancelError extends Error {}
let completionCallback: Promise<void> | undefined;
const download = jest
.fn()
.mockImplementationOnce(
async (
_window: unknown,
_url: string,
options: {
onCompleted: (file: {
fileSize: number;
filename: string;
path: string;
}) => Promise<void>;
}
) => {
completionCallback = options.onCompleted({
fileSize: 1,
filename: 'movie.mp4',
path: '/downloads/movie.mp4',
});
}
)
.mockImplementationOnce(
async (
_window: unknown,
_url: string,
options: {
onCompleted: (file: {
fileSize: number;
filename: string;
path: string;
}) => Promise<void>;
}
) => {
await options.onCompleted({
fileSize: 1,
filename: 'second.mp4',
path: '/downloads/second.mp4',
});
}
);
const where = jest
.fn()
.mockResolvedValueOnce(undefined)
.mockResolvedValueOnce(undefined)
.mockRejectedValueOnce(new Error('database is busy'))
.mockResolvedValue(undefined);
const db = {
update: jest.fn(() => ({
set: jest.fn(() => ({ where })),
})),
};
jest.doMock('electron-dl', () => ({
CancelError: TestCancelError,
download,
}));
jest.doMock('../../database/connection', () => ({
getDatabase: jest.fn().mockResolvedValue(db),
}));
jest.doMock('./download-file-path', () => ({
removePartialDownload: jest.fn(),
reserveAvailableDownloadFile: jest.fn(
(directory: string, filename: string) => ({
filename,
path: `${directory}/${filename}`,
})
),
}));
const consoleError = jest
.spyOn(console, 'error')
.mockImplementation(() => undefined);
try {
const runtime = await import('./download-runtime');
runtime.setMainWindow({
isDestroyed: () => false,
webContents: { send: jest.fn() },
} as never);
runtime.enqueueDownload(createTask());
runtime.enqueueDownload({
...createTask(),
fileName: 'second.mp4',
id: 43,
});
await waitForCallCount(download, 2);
await expect(completionCallback).resolves.toBeUndefined();
} finally {
consoleError.mockRestore();
}
});
it('continues the queue when initial database access fails', async () => {
jest.resetModules();
class TestCancelError extends Error {}
const download = jest.fn(
async (
_window: unknown,
_url: string,
options: {
onCompleted: (file: {
fileSize: number;
filename: string;
path: string;
}) => Promise<void>;
}
) => {
await options.onCompleted({
fileSize: 1,
filename: 'second.mp4',
path: '/downloads/second.mp4',
});
}
);
const where = jest.fn().mockResolvedValue(undefined);
const db = {
update: jest.fn(() => ({
set: jest.fn(() => ({ where })),
})),
};
const getDatabase = jest
.fn()
.mockRejectedValueOnce(new Error('database unavailable'))
.mockResolvedValue(db);
jest.doMock('electron-dl', () => ({
CancelError: TestCancelError,
download,
}));
jest.doMock('../../database/connection', () => ({ getDatabase }));
jest.doMock('./download-file-path', () => ({
removePartialDownload: jest.fn(),
reserveAvailableDownloadFile: jest.fn(
(directory: string, filename: string) => ({
filename,
path: `${directory}/${filename}`,
})
),
}));
const consoleError = jest
.spyOn(console, 'error')
.mockImplementation(() => undefined);
try {
const runtime = await import('./download-runtime');
runtime.setMainWindow({
isDestroyed: () => false,
webContents: { send: jest.fn() },
} as never);
runtime.enqueueDownload(createTask());
runtime.enqueueDownload({
...createTask(),
fileName: 'second.mp4',
id: 43,
});
await waitForCallCount(download, 1);
expect(download).toHaveBeenCalledWith(
expect.anything(),
'https://example.test/movie.mp4',
expect.objectContaining({ filename: 'second.mp4' })
);
} finally {
consoleError.mockRestore();
}
});
});
@@ -0,0 +1,295 @@
import { eq, sql } from 'drizzle-orm';
import type { BrowserWindow } from 'electron';
import { CancelError, download } from 'electron-dl';
import { getDatabase } from '../../database/connection';
import * as schema from '../../database/schema';
import {
removePartialDownload,
reserveAvailableDownloadFile,
} from './download-file-path';
import {
attachDownloadItem,
requestDownloadCancellation,
type DownloadTask,
} from './download-task';
const downloadQueue: DownloadTask[] = [];
let activeDownload: DownloadTask | null = null;
let mainWindow: BrowserWindow | null = null;
export function setMainWindow(win: BrowserWindow): void {
mainWindow = win;
}
export function broadcastDownloadUpdate(): void {
if (mainWindow && !mainWindow.isDestroyed()) {
mainWindow.webContents.send('DOWNLOADS_UPDATE_EVENT');
}
}
export function enqueueDownload(task: DownloadTask): void {
downloadQueue.push(task);
broadcastDownloadUpdate();
void processQueue();
}
export async function cancelDownload(downloadId: number): Promise<boolean> {
if (activeDownload?.id === downloadId) {
requestDownloadCancellation(activeDownload);
return true;
}
const queueIndex = downloadQueue.findIndex(
(task) => task.id === downloadId
);
if (queueIndex === -1) {
return false;
}
downloadQueue.splice(queueIndex, 1);
const db = await getDatabase();
await db
.update(schema.downloads)
.set({
errorMessage: null,
filePath: null,
status: 'canceled',
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, downloadId));
broadcastDownloadUpdate();
return true;
}
export function removeDownloadFromRuntime(downloadId: number): void {
if (activeDownload?.id === downloadId) {
requestDownloadCancellation(activeDownload);
}
const queueIndex = downloadQueue.findIndex(
(task) => task.id === downloadId
);
if (queueIndex !== -1) {
downloadQueue.splice(queueIndex, 1);
}
}
async function processQueue(): Promise<void> {
if (activeDownload || downloadQueue.length === 0) {
return;
}
const task = downloadQueue.shift();
if (!task) {
return;
}
activeDownload = task;
try {
await startDownload(task);
} catch (error) {
console.error(
`[Downloads] Unhandled error for ${task.fileName}:`,
error
);
finishTask(task);
}
}
function finishTask(task: DownloadTask): void {
if (activeDownload === task) {
activeDownload = null;
}
broadcastDownloadUpdate();
void processQueue();
}
function createCancellationHandler(
task: DownloadTask,
db: Awaited<ReturnType<typeof getDatabase>>,
reservation: ReturnType<typeof reserveAvailableDownloadFile>
): (item: Parameters<typeof removePartialDownload>[0]) => Promise<void> {
let cancellationPromise: Promise<void> | undefined;
return (item) => {
cancellationPromise ??= (async () => {
console.log(`[Downloads] Canceled: ${reservation.filename}`);
removePartialFile(item);
try {
await db
.update(schema.downloads)
.set({
errorMessage: null,
filePath: null,
status: 'canceled',
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
} catch (error) {
console.error(
'[Downloads] Failed to persist cancellation:',
error
);
} finally {
finishTask(task);
}
})();
return cancellationPromise;
};
}
async function startDownload(task: DownloadTask): Promise<void> {
const db = await getDatabase();
await db
.update(schema.downloads)
.set({
errorMessage: null,
filePath: null,
status: 'downloading',
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
broadcastDownloadUpdate();
if (!mainWindow || mainWindow.isDestroyed()) {
console.error('[Downloads] No main window available');
await db
.update(schema.downloads)
.set({
errorMessage: 'No window available for download',
status: 'failed',
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
finishTask(task);
return;
}
let lastProgressUpdate = 0;
const progressThrottleMs = 500;
let handleCancellation:
| ReturnType<typeof createCancellationHandler>
| undefined;
try {
const reservation = reserveAvailableDownloadFile(
task.directory,
task.fileName
);
task.reservedPath = reservation.path;
handleCancellation = createCancellationHandler(task, db, reservation);
await db
.update(schema.downloads)
.set({
errorMessage: null,
fileName: reservation.filename,
filePath: reservation.path,
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
const downloadOptions: Parameters<typeof download>[2] = {
directory: task.directory,
filename: reservation.filename,
onStarted: (item) => {
console.log(`[Downloads] Started: ${reservation.filename}`);
attachDownloadItem(task, item);
},
onProgress: async (progress) => {
const now = Date.now();
if (now - lastProgressUpdate < progressThrottleMs) {
return;
}
lastProgressUpdate = now;
await db
.update(schema.downloads)
.set({
bytesDownloaded: progress.transferredBytes,
totalBytes: progress.totalBytes,
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
broadcastDownloadUpdate();
},
onCompleted: async (file) => {
console.log(`[Downloads] Completed: ${file.filename}`);
try {
await db
.update(schema.downloads)
.set({
bytesDownloaded: file.fileSize,
errorMessage: null,
fileName: file.filename,
filePath: file.path,
status: 'completed',
totalBytes: file.fileSize,
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
} catch (error) {
console.error(
'[Downloads] Failed to persist completion:',
error
);
} finally {
finishTask(task);
}
},
onCancel: handleCancellation,
};
if (task.headers) {
(
downloadOptions as typeof downloadOptions & {
headers: Record<string, string>;
}
).headers = task.headers;
}
await download(mainWindow, task.url, downloadOptions);
} catch (error) {
if (error instanceof CancelError) {
const reservedPath = task.reservedPath;
if (handleCancellation) {
await handleCancellation(
task.downloadItem ??
(reservedPath
? { getSavePath: () => reservedPath }
: undefined)
);
} else {
finishTask(task);
}
return;
}
console.error(`[Downloads] Error downloading ${task.fileName}:`, error);
const reservedPath = task.reservedPath;
removePartialFile(
task.downloadItem ??
(reservedPath ? { getSavePath: () => reservedPath } : undefined)
);
await db
.update(schema.downloads)
.set({
errorMessage:
error instanceof Error ? error.message : String(error),
filePath: null,
status: 'failed',
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
finishTask(task);
}
}
function removePartialFile(
item: Parameters<typeof removePartialDownload>[0]
): void {
try {
removePartialDownload(item);
} catch (error) {
console.error('[Downloads] Failed to delete partial file:', error);
}
}
@@ -0,0 +1,27 @@
import type { DownloadItem } from 'electron';
export interface DownloadTask {
id: number;
url: string;
fileName: string;
directory: string;
headers?: Record<string, string>;
cancelRequested?: boolean;
downloadItem?: DownloadItem;
reservedPath?: string;
}
export function attachDownloadItem(
task: DownloadTask,
item: DownloadItem
): void {
task.downloadItem = item;
if (task.cancelRequested) {
item.cancel();
}
}
export function requestDownloadCancellation(task: DownloadTask): void {
task.cancelRequested = true;
task.downloadItem?.cancel();
}
@@ -1,373 +1,90 @@
/**
* Downloads IPC event handlers
* Uses electron-dl for download management with queue control and progress tracking
*/
import { and, eq, inArray, sql } from 'drizzle-orm';
import { app, BrowserWindow, dialog, ipcMain, shell } from 'electron';
import { download, File as ElectronDlFile } from 'electron-dl';
import { existsSync, unlinkSync } from 'fs';
import { basename, extname, join } from 'path';
import { and, eq, inArray } from 'drizzle-orm';
import { app, dialog, ipcMain, shell } from 'electron';
import { existsSync } from 'node:fs';
import { mkdir, readFile, rename, writeFile } from 'node:fs/promises';
import { join } from 'node:path';
import { getDatabase } from '../../database/connection';
import * as schema from '../../database/schema';
import { DownloadDirectoryAuthorizer } from './download-directory-authorization';
import {
retryDownloadRequest,
startDownloadRequest,
type StartDownloadRequest,
} from './download-requests';
import { resetStaleDownloads } from './download-recovery';
import {
broadcastDownloadUpdate,
cancelDownload,
removeDownloadFromRuntime,
setMainWindow,
} from './download-runtime';
type DownloadStatus = 'queued' | 'downloading' | 'completed' | 'failed' | 'canceled';
interface DownloadTask {
id: number;
url: string;
fileName: string;
directory: string;
headers?: Record<string, string>;
downloadItem?: ElectronDlFile;
function getDownloadAuthorizationPath(): string {
return join(
app.getPath('userData'),
'download-directory-authorization.json'
);
}
// Download queue management
const downloadQueue: DownloadTask[] = [];
let activeDownload: DownloadTask | null = null;
let mainWindow: BrowserWindow | null = null;
/**
* Set the main window reference for sending updates
*/
export function setMainWindow(win: BrowserWindow) {
mainWindow = win;
}
/**
* Broadcast download updates to renderer
*/
function broadcastUpdate() {
if (mainWindow && !mainWindow.isDestroyed()) {
mainWindow.webContents.send('DOWNLOADS_UPDATE_EVENT');
}
}
/**
* Get file extension from URL
*/
function getExtensionFromUrl(url: string): string {
try {
const urlObj = new URL(url);
const pathname = urlObj.pathname;
const ext = extname(pathname);
return ext || '.mp4';
} catch {
return '.mp4';
}
}
/**
* Sanitize filename for filesystem
*/
function sanitizeFilename(name: string): string {
return name.replace(/[<>:"/\\|?*]/g, '_').trim();
}
/**
* Process the download queue - starts next download if none active
*/
async function processQueue() {
if (activeDownload || downloadQueue.length === 0) {
return;
}
const task = downloadQueue.shift();
if (!task) return;
activeDownload = task;
await startDownload(task);
}
/**
* Start a download using electron-dl
*/
async function startDownload(task: DownloadTask) {
const db = await getDatabase();
// Update status to downloading
await db
.update(schema.downloads)
.set({
status: 'downloading',
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
broadcastUpdate();
if (!mainWindow || mainWindow.isDestroyed()) {
console.error('[Downloads] No main window available');
await db
.update(schema.downloads)
.set({
status: 'failed',
errorMessage: 'No window available for download',
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
activeDownload = null;
broadcastUpdate();
processQueue();
return;
}
let lastProgressUpdate = 0;
const PROGRESS_THROTTLE_MS = 500;
try {
const downloadOptions: Parameters<typeof download>[2] = {
directory: task.directory,
filename: task.fileName,
overwrite: true,
onStarted: (item) => {
console.log(`[Downloads] Started: ${task.fileName}`);
task.downloadItem = item as unknown as ElectronDlFile;
},
onProgress: async (progress) => {
const now = Date.now();
if (now - lastProgressUpdate < PROGRESS_THROTTLE_MS) {
return;
}
lastProgressUpdate = now;
await db
.update(schema.downloads)
.set({
bytesDownloaded: progress.transferredBytes,
totalBytes: progress.totalBytes,
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
broadcastUpdate();
},
onCompleted: async (file) => {
console.log(`[Downloads] Completed: ${task.fileName}`);
await db
.update(schema.downloads)
.set({
status: 'completed',
filePath: file.path,
fileName: file.filename,
bytesDownloaded: file.fileSize,
totalBytes: file.fileSize,
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
activeDownload = null;
broadcastUpdate();
processQueue();
},
onCancel: async () => {
console.log(`[Downloads] Canceled: ${task.fileName}`);
// Delete partial file if it exists
const partialPath = join(task.directory, task.fileName);
if (existsSync(partialPath)) {
try {
unlinkSync(partialPath);
} catch (e) {
console.error('[Downloads] Failed to delete partial file:', e);
}
}
await db
.update(schema.downloads)
.set({
status: 'canceled',
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
activeDownload = null;
broadcastUpdate();
processQueue();
},
};
// Add headers if provided (user-agent, referer, origin)
if (task.headers) {
(downloadOptions as any).headers = task.headers;
}
await download(mainWindow, task.url, downloadOptions);
} catch (error) {
console.error(`[Downloads] Error downloading ${task.fileName}:`, error);
// Delete partial file if it exists
const partialPath = join(task.directory, task.fileName);
if (existsSync(partialPath)) {
try {
unlinkSync(partialPath);
} catch (e) {
console.error('[Downloads] Failed to delete partial file:', e);
const downloadDirectoryAuthorizer = new DownloadDirectoryAuthorizer({
getDefaultDirectory: () => app.getPath('downloads'),
loadSelectedDirectory: async () => {
try {
const stored = JSON.parse(
await readFile(getDownloadAuthorizationPath(), 'utf-8')
) as { directory?: unknown };
return typeof stored.directory === 'string'
? stored.directory
: null;
} catch (error) {
if ((error as NodeJS.ErrnoException).code !== 'ENOENT') {
console.warn(
'[Downloads] Ignoring invalid folder authorization:',
error
);
}
return null;
}
},
saveSelectedDirectory: async (directory) => {
const authorizationPath = getDownloadAuthorizationPath();
const temporaryPath = `${authorizationPath}.${process.pid}.tmp`;
await mkdir(app.getPath('userData'), { recursive: true });
await writeFile(
temporaryPath,
JSON.stringify({ directory, version: 1 }),
'utf-8'
);
await rename(temporaryPath, authorizationPath);
},
});
await db
.update(schema.downloads)
.set({
status: 'failed',
errorMessage: error instanceof Error ? error.message : String(error),
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, task.id));
activeDownload = null;
broadcastUpdate();
processQueue();
async function isManagedDownloadFile(filePath: string): Promise<boolean> {
if (!filePath) {
return false;
}
try {
const db = await getDatabase();
const rows = await db
.select({ id: schema.downloads.id })
.from(schema.downloads)
.where(eq(schema.downloads.filePath, filePath))
.limit(1);
return rows.length > 0;
} catch (error) {
console.error('Error verifying managed download path:', error);
return false;
}
}
/**
* Enqueue a new download
*/
ipcMain.handle(
'DOWNLOADS_START',
async (
_event,
data: {
playlistId: string;
xtreamId: number;
contentType: 'vod' | 'episode';
title: string;
url: string;
posterUrl?: string;
downloadFolder: string;
headers?: { userAgent?: string; referer?: string; origin?: string };
seriesXtreamId?: number;
seasonNumber?: number;
episodeNumber?: number;
// Playlist info for auto-creation if needed
playlistName?: string;
playlistType?: 'xtream' | 'stalker' | 'm3u-file' | 'm3u-text' | 'm3u-url';
serverUrl?: string;
portalUrl?: string;
macAddress?: string;
}
) => {
async (_event, data: StartDownloadRequest) => {
try {
console.log('[Downloads] Enqueue download:', data.title);
const db = await getDatabase();
// Ensure playlist exists in database (required for foreign key constraint)
if (data.playlistId) {
const existingPlaylist = await db
.select()
.from(schema.playlists)
.where(eq(schema.playlists.id, data.playlistId))
.limit(1);
if (existingPlaylist.length === 0) {
// Create playlist entry for downloads to work
console.log('[Downloads] Creating playlist entry for:', data.playlistId);
await db.insert(schema.playlists).values({
id: data.playlistId,
name: data.playlistName || 'Unknown Playlist',
type: data.playlistType || 'stalker',
serverUrl: data.serverUrl,
macAddress: data.macAddress,
url: data.portalUrl,
});
}
} else {
throw new Error('playlistId is required for downloads');
}
// Check if already exists
const existing = await db
.select()
.from(schema.downloads)
.where(
and(
eq(schema.downloads.playlistId, data.playlistId),
eq(schema.downloads.xtreamId, data.xtreamId),
eq(schema.downloads.contentType, data.contentType)
)
)
.limit(1);
if (existing.length > 0) {
const item = existing[0];
// If completed or failed/canceled, allow retry by updating
if (['completed', 'failed', 'canceled'].includes(item.status)) {
await db
.update(schema.downloads)
.set({
status: 'queued',
url: data.url,
bytesDownloaded: 0,
totalBytes: null,
errorMessage: null,
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, item.id));
const ext = getExtensionFromUrl(data.url);
const fileName = sanitizeFilename(data.title) + ext;
downloadQueue.push({
id: item.id,
url: data.url,
fileName,
directory: data.downloadFolder,
headers: data.headers
? {
'User-Agent': data.headers.userAgent || '',
Referer: data.headers.referer || '',
Origin: data.headers.origin || '',
}
: undefined,
});
broadcastUpdate();
processQueue();
return { success: true, id: item.id };
}
// Already queued or downloading
return { success: false, error: 'Download already in progress', id: item.id };
}
// Create new download entry
const ext = getExtensionFromUrl(data.url);
const fileName = sanitizeFilename(data.title) + ext;
const result = await db.insert(schema.downloads).values({
playlistId: data.playlistId,
xtreamId: data.xtreamId,
contentType: data.contentType,
title: data.title,
url: data.url,
posterUrl: data.posterUrl,
fileName,
status: 'queued',
seriesXtreamId: data.seriesXtreamId,
seasonNumber: data.seasonNumber,
episodeNumber: data.episodeNumber,
});
const insertedId = Number(result.lastInsertRowid);
downloadQueue.push({
id: insertedId,
url: data.url,
fileName,
directory: data.downloadFolder,
headers: data.headers
? {
'User-Agent': data.headers.userAgent || '',
Referer: data.headers.referer || '',
Origin: data.headers.origin || '',
}
: undefined,
});
broadcastUpdate();
processQueue();
return { success: true, id: insertedId };
return await startDownloadRequest(
data,
downloadDirectoryAuthorizer
);
} catch (error) {
console.error('[Downloads] Error enqueuing download:', error);
throw error;
@@ -375,95 +92,27 @@ ipcMain.handle(
}
);
/**
* Cancel a download
*/
ipcMain.handle('DOWNLOADS_CANCEL', async (_event, downloadId: number) => {
try {
console.log('[Downloads] Cancel download:', downloadId);
const db = await getDatabase();
// Check if it's the active download
if (activeDownload && activeDownload.id === downloadId) {
if (activeDownload.downloadItem) {
(activeDownload.downloadItem as any).cancel?.();
}
return { success: true };
}
// Remove from queue
const queueIndex = downloadQueue.findIndex((t) => t.id === downloadId);
if (queueIndex !== -1) {
downloadQueue.splice(queueIndex, 1);
await db
.update(schema.downloads)
.set({
status: 'canceled',
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, downloadId));
broadcastUpdate();
return { success: true };
}
return { success: false, error: 'Download not found in queue' };
return (await cancelDownload(downloadId))
? { success: true }
: { error: 'Download not found in queue', success: false };
} catch (error) {
console.error('[Downloads] Error canceling download:', error);
throw error;
}
});
/**
* Retry a failed/canceled download
*/
ipcMain.handle(
'DOWNLOADS_RETRY',
async (_event, downloadId: number, downloadFolder: string) => {
try {
console.log('[Downloads] Retry download:', downloadId);
const db = await getDatabase();
const existing = await db
.select()
.from(schema.downloads)
.where(eq(schema.downloads.id, downloadId))
.limit(1);
if (existing.length === 0) {
return { success: false, error: 'Download not found' };
}
const item = existing[0];
if (!['failed', 'canceled'].includes(item.status)) {
return { success: false, error: 'Can only retry failed or canceled downloads' };
}
await db
.update(schema.downloads)
.set({
status: 'queued',
bytesDownloaded: 0,
totalBytes: null,
errorMessage: null,
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(eq(schema.downloads.id, downloadId));
const ext = getExtensionFromUrl(item.url);
const fileName = sanitizeFilename(item.title) + ext;
downloadQueue.push({
id: item.id,
url: item.url,
fileName,
directory: downloadFolder,
});
broadcastUpdate();
processQueue();
return { success: true };
return await retryDownloadRequest(
downloadId,
downloadFolder,
downloadDirectoryAuthorizer
);
} catch (error) {
console.error('[Downloads] Error retrying download:', error);
throw error;
@@ -471,31 +120,15 @@ ipcMain.handle(
}
);
/**
* Remove a download from the list
*/
ipcMain.handle('DOWNLOADS_REMOVE', async (_event, downloadId: number) => {
try {
console.log('[Downloads] Remove download:', downloadId);
removeDownloadFromRuntime(downloadId);
const db = await getDatabase();
// Cancel if active
if (activeDownload && activeDownload.id === downloadId) {
if (activeDownload.downloadItem) {
(activeDownload.downloadItem as any).cancel?.();
}
}
// Remove from queue
const queueIndex = downloadQueue.findIndex((t) => t.id === downloadId);
if (queueIndex !== -1) {
downloadQueue.splice(queueIndex, 1);
}
// Delete from database
await db.delete(schema.downloads).where(eq(schema.downloads.id, downloadId));
broadcastUpdate();
await db
.delete(schema.downloads)
.where(eq(schema.downloads.id, downloadId));
broadcastDownloadUpdate();
return { success: true };
} catch (error) {
console.error('[Downloads] Error removing download:', error);
@@ -503,36 +136,21 @@ ipcMain.handle('DOWNLOADS_REMOVE', async (_event, downloadId: number) => {
}
});
/**
* Get all downloads for a playlist
*/
ipcMain.handle('DOWNLOADS_GET_LIST', async (_event, playlistId?: string) => {
try {
const db = await getDatabase();
if (playlistId) {
const result = await db
.select()
.from(schema.downloads)
.where(eq(schema.downloads.playlistId, playlistId))
.orderBy(schema.downloads.createdAt);
return result;
}
const result = await db
.select()
.from(schema.downloads)
.orderBy(schema.downloads.createdAt);
return result;
const query = db.select().from(schema.downloads);
return playlistId
? query
.where(eq(schema.downloads.playlistId, playlistId))
.orderBy(schema.downloads.createdAt)
: query.orderBy(schema.downloads.createdAt);
} catch (error) {
console.error('[Downloads] Error getting download list:', error);
throw error;
}
});
/**
* Get download by ID
*/
ipcMain.handle('DOWNLOADS_GET', async (_event, downloadId: number) => {
try {
const db = await getDatabase();
@@ -548,102 +166,67 @@ ipcMain.handle('DOWNLOADS_GET', async (_event, downloadId: number) => {
}
});
/**
* Get default download folder
*/
ipcMain.handle('DOWNLOADS_GET_DEFAULT_FOLDER', async () => {
return app.getPath('downloads');
return downloadDirectoryAuthorizer.getPreferredDirectory();
});
/**
* Select download folder via dialog
*/
ipcMain.handle('DOWNLOADS_SELECT_FOLDER', async () => {
const result = await dialog.showOpenDialog({
defaultPath: app.getPath('downloads'),
properties: ['openDirectory', 'createDirectory'],
title: 'Select Download Folder',
defaultPath: app.getPath('downloads'),
});
if (result.canceled || result.filePaths.length === 0) {
return null;
}
return result.filePaths[0];
return downloadDirectoryAuthorizer.authorizeSelectedDirectory(
result.filePaths[0]
);
});
/**
* Reveal file in system file manager
*/
ipcMain.handle('DOWNLOADS_REVEAL_FILE', async (_event, filePath: string) => {
if (existsSync(filePath)) {
shell.showItemInFolder(filePath);
return { success: true };
if (!(await isManagedDownloadFile(filePath)) || !existsSync(filePath)) {
return { error: 'File not found', success: false };
}
return { success: false, error: 'File not found' };
shell.showItemInFolder(filePath);
return { success: true };
});
/**
* Play downloaded file
*/
ipcMain.handle('DOWNLOADS_PLAY_FILE', async (_event, filePath: string) => {
if (existsSync(filePath)) {
await shell.openPath(filePath);
return { success: true };
if (!(await isManagedDownloadFile(filePath)) || !existsSync(filePath)) {
return { error: 'File not found', success: false };
}
return { success: false, error: 'File not found' };
await shell.openPath(filePath);
return { success: true };
});
/**
* Clear all completed downloads
*/
ipcMain.handle('DOWNLOADS_CLEAR_COMPLETED', async (_event, playlistId?: string) => {
try {
const db = await getDatabase();
if (playlistId) {
ipcMain.handle(
'DOWNLOADS_CLEAR_COMPLETED',
async (_event, playlistId?: string) => {
try {
const db = await getDatabase();
const terminalStatus = inArray(schema.downloads.status, [
'completed',
'failed',
'canceled',
]);
await db
.delete(schema.downloads)
.where(
and(
eq(schema.downloads.playlistId, playlistId),
inArray(schema.downloads.status, ['completed', 'failed', 'canceled'])
)
);
} else {
await db
.delete(schema.downloads)
.where(
inArray(schema.downloads.status, ['completed', 'failed', 'canceled'])
playlistId
? and(
eq(schema.downloads.playlistId, playlistId),
terminalStatus
)
: terminalStatus
);
broadcastDownloadUpdate();
return { success: true };
} catch (error) {
console.error('[Downloads] Error clearing completed:', error);
throw error;
}
broadcastUpdate();
return { success: true };
} catch (error) {
console.error('[Downloads] Error clearing completed:', error);
throw error;
}
});
);
/**
* Reset stale downloads on startup (downloading -> failed)
*/
export async function resetStaleDownloads() {
try {
const db = await getDatabase();
await db
.update(schema.downloads)
.set({
status: 'failed',
errorMessage: 'Download interrupted by application restart',
updatedAt: sql`CURRENT_TIMESTAMP`,
})
.where(
inArray(schema.downloads.status, ['queued', 'downloading'])
);
console.log('[Downloads] Reset stale downloads');
} catch (error) {
console.error('[Downloads] Error resetting stale downloads:', error);
}
}
export { resetStaleDownloads, setMainWindow };
@@ -0,0 +1,29 @@
import { removePartialDownload } from './download-file-path';
interface StaleDownloadFile {
filePath: string | null;
}
function removePersistedPartial(filePath: string): void {
removePartialDownload({ getSavePath: () => filePath });
}
export function cleanupStaleDownloadFiles(
downloads: readonly StaleDownloadFile[],
removeFile: (filePath: string) => void = removePersistedPartial
): void {
for (const download of downloads) {
if (!download.filePath) {
continue;
}
try {
removeFile(download.filePath);
} catch (error) {
console.error(
'[Downloads] Failed to delete stale partial file:',
download.filePath,
error
);
}
}
}
@@ -0,0 +1,63 @@
import type { Playlist } from '@iptvnator/shared/interfaces';
import {
createPlaylistObject,
getFilenameFromUrl,
} from '@iptvnator/shared/m3u-utils';
import { parse } from 'iptv-playlist-parser';
import { readFile } from 'node:fs/promises';
import { basename } from 'node:path';
import { createPlaylistAgentFactory } from '../util/secure-https';
import { requestWithValidatedRedirects } from '../util/validated-axios';
export async function fetchPlaylistFromUrl(
url: string,
title?: string
): Promise<Playlist> {
const result = await requestWithValidatedRedirects<string>(
url,
{
agentFactory: createPlaylistAgentFactory(),
method: 'GET',
},
{ allowPrivateNetworks: true }
);
const parsedPlaylist = parse(result.data);
const extractedName = url && url.length > 1 ? getFilenameFromUrl(url) : '';
const playlistName =
!extractedName || extractedName === 'Untitled playlist'
? 'Imported from URL'
: extractedName;
return createPlaylistObject(
title ?? playlistName,
parsedPlaylist,
url,
'URL'
);
}
export async function fetchPlaylistFromFile(
filePath: string,
title: string
): Promise<Playlist> {
const fileContent = await readFile(filePath, 'utf-8');
return createPlaylistObject(title, parse(fileContent), filePath, 'FILE');
}
export function derivePlaylistTitleFromFilePath(filePath: string): string {
const filename = basename(filePath);
return filename.replace(/\.(m3u8?|pls|txt)$/i, '') || 'from file';
}
export function preserveAutoUpdatedPlaylistFields(
playlistObject: Playlist,
playlist: Playlist
): Playlist {
return {
...playlistObject,
_id: playlist._id,
autoRefresh: playlist.autoRefresh,
favorites: playlist.favorites || [],
userAgent: playlist.userAgent,
};
}
@@ -0,0 +1,36 @@
import { resolve } from 'node:path';
const MAX_AUTHORIZED_WRITE_PATHS = 32;
export class PlaylistWriteAuthorizer {
private readonly pathsBySender = new Map<number, Set<string>>();
authorize(senderId: number, filePath: string): void {
const senderPaths =
this.pathsBySender.get(senderId) ?? new Set<string>();
senderPaths.add(resolve(filePath));
this.pathsBySender.set(senderId, senderPaths);
while (senderPaths.size > MAX_AUTHORIZED_WRITE_PATHS) {
const oldest = senderPaths.values().next().value;
if (oldest === undefined) {
break;
}
senderPaths.delete(oldest);
}
}
consume(senderId: number, filePath: string): string {
const normalizedPath = resolve(String(filePath ?? ''));
const senderPaths = this.pathsBySender.get(senderId);
if (!senderPaths?.delete(normalizedPath)) {
throw new Error(
'Refusing to write to a path not authorized by a save dialog'
);
}
if (senderPaths.size === 0) {
this.pathsBySender.delete(senderId);
}
return normalizedPath;
}
}
@@ -9,11 +9,13 @@ import {
PLAYLIST_REFRESH,
PLAYLIST_REFRESH_EVENT,
} from '@iptvnator/shared/interfaces';
import { resolve } from 'node:path';
type IpcHandler = (event: MockIpcEvent, ...args: unknown[]) => Promise<unknown>;
type MockIpcEvent = {
sender: {
id: number;
isDestroyed: jest.Mock<boolean, []>;
send: jest.Mock;
};
@@ -57,9 +59,7 @@ jest.mock('electron', () => ({
jest.mock('axios', () => ({
__esModule: true,
default: {
get: (...args: unknown[]) => mockAxiosGet(...args),
},
default: (...args: unknown[]) => mockAxiosGet(...args),
}));
jest.mock('iptv-playlist-parser', () => ({
@@ -118,9 +118,10 @@ function createPlaylist(overrides: Partial<Playlist> = {}): Playlist {
};
}
function createIpcEvent(): MockIpcEvent {
function createIpcEvent(senderId = 1): MockIpcEvent {
return {
sender: {
id: senderId,
isDestroyed: jest.fn(() => false),
send: jest.fn(),
},
@@ -190,8 +191,12 @@ describe('playlist IPC events', () => {
);
expect(mockAxiosGet).toHaveBeenCalledWith(
'https://example.test/remote.m3u',
{ httpsAgent: expect.any(Object) }
expect.objectContaining({
httpsAgent: expect.any(Object),
maxRedirects: 0,
method: 'GET',
url: 'https://example.test/remote.m3u',
})
);
expect(mockParse).toHaveBeenCalledWith('#EXTM3U');
expect(mockCreatePlaylistObject).toHaveBeenCalledWith(
@@ -323,8 +328,12 @@ describe('playlist IPC events', () => {
}),
]);
expect(mockAxiosGet).toHaveBeenCalledWith(
'https://example.test/list.m3u',
{ httpsAgent: expect.any(Object) }
expect.objectContaining({
httpsAgent: expect.any(Object),
maxRedirects: 0,
method: 'GET',
url: 'https://example.test/list.m3u',
})
);
expect(mockReadFile).toHaveBeenCalledWith(
'/playlists/local.m3u',
@@ -521,9 +530,63 @@ describe('playlist IPC events', () => {
)
).resolves.toEqual({ success: true });
expect(mockWriteFile).toHaveBeenCalledWith(
'/exports/list.m3u',
resolve('/exports/list.m3u'),
'#EXTM3U',
'utf-8'
);
});
it('rejects write-file for a path not authorized by a save dialog', async () => {
mockWriteFile.mockClear();
await expect(
getHandler('write-file')(
createIpcEvent(),
'/unauthorized/evil.sh',
'payload'
)
).rejects.toThrow(/not authorized/i);
expect(mockWriteFile).not.toHaveBeenCalled();
});
it('scopes save-dialog write authorization to the requesting renderer', async () => {
mockShowSaveDialog.mockResolvedValue({
canceled: false,
filePath: '/exports/private.m3u',
});
await getHandler('save-file-dialog')(
createIpcEvent(1),
'/exports/private.m3u',
[]
);
await expect(
getHandler('write-file')(
createIpcEvent(2),
'/exports/private.m3u',
'#EXTM3U'
)
).rejects.toThrow(/not authorized/i);
expect(mockWriteFile).not.toHaveBeenCalled();
});
it('consumes write authorization even when the filesystem write fails', async () => {
mockShowSaveDialog.mockResolvedValue({
canceled: false,
filePath: '/exports/failure.m3u',
});
mockWriteFile.mockRejectedValueOnce(new Error('disk full'));
const event = createIpcEvent(3);
await getHandler('save-file-dialog')(event, '/exports/failure.m3u', []);
await expect(
getHandler('write-file')(event, '/exports/failure.m3u', '#EXTM3U')
).rejects.toThrow('disk full');
await expect(
getHandler('write-file')(event, '/exports/failure.m3u', '#EXTM3U')
).rejects.toThrow(/not authorized/i);
expect(mockWriteFile).toHaveBeenCalledTimes(1);
});
});
@@ -3,12 +3,8 @@
* between the frontend to the electron backend.
*/
import axios from 'axios';
import { app, dialog, ipcMain, WebContents } from 'electron';
import { parse } from 'iptv-playlist-parser';
import { createPlaylistObject, getFilenameFromUrl } from '@iptvnator/shared/m3u-utils';
import { readFile, writeFile } from 'node:fs/promises';
import { basename } from 'node:path';
import { writeFile } from 'node:fs/promises';
import { pathToFileURL } from 'url';
import { Worker } from 'worker_threads';
import {
@@ -25,6 +21,13 @@ import type {
PlaylistRefreshWorkerMessage,
PlaylistRefreshWorkerResponseMessage,
} from '../workers/playlist-refresh.worker.types';
import {
derivePlaylistTitleFromFilePath,
fetchPlaylistFromFile,
fetchPlaylistFromUrl,
preserveAutoUpdatedPlaylistFields,
} from './playlist-source';
import { PlaylistWriteAuthorizer } from './playlist-write-authorization';
export default class PlaylistEvents {
static bootstrapPlaylistEvents(): Electron.IpcMain {
@@ -32,7 +35,7 @@ export default class PlaylistEvents {
}
}
const https = require('https');
const playlistWriteAuthorizer = new PlaylistWriteAuthorizer();
type ActivePlaylistRefresh = {
reject: (reason?: unknown) => void;
@@ -43,68 +46,13 @@ type ActivePlaylistRefresh = {
const activePlaylistRefreshes = new Map<string, ActivePlaylistRefresh>();
/**
* Fetches and parses a playlist from a URL
* @param url - The URL to fetch the playlist from
* @param title - Optional title for the playlist
* @returns Parsed playlist object
*/
async function fetchPlaylistFromUrl(
url: string,
title?: string
): Promise<Playlist> {
const agent = new https.Agent({
rejectUnauthorized: false,
});
const result = await axios.get(url, { httpsAgent: agent });
const parsedPlaylist = parse(result.data);
const extractedName =
url && url.length > 1 ? getFilenameFromUrl(url) : '';
const playlistName =
!extractedName || extractedName === 'Untitled playlist'
? 'Imported from URL'
: extractedName;
const playlistObject = createPlaylistObject(
title ?? playlistName,
parsedPlaylist,
url,
'URL'
);
return playlistObject;
}
/**
* Reads and parses a playlist from a file path
* @param filePath - The path to the playlist file
* @param title - Title for the playlist
* @returns Parsed playlist object
*/
async function fetchPlaylistFromFile(
filePath: string,
title: string
): Promise<Playlist> {
const fileContent = await readFile(filePath, 'utf-8');
const parsedPlaylist = parse(fileContent);
const playlistObject = createPlaylistObject(
title,
parsedPlaylist,
filePath,
'FILE'
);
return playlistObject;
}
function resolvePlaylistRefreshWorker(): Worker {
const bootstrap = resolveWorkerRuntimeBootstrap({
isPackaged: app.isPackaged,
workerFilename: 'playlist-refresh.worker.js',
developmentWorkerDir: __dirname + '/workers',
resourcesPath: (
process as NodeJS.Process & { resourcesPath?: string }
).resourcesPath,
resourcesPath: (process as NodeJS.Process & { resourcesPath?: string })
.resourcesPath,
appPath: app.getAppPath(),
});
@@ -137,24 +85,6 @@ function createPlaylistRefreshError(error: {
return workerError;
}
function derivePlaylistTitleFromFilePath(filePath: string): string {
const filename = basename(filePath);
return filename.replace(/\.(m3u8?|pls|txt)$/i, '') || 'from file';
}
function preserveAutoUpdatedPlaylistFields(
playlistObject: Playlist,
playlist: Playlist
): Playlist {
return {
...playlistObject,
_id: playlist._id,
autoRefresh: playlist.autoRefresh,
favorites: playlist.favorites || [],
userAgent: playlist.userAgent,
};
}
ipcMain.handle('fetch-playlist-by-url', async (event, url, title?: string) => {
try {
return await fetchPlaylistFromUrl(url, title);
@@ -254,75 +184,83 @@ ipcMain.handle(AUTO_UPDATE_PLAYLISTS, async (event, playlists) => {
return updatedPlaylists;
});
ipcMain.handle(PLAYLIST_REFRESH, async (event, payload: PlaylistRefreshPayload) => {
const worker = resolvePlaylistRefreshWorker();
ipcMain.handle(
PLAYLIST_REFRESH,
async (event, payload: PlaylistRefreshPayload) => {
const worker = resolvePlaylistRefreshWorker();
return await new Promise<Playlist>((resolve, reject) => {
const cleanup = async (): Promise<void> => {
activePlaylistRefreshes.delete(payload.operationId);
worker.removeAllListeners();
await worker.terminate().catch(() => undefined);
};
return await new Promise<Playlist>((resolve, reject) => {
const cleanup = async (): Promise<void> => {
activePlaylistRefreshes.delete(payload.operationId);
worker.removeAllListeners();
await worker.terminate().catch(() => undefined);
};
activePlaylistRefreshes.set(payload.operationId, {
worker,
sender: event.sender,
resolve,
reject,
});
activePlaylistRefreshes.set(payload.operationId, {
worker,
sender: event.sender,
resolve,
reject,
});
worker.on('message', async (message: PlaylistRefreshWorkerMessage<Playlist>) => {
if (message.type === 'ready') {
worker.postMessage({
type: 'request',
payload,
});
return;
}
if (message.type === 'event') {
emitPlaylistRefreshEvent(event.sender, message.event);
return;
}
await cleanup();
const response = message as PlaylistRefreshWorkerResponseMessage<Playlist>;
if (response.success && response.result) {
resolve(response.result);
return;
}
reject(
createPlaylistRefreshError(
response.error ?? {
message: 'Playlist refresh worker request failed',
worker.on(
'message',
async (message: PlaylistRefreshWorkerMessage<Playlist>) => {
if (message.type === 'ready') {
worker.postMessage({
type: 'request',
payload,
});
return;
}
)
if (message.type === 'event') {
emitPlaylistRefreshEvent(event.sender, message.event);
return;
}
await cleanup();
const response =
message as PlaylistRefreshWorkerResponseMessage<Playlist>;
if (response.success && response.result) {
resolve(response.result);
return;
}
reject(
createPlaylistRefreshError(
response.error ?? {
message:
'Playlist refresh worker request failed',
}
)
);
}
);
});
worker.on('error', async (error) => {
await cleanup();
reject(error);
});
worker.on('error', async (error) => {
await cleanup();
reject(error);
});
worker.on('exit', async (code) => {
if (!activePlaylistRefreshes.has(payload.operationId)) {
return;
}
worker.on('exit', async (code) => {
if (!activePlaylistRefreshes.has(payload.operationId)) {
return;
}
await cleanup();
reject(
new Error(
code === 0
? 'Playlist refresh worker exited unexpectedly'
: `Playlist refresh worker stopped with exit code ${code}`
)
);
await cleanup();
reject(
new Error(
code === 0
? 'Playlist refresh worker exited unexpectedly'
: `Playlist refresh worker stopped with exit code ${code}`
)
);
});
});
});
});
}
);
ipcMain.handle(
PLAYLIST_CANCEL_REFRESH,
@@ -345,15 +283,14 @@ ipcMain.handle('save-file-dialog', async (event, defaultPath, filters) => {
try {
const { canceled, filePath } = await dialog.showSaveDialog({
defaultPath,
filters: filters || [
{ name: 'All Files', extensions: ['*'] },
],
filters: filters || [{ name: 'All Files', extensions: ['*'] }],
});
if (canceled || !filePath) {
return null;
}
playlistWriteAuthorizer.authorize(event.sender.id, filePath);
return filePath;
} catch (error) {
console.error('Error showing save dialog:', error);
@@ -362,8 +299,18 @@ ipcMain.handle('save-file-dialog', async (event, defaultPath, filters) => {
});
ipcMain.handle('write-file', async (event, filePath, content) => {
let normalizedPath: string;
try {
await writeFile(filePath, content, 'utf-8');
normalizedPath = playlistWriteAuthorizer.consume(
event.sender.id,
filePath
);
} catch (error) {
console.error('Blocked unauthorized write-file path:', filePath);
throw error;
}
try {
await writeFile(normalizedPath, content, 'utf-8');
return { success: true };
} catch (error) {
console.error('Error writing file:', error);
@@ -5,10 +5,15 @@
import axios, { AxiosRequestConfig } from 'axios';
import { ipcMain } from 'electron';
import { PortalDebugEvent, STALKER_REQUEST } from '@iptvnator/shared/interfaces';
import {
PortalDebugEvent,
STALKER_REQUEST,
} from '@iptvnator/shared/interfaces';
import { rememberStalkerPlaybackContext } from '../services/stalker-playback-context.service';
import { emitPortalDebugEvent } from './portal-debug.events';
import { buildStalkerIdentityRequestContext } from './stalker-identity';
import { assertRemoteUrlAllowed } from './url-safety';
import { requestWithValidatedRedirects } from '../util/validated-axios';
export default class StalkerEvents {
static bootstrapStalkerEvents(): Electron.IpcMain {
@@ -16,6 +21,13 @@ export default class StalkerEvents {
}
}
type StalkerResponseData = {
js?: {
cmd?: unknown;
};
[key: string]: unknown;
};
/**
* Handle Stalker API requests with MAC address cookie and optional Bearer token
*/
@@ -48,14 +60,24 @@ ipcMain.handle(
// Build URL with query parameters
// Note: For 'cmd' parameter, we need to use encodeURI (not encodeURIComponent)
// to preserve forward slashes, matching stalker-to-m3u implementation
// SSRF/LFI guard: block non-http(s)/credentialed portal URLs.
// Private/LAN targets remain allowed (users run local Stalker servers).
await assertRemoteUrlAllowed(url, { allowPrivateNetworks: true });
const urlObject = new URL(url);
const queryParts: string[] = [];
Object.entries(requestParams).forEach(([key, value]) => {
if (key === 'cmd') {
// Don't encode cmd - it's already a path like /media/12345.mpg
// Encoding would break the path format expected by the server
queryParts.push(`${key}=${String(value)}`);
// Encode cmd but preserve forward slashes so the path format
// (e.g. /media/12345.mpg) the server expects still survives.
// Encoding the remaining characters prevents a malicious portal
// from injecting extra query parameters (&, =, #) into the URL.
queryParts.push(
`${key}=${encodeURIComponent(String(value)).replace(
/%2F/gi,
'/'
)}`
);
} else {
// Use encodeURIComponent for other params
queryParts.push(
@@ -93,7 +115,12 @@ ipcMain.handle(
params: requestParams,
};
const response = await axios(config);
const response =
await requestWithValidatedRedirects<StalkerResponseData>(
fullUrl,
config,
{ allowPrivateNetworks: true }
);
// Check if response is successful
if (response.status >= 400) {
@@ -0,0 +1,124 @@
import {
assertRemoteUrlAllowed,
isPrivateOrReservedIp,
isPrivateNetworkUrlAccessAllowed,
UnsafeUrlError,
validateRemoteUrl,
} from './url-safety';
describe('url-safety', () => {
describe('isPrivateOrReservedIp', () => {
it.each([
'127.0.0.1',
'10.0.0.5',
'192.168.1.1',
'172.16.0.1',
'169.254.169.254',
'::1',
'fd00::1',
'fe80::1',
'febf::1',
'::ffff:127.0.0.1',
'::ffff:7f00:1',
'::ffff:c0a8:101',
])('flags %s as private/reserved', (ip) => {
expect(isPrivateOrReservedIp(ip)).toBe(true);
});
it.each(['93.184.216.34', '8.8.8.8', '1.1.1.1'])(
'allows public address %s',
(ip) => {
expect(isPrivateOrReservedIp(ip)).toBe(false);
}
);
});
describe('assertRemoteUrlAllowed', () => {
const publicResolver = async () => ['93.184.216.34'];
it('rejects non-http(s) protocols', async () => {
await expect(
assertRemoteUrlAllowed('file:///etc/passwd')
).rejects.toBeInstanceOf(UnsafeUrlError);
});
it('rejects embedded credentials', async () => {
await expect(
assertRemoteUrlAllowed('http://user:pass@example.com', {
resolveHostname: publicResolver,
})
).rejects.toBeInstanceOf(UnsafeUrlError);
});
it('rejects loopback literal', async () => {
await expect(
assertRemoteUrlAllowed('http://127.0.0.1/admin')
).rejects.toBeInstanceOf(UnsafeUrlError);
});
it('rejects localhost hostname', async () => {
await expect(
assertRemoteUrlAllowed('http://localhost:8080')
).rejects.toBeInstanceOf(UnsafeUrlError);
});
it('rejects the cloud metadata address', async () => {
await expect(
assertRemoteUrlAllowed(
'http://169.254.169.254/latest/meta-data'
)
).rejects.toBeInstanceOf(UnsafeUrlError);
});
it('rejects hex-form IPv4-mapped IPv6 loopback literals', async () => {
await expect(
assertRemoteUrlAllowed('http://[::ffff:7f00:1]/admin')
).rejects.toBeInstanceOf(UnsafeUrlError);
});
it('rejects a hostname that resolves to a private IP', async () => {
await expect(
assertRemoteUrlAllowed('http://rebind.example', {
resolveHostname: async () => ['10.0.0.1'],
})
).rejects.toBeInstanceOf(UnsafeUrlError);
});
it('allows a public URL', async () => {
const target = await validateRemoteUrl(
'https://example.com/playlist.m3u',
{ resolveHostname: publicResolver }
);
expect(target.url.hostname).toBe('example.com');
expect(target.addresses).toEqual(['93.184.216.34']);
});
it('honors the allowPrivateNetworks override', async () => {
const url = await assertRemoteUrlAllowed('http://127.0.0.1/probe', {
allowPrivateNetworks: true,
});
expect(url.hostname).toBe('127.0.0.1');
});
});
describe('isPrivateNetworkUrlAccessAllowed', () => {
const originalValue = process.env.IPTVNATOR_ALLOW_PRIVATE_NETWORK_URLS;
afterEach(() => {
if (originalValue === undefined) {
delete process.env.IPTVNATOR_ALLOW_PRIVATE_NETWORK_URLS;
} else {
process.env.IPTVNATOR_ALLOW_PRIVATE_NETWORK_URLS =
originalValue;
}
});
it('requires an explicit environment opt-in', () => {
delete process.env.IPTVNATOR_ALLOW_PRIVATE_NETWORK_URLS;
expect(isPrivateNetworkUrlAccessAllowed()).toBe(false);
process.env.IPTVNATOR_ALLOW_PRIVATE_NETWORK_URLS = ' TRUE ';
expect(isPrivateNetworkUrlAccessAllowed()).toBe(true);
});
});
});
@@ -0,0 +1,272 @@
import { lookup } from 'node:dns/promises';
import { isIP } from 'node:net';
const PRIVATE_NETWORK_URLS_ENV = 'IPTVNATOR_ALLOW_PRIVATE_NETWORK_URLS';
/**
* Raised when a renderer-supplied remote URL is rejected by the SSRF guard.
*/
export class UnsafeUrlError extends Error {
readonly status: number;
constructor(message: string, status = 400) {
super(message);
this.name = 'UnsafeUrlError';
this.status = status;
}
}
export interface RemoteUrlPolicy {
/** When true, private/reserved targets are permitted (default false). */
allowPrivateNetworks?: boolean;
/** Injectable DNS resolver. Defaults to `dns.lookup`; overridable in tests. */
resolveHostname?: (hostname: string) => Promise<readonly string[]>;
}
export interface ValidatedRemoteUrl {
/** Parsed URL, retaining the original hostname for TLS SNI and Host. */
url: URL;
/** Validated addresses that socket-level DNS lookup must be pinned to. */
addresses?: readonly string[];
}
function isTruthyEnvironmentOptIn(value: string | undefined): boolean {
const normalized = value?.trim().toLowerCase();
return normalized === '1' || normalized === 'true';
}
export function isPrivateNetworkUrlAccessAllowed(): boolean {
return isTruthyEnvironmentOptIn(process.env[PRIVATE_NETWORK_URLS_ENV]);
}
function normalizeHostname(hostname: string): string {
let host = hostname.trim().toLowerCase();
if (host.startsWith('[') && host.endsWith(']')) {
host = host.slice(1, -1);
}
return host;
}
export function isLocalHostname(hostname: string): boolean {
return hostname === 'localhost' || hostname.endsWith('.localhost');
}
export function isPrivateOrReservedIpv4(address: string): boolean {
const parts = address.split('.').map((part) => Number(part));
if (
parts.length !== 4 ||
parts.some((part) => !Number.isInteger(part) || part < 0 || part > 255)
) {
return true;
}
const [first, second, third] = parts;
return (
first === 0 ||
first === 10 ||
first === 127 ||
(first === 100 && second >= 64 && second <= 127) ||
(first === 169 && second === 254) ||
(first === 172 && second >= 16 && second <= 31) ||
(first === 192 && second === 168) ||
(first === 192 && second === 0) ||
(first === 198 && (second === 18 || second === 19)) ||
(first === 198 && second === 51 && third === 100) ||
(first === 203 && second === 0 && third === 113) ||
first >= 224
);
}
function parseIpv6Words(address: string): number[] | null {
const halves = address.split('::');
if (halves.length > 2) {
return null;
}
const parseHalf = (half: string): number[] | null => {
if (!half) {
return [];
}
const words: number[] = [];
for (const part of half.split(':')) {
if (part.includes('.')) {
const octets = part.split('.').map((octet) => Number(octet));
if (
octets.length !== 4 ||
octets.some(
(octet) =>
!Number.isInteger(octet) || octet < 0 || octet > 255
)
) {
return null;
}
words.push(
(octets[0] << 8) | octets[1],
(octets[2] << 8) | octets[3]
);
continue;
}
if (!/^[0-9a-f]{1,4}$/i.test(part)) {
return null;
}
words.push(Number.parseInt(part, 16));
}
return words;
};
const left = parseHalf(halves[0]);
const right = parseHalf(halves[1] ?? '');
if (!left || !right) {
return null;
}
if (halves.length === 1) {
return left.length === 8 ? left : null;
}
const omittedWordCount = 8 - left.length - right.length;
if (omittedWordCount < 1) {
return null;
}
return [...left, ...Array<number>(omittedWordCount).fill(0), ...right];
}
function getMappedIpv4Address(address: string): string | null {
const words = parseIpv6Words(address);
if (
!words ||
words.slice(0, 5).some((word) => word !== 0) ||
words[5] !== 0xffff
) {
return null;
}
return [
words[6] >> 8,
words[6] & 0xff,
words[7] >> 8,
words[7] & 0xff,
].join('.');
}
export function isPrivateOrReservedIpv6(address: string): boolean {
const normalized = address.toLowerCase();
if (
normalized === '::' ||
normalized === '::1' ||
normalized.startsWith('fc') ||
normalized.startsWith('fd') ||
// Link-local fe80::/10 spans fe80:: through febf:: (3rd nibble 8-b).
/^fe[89ab]/.test(normalized)
) {
return true;
}
const mappedIpv4 = getMappedIpv4Address(normalized);
return mappedIpv4 ? isPrivateOrReservedIpv4(mappedIpv4) : false;
}
export function isPrivateOrReservedIp(address: string): boolean {
const version = isIP(address);
if (version === 4) {
return isPrivateOrReservedIpv4(address);
}
if (version === 6) {
return isPrivateOrReservedIpv6(address);
}
return false;
}
async function defaultResolveHostname(
hostname: string
): Promise<readonly string[]> {
const records = await lookup(hostname, { all: true, verbatim: true });
return records.map((record) => record.address);
}
/**
* Validates a renderer-supplied remote URL before the Electron main process
* fetches it, preventing SSRF to loopback / private / reserved network
* targets (e.g. `http://127.0.0.1`, `http://169.254.169.254` metadata).
*
* Returns the parsed {@link URL} when allowed, otherwise throws
* {@link UnsafeUrlError}.
*/
export async function validateRemoteUrl(
rawUrl: string,
policy: RemoteUrlPolicy = {}
): Promise<ValidatedRemoteUrl> {
let url: URL;
try {
url = new URL(rawUrl);
} catch {
throw new UnsafeUrlError('Invalid URL');
}
if (url.protocol !== 'http:' && url.protocol !== 'https:') {
throw new UnsafeUrlError('Only http and https URLs are supported');
}
if (url.username || url.password) {
throw new UnsafeUrlError('URL credentials are not supported');
}
if (policy.allowPrivateNetworks) {
return { url };
}
const hostname = normalizeHostname(url.hostname);
if (isLocalHostname(hostname) || isPrivateOrReservedIp(hostname)) {
throw new UnsafeUrlError(
'URL points to a private or local network address'
);
}
if (isIP(hostname) !== 0) {
return { url, addresses: [hostname] };
}
{
const resolveHostname =
policy.resolveHostname ?? defaultResolveHostname;
let addresses: readonly string[];
try {
addresses = await resolveHostname(hostname);
} catch {
throw new UnsafeUrlError('URL host could not be resolved');
}
if (
addresses.length === 0 ||
addresses.some((address) => {
const normalizedAddress = normalizeHostname(address);
return (
isIP(normalizedAddress) === 0 ||
isPrivateOrReservedIp(normalizedAddress)
);
})
) {
throw new UnsafeUrlError(
'URL points to a private or local network address'
);
}
return {
url,
addresses: addresses.map((address) => normalizeHostname(address)),
};
}
}
/**
* Validates a URL and returns its parsed representation.
*/
export async function assertRemoteUrlAllowed(
rawUrl: string,
policy: RemoteUrlPolicy = {}
): Promise<URL> {
return (await validateRemoteUrl(rawUrl, policy)).url;
}
@@ -18,9 +18,11 @@ function createDeferred<T>() {
jest.mock('electron', () => ({
ipcMain: {
handle: jest.fn((channel: string, handler: (...args: unknown[]) => unknown) => {
registeredHandlers.set(channel, handler);
}),
handle: jest.fn(
(channel: string, handler: (...args: unknown[]) => unknown) => {
registeredHandlers.set(channel, handler);
}
),
},
}));
@@ -53,7 +55,10 @@ describe('XtreamEvents session cancellation', () => {
it('aborts requests that were registered with only a session id', async () => {
const requestHandler = registeredHandlers.get('XTREAM_REQUEST');
const cancelHandler = registeredHandlers.get(XTREAM_CANCEL_SESSION);
const pendingRequest = createDeferred<{ status: number; data: unknown }>();
const pendingRequest = createDeferred<{
status: number;
data: unknown;
}>();
const cancelError = Object.assign(new Error('cancelled'), {
code: 'ERR_CANCELED',
});
@@ -70,23 +75,27 @@ describe('XtreamEvents session cancellation', () => {
(error: unknown) => error === cancelError
);
const requestPromise = requestHandler?.({}, {
url: 'http://localhost:3211',
params: {
action: 'get_live_categories',
password: 'secret',
username: 'user1',
},
sessionId: 'session-1',
suppressErrorLog: true,
}) as Promise<unknown>;
const requestPromise = requestHandler?.(
{},
{
url: 'http://localhost:3211',
params: {
action: 'get_live_categories',
password: 'secret',
username: 'user1',
},
sessionId: 'session-1',
suppressErrorLog: true,
}
) as Promise<unknown>;
await Promise.resolve();
expect(abortSignal?.aborted).toBe(false);
const cancelResult = (await cancelHandler?.(
{},
'session-1'
)) as { success: boolean; cancelled: number };
const cancelResult = (await cancelHandler?.({}, 'session-1')) as {
success: boolean;
cancelled: number;
};
expect(cancelResult).toEqual({ success: true, cancelled: 1 });
expect(abortSignal?.aborted).toBe(true);
@@ -102,8 +111,14 @@ describe('XtreamEvents session cancellation', () => {
it('counts every matching in-flight request for the same session', async () => {
const requestHandler = registeredHandlers.get('XTREAM_REQUEST');
const cancelHandler = registeredHandlers.get(XTREAM_CANCEL_SESSION);
const firstRequest = createDeferred<{ status: number; data: unknown }>();
const secondRequest = createDeferred<{ status: number; data: unknown }>();
const firstRequest = createDeferred<{
status: number;
data: unknown;
}>();
const secondRequest = createDeferred<{
status: number;
data: unknown;
}>();
const cancelError = Object.assign(new Error('cancelled'), {
code: 'ERR_CANCELED',
});
@@ -129,31 +144,38 @@ describe('XtreamEvents session cancellation', () => {
(error: unknown) => error === cancelError
);
const firstPromise = requestHandler?.({}, {
url: 'http://localhost:3211',
params: {
action: 'get_live_categories',
password: 'secret',
username: 'user1',
},
sessionId: 'session-2',
suppressErrorLog: true,
}) as Promise<unknown>;
const secondPromise = requestHandler?.({}, {
url: 'http://localhost:3211',
params: {
action: 'get_vod_streams',
password: 'secret',
username: 'user1',
},
sessionId: 'session-2',
suppressErrorLog: true,
}) as Promise<unknown>;
const cancelResult = (await cancelHandler?.(
const firstPromise = requestHandler?.(
{},
'session-2'
)) as { success: boolean; cancelled: number };
{
url: 'http://localhost:3211',
params: {
action: 'get_live_categories',
password: 'secret',
username: 'user1',
},
sessionId: 'session-2',
suppressErrorLog: true,
}
) as Promise<unknown>;
const secondPromise = requestHandler?.(
{},
{
url: 'http://localhost:3211',
params: {
action: 'get_vod_streams',
password: 'secret',
username: 'user1',
},
sessionId: 'session-2',
suppressErrorLog: true,
}
) as Promise<unknown>;
await Promise.resolve();
const cancelResult = (await cancelHandler?.({}, 'session-2')) as {
success: boolean;
cancelled: number;
};
expect(cancelResult).toEqual({ success: true, cancelled: 2 });
expect(abortSignals).toHaveLength(2);
@@ -162,7 +184,11 @@ describe('XtreamEvents session cancellation', () => {
firstRequest.reject(cancelError);
secondRequest.reject(cancelError);
await expect(firstPromise).rejects.toMatchObject({ name: 'AbortError' });
await expect(secondPromise).rejects.toMatchObject({ name: 'AbortError' });
await expect(firstPromise).rejects.toMatchObject({
name: 'AbortError',
});
await expect(secondPromise).rejects.toMatchObject({
name: 'AbortError',
});
});
});
@@ -5,8 +5,13 @@
import axios, { AxiosRequestConfig } from 'axios';
import { ipcMain } from 'electron';
import { PortalDebugEvent, XTREAM_CANCEL_SESSION } from '@iptvnator/shared/interfaces';
import {
PortalDebugEvent,
XTREAM_CANCEL_SESSION,
} from '@iptvnator/shared/interfaces';
import { emitPortalDebugEvent } from './portal-debug.events';
import { assertRemoteUrlAllowed, UnsafeUrlError } from './url-safety';
import { requestWithValidatedRedirects } from '../util/validated-axios';
export default class XtreamEvents {
static bootstrapXtreamEvents(): Electron.IpcMain {
@@ -14,7 +19,11 @@ export default class XtreamEvents {
}
}
function formatXtreamError(error: unknown, requestUrl: string, action?: string) {
function formatXtreamError(
error: unknown,
requestUrl: string,
action?: string
) {
const parsedUrl = new URL(requestUrl);
const base = {
action,
@@ -101,7 +110,11 @@ ipcMain.handle(
signal: controller.signal,
};
const response = await axios(config);
const response = await requestWithValidatedRedirects<unknown>(
apiUrl.toString(),
config,
{ allowPrivateNetworks: true }
);
// Check if response is successful
if (response.status >= 400) {
@@ -172,7 +185,11 @@ ipcMain.handle(
if (!payload.suppressErrorLog) {
console.error(
'[XTREAM_REQUEST] Failed',
formatXtreamError(error, payload.url, payload.params?.action)
formatXtreamError(
error,
payload.url,
payload.params?.action
)
);
}
@@ -218,7 +235,10 @@ ipcMain.handle(
ipcMain.handle(
XTREAM_CANCEL_SESSION,
async (_event, sessionId: string): Promise<{ success: boolean; cancelled: number }> => {
async (
_event,
sessionId: string
): Promise<{ success: boolean; cancelled: number }> => {
if (!sessionId) {
return { success: false, cancelled: 0 };
}
@@ -255,6 +275,26 @@ ipcMain.handle(
method?: 'GET' | 'HEAD';
}
) => {
// Guard against SSRF: a malicious portal/playlist could ask the main
// process to probe loopback/private/metadata addresses. Only allow
// public http(s) targets.
try {
// Private/LAN targets are allowed (users probe self-hosted Xtream
// servers); maxRedirects:0 below blocks redirect-based SSRF.
await assertRemoteUrlAllowed(payload.url, {
allowPrivateNetworks: true,
});
} catch (error) {
return {
status: 0,
url: payload.url,
error:
error instanceof UnsafeUrlError
? error.message
: 'Invalid URL',
};
}
const config: AxiosRequestConfig = {
method: payload.method ?? 'HEAD',
url: payload.url,
@@ -263,7 +303,11 @@ ipcMain.handle(
'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36',
},
timeout: 10000,
maxRedirects: 5,
// SSRF hardening: do NOT follow redirects. assertRemoteUrlAllowed
// validated payload.url, but a validated public host could 3xx to an
// internal address that would never be re-checked. validateStatus
// below returns the 3xx as the probe result instead of following it.
maxRedirects: 0,
validateStatus: () => true,
};
@@ -0,0 +1,71 @@
import { Agent } from 'node:https';
import type { LookupFunction } from 'node:net';
import {
createPlaylistAgentFactory,
isInsecureTlsAllowed,
} from './secure-https';
jest.mock('node:https', () => ({
Agent: jest.fn(() => ({ protocol: 'https:' })),
}));
describe('secure-https', () => {
const originalValue = process.env.IPTVNATOR_ALLOW_INSECURE_TLS;
const agentConstructorMock = Agent as unknown as jest.Mock;
const lookup = jest.fn() as unknown as LookupFunction;
beforeEach(() => {
agentConstructorMock.mockClear();
});
afterEach(() => {
if (originalValue === undefined) {
delete process.env.IPTVNATOR_ALLOW_INSECURE_TLS;
} else {
process.env.IPTVNATOR_ALLOW_INSECURE_TLS = originalValue;
}
});
it('validates certificates by default', () => {
delete process.env.IPTVNATOR_ALLOW_INSECURE_TLS;
expect(isInsecureTlsAllowed()).toBe(false);
createPlaylistAgentFactory().createHttpsAgent(lookup);
expect(agentConstructorMock).toHaveBeenCalledWith({
lookup,
rejectUnauthorized: true,
});
});
it.each(['1', 'true', ' TRUE '])(
'allows an explicit insecure TLS opt-in via %s',
(value) => {
process.env.IPTVNATOR_ALLOW_INSECURE_TLS = value;
expect(isInsecureTlsAllowed()).toBe(true);
createPlaylistAgentFactory().createHttpsAgent(lookup);
expect(agentConstructorMock).toHaveBeenCalledWith({
lookup,
rejectUnauthorized: false,
});
}
);
it('does not accept unrelated truthy values', () => {
process.env.IPTVNATOR_ALLOW_INSECURE_TLS = 'yes';
expect(isInsecureTlsAllowed()).toBe(false);
});
it('preserves TLS policy when no pinned lookup is required', () => {
delete process.env.IPTVNATOR_ALLOW_INSECURE_TLS;
createPlaylistAgentFactory().createHttpsAgent();
expect(agentConstructorMock).toHaveBeenCalledWith({
rejectUnauthorized: true,
});
});
});
@@ -0,0 +1,36 @@
import { Agent } from 'node:https';
import type { LookupFunction } from 'node:net';
import type { ValidatedRequestAgentFactory } from './validated-axios';
const INSECURE_TLS_ENV = 'IPTVNATOR_ALLOW_INSECURE_TLS';
/**
* Whether TLS certificate validation is disabled for remote playlist fetches.
*
* Secure by default. Validation is only disabled when the operator explicitly
* opts in via `IPTVNATOR_ALLOW_INSECURE_TLS=1` (or `=true`). Some IPTV
* providers serve self-signed or otherwise invalid certificates, so the
* opt-out exists — but it must never be the default, otherwise every playlist
* fetch is silently exposed to network MITM.
*/
export function isInsecureTlsAllowed(): boolean {
const value = process.env[INSECURE_TLS_ENV]?.trim().toLowerCase();
return value === '1' || value === 'true';
}
/**
* Builds the agent factory used for remote playlist fetches.
*
* Certificates are validated unless insecure TLS has been explicitly opted in
* (see {@link isInsecureTlsAllowed}). The validated request layer supplies a
* DNS lookup pinned to the addresses approved for each redirect hop.
*/
export function createPlaylistAgentFactory(): ValidatedRequestAgentFactory & {
createHttpsAgent(lookup?: LookupFunction): Agent;
} {
const rejectUnauthorized = !isInsecureTlsAllowed();
return {
createHttpsAgent: (lookup) => new Agent({ lookup, rejectUnauthorized }),
};
}
@@ -0,0 +1,384 @@
import axios from 'axios';
import type { LookupAddress, LookupOptions } from 'node:dns';
import { Agent as HttpAgent } from 'node:http';
import { Agent as HttpsAgent } from 'node:https';
import type { LookupFunction } from 'node:net';
import { UnsafeUrlError } from '../events/url-safety';
import {
requestWithValidatedRedirects,
type ValidatedRequestAgentFactory,
} from './validated-axios';
jest.mock('axios', () => ({
__esModule: true,
default: jest.fn(),
}));
const axiosMock = axios as unknown as jest.Mock;
describe('requestWithValidatedRedirects', () => {
const publicResolver = async () => ['93.184.216.34'];
function createCapturingAgentFactory() {
const lookups: LookupFunction[] = [];
const factory: ValidatedRequestAgentFactory = {
createHttpsAgent: jest.fn((lookup) => {
if (!lookup) {
throw new Error('Expected a pinned lookup');
}
lookups.push(lookup);
return new HttpsAgent({ lookup });
}),
};
return { factory, lookups };
}
beforeEach(() => {
axiosMock.mockReset();
});
it('rejects a redirect target that violates the URL policy', async () => {
axiosMock.mockResolvedValueOnce({
status: 302,
headers: {
location: 'http://169.254.169.254/latest/meta-data',
},
});
await expect(
requestWithValidatedRedirects(
'https://example.com/player_api.php',
{ method: 'GET' },
{ resolveHostname: publicResolver }
)
).rejects.toBeInstanceOf(UnsafeUrlError);
expect(axiosMock).toHaveBeenCalledTimes(1);
});
it('pins the socket lookup to the address validated by the URL policy', async () => {
axiosMock.mockResolvedValueOnce({
status: 200,
headers: {},
data: '#EXTM3U',
});
const { factory, lookups } = createCapturingAgentFactory();
await requestWithValidatedRedirects(
'https://epg.example/guide.xml',
{ agentFactory: factory, method: 'GET' },
{
resolveHostname: async () => ['93.184.216.34'],
}
);
const requestConfig = axiosMock.mock.calls[0][0];
const resolvedAddress = await new Promise<{
address: string | LookupAddress[];
family?: number;
}>((resolve, reject) => {
lookups[0](
'epg.example',
{ all: false, family: 0 } as LookupOptions,
(error, address, family) => {
if (error) {
reject(error);
return;
}
resolve({ address, family });
}
);
});
expect(factory.createHttpsAgent).toHaveBeenCalledWith(lookups[0]);
expect(requestConfig.httpsAgent).toBeInstanceOf(HttpsAgent);
expect(requestConfig).not.toHaveProperty('agentFactory');
expect(resolvedAddress).toEqual({
address: '93.184.216.34',
family: 4,
});
expect(requestConfig.proxy).toBe(false);
expect(requestConfig.url).toBe('https://epg.example/guide.xml');
});
it('supplies the validated lookup to an explicit HTTP agent factory', async () => {
axiosMock.mockResolvedValueOnce({
status: 200,
headers: {},
data: '#EXTM3U',
});
let pinnedLookup: LookupFunction | undefined;
const factory: ValidatedRequestAgentFactory = {
createHttpAgent: jest.fn((lookup) => {
pinnedLookup = lookup;
return new HttpAgent({ lookup });
}),
};
await requestWithValidatedRedirects(
'http://playlist.example/list.m3u',
{ agentFactory: factory, method: 'GET' },
{ resolveHostname: publicResolver }
);
const requestConfig = axiosMock.mock.calls[0][0];
expect(factory.createHttpAgent).toHaveBeenCalledWith(pinnedLookup);
expect(requestConfig.httpAgent).toBeInstanceOf(HttpAgent);
expect(requestConfig).not.toHaveProperty('agentFactory');
});
it('uses a custom HTTPS agent factory when private networks are allowed', async () => {
axiosMock.mockResolvedValueOnce({
status: 200,
headers: {},
data: '#EXTM3U',
});
const factory: ValidatedRequestAgentFactory = {
createHttpsAgent: jest.fn((lookup) => new HttpsAgent({ lookup })),
};
await requestWithValidatedRedirects(
'https://192.168.1.10/list.m3u',
{ agentFactory: factory, method: 'GET' },
{ allowPrivateNetworks: true }
);
const requestConfig = axiosMock.mock.calls[0][0];
expect(factory.createHttpsAgent).toHaveBeenCalledWith();
expect(requestConfig.httpsAgent).toBeInstanceOf(HttpsAgent);
expect(requestConfig).not.toHaveProperty('agentFactory');
});
it('revalidates and repins every redirect hop', async () => {
axiosMock
.mockResolvedValueOnce({
status: 302,
headers: { location: 'https://cdn.example/guide.xml' },
})
.mockResolvedValueOnce({
status: 200,
headers: {},
data: '<tv />',
});
const { factory, lookups } = createCapturingAgentFactory();
const resolveHostname = jest.fn(async (hostname: string) =>
hostname === 'epg.example' ? ['93.184.216.34'] : ['142.250.191.110']
);
await requestWithValidatedRedirects(
'https://epg.example/guide.xml',
{ agentFactory: factory, method: 'GET' },
{ resolveHostname }
);
const resolvePinnedAddress = async (callIndex: number) => {
return new Promise<string | LookupAddress[]>((resolve, reject) => {
lookups[callIndex](
'ignored.example',
{ all: false, family: 0 } as LookupOptions,
(error, address) => {
if (error) {
reject(error);
return;
}
resolve(address);
}
);
});
};
expect(resolveHostname).toHaveBeenNthCalledWith(1, 'epg.example');
expect(resolveHostname).toHaveBeenNthCalledWith(2, 'cdn.example');
expect(factory.createHttpsAgent).toHaveBeenCalledTimes(2);
await expect(resolvePinnedAddress(0)).resolves.toBe('93.184.216.34');
await expect(resolvePinnedAddress(1)).resolves.toBe('142.250.191.110');
});
it('removes sensitive headers when a redirect changes origin', async () => {
axiosMock
.mockResolvedValueOnce({
status: 302,
headers: { location: 'https://cdn.example.net/playlist.m3u' },
})
.mockResolvedValueOnce({
status: 200,
headers: {},
data: '#EXTM3U',
});
await requestWithValidatedRedirects(
'https://example.com/playlist.m3u',
{
headers: {
Authorization: 'Bearer secret',
Cookie: 'session=secret',
Accept: 'text/plain',
},
method: 'GET',
},
{ resolveHostname: publicResolver }
);
const redirectedConfig = axiosMock.mock.calls[1][0];
expect(redirectedConfig.headers).toMatchObject({
Accept: 'text/plain',
});
expect(redirectedConfig.headers).not.toHaveProperty('Authorization');
expect(redirectedConfig.headers).not.toHaveProperty('Cookie');
});
it('does not forward axios params to a cross-origin redirect', async () => {
axiosMock
.mockResolvedValueOnce({
status: 302,
headers: { location: 'https://cdn.example.net/playlist.m3u' },
})
.mockResolvedValueOnce({
status: 200,
headers: {},
data: '#EXTM3U',
});
await requestWithValidatedRedirects(
'https://example.com/playlist.m3u',
{
method: 'GET',
params: { token: 'secret' },
},
{ resolveHostname: publicResolver }
);
expect(axiosMock.mock.calls[1][0].params).toBeUndefined();
});
it('does not forward axios basic auth to a cross-origin redirect', async () => {
axiosMock
.mockResolvedValueOnce({
status: 302,
headers: { location: 'https://cdn.example.net/playlist.m3u' },
})
.mockResolvedValueOnce({
status: 200,
headers: {},
data: '#EXTM3U',
});
await requestWithValidatedRedirects(
'https://example.com/playlist.m3u',
{
auth: { password: 'secret', username: 'provider' },
method: 'GET',
},
{ resolveHostname: publicResolver }
);
expect(axiosMock.mock.calls[1][0].auth).toBeUndefined();
});
it('rejects a cross-origin redirect that would replay a request body', async () => {
axiosMock.mockResolvedValueOnce({
status: 307,
headers: { location: 'https://other.example/submit' },
});
await expect(
requestWithValidatedRedirects(
'https://example.com/submit',
{ data: { secret: true }, method: 'POST' },
{ resolveHostname: publicResolver }
)
).rejects.toThrow(/request bodies/i);
expect(axiosMock).toHaveBeenCalledTimes(1);
});
it('rejects redirects without a location header', async () => {
axiosMock.mockResolvedValueOnce({
status: 302,
headers: {},
});
await expect(
requestWithValidatedRedirects(
'https://example.com/start',
{ method: 'GET' },
{ resolveHostname: publicResolver }
)
).rejects.toMatchObject({
message: expect.stringMatching(/location/i),
status: 502,
});
});
it('stops after the configured redirect limit', async () => {
axiosMock.mockResolvedValue({
status: 302,
headers: { location: '/again' },
});
await expect(
requestWithValidatedRedirects(
'https://example.com/start',
{ method: 'GET' },
{ resolveHostname: publicResolver },
1
)
).rejects.toMatchObject({
message: expect.stringMatching(/too many redirects/i),
status: 502,
});
expect(axiosMock).toHaveBeenCalledTimes(2);
});
it('converts a 303 redirect to GET and removes the request body', async () => {
axiosMock
.mockResolvedValueOnce({
status: 303,
headers: { location: '/result' },
})
.mockResolvedValueOnce({
status: 200,
headers: {},
data: 'ok',
});
await requestWithValidatedRedirects(
'https://example.com/submit',
{ data: { value: 1 }, method: 'POST' },
{ resolveHostname: publicResolver }
);
expect(axiosMock.mock.calls[1][0]).toMatchObject({
data: undefined,
method: 'GET',
url: 'https://example.com/result',
});
});
it.each([301, 302])(
'converts a POST redirected with %i to GET and removes the request body',
async (status) => {
axiosMock
.mockResolvedValueOnce({
status,
headers: { location: '/result' },
})
.mockResolvedValueOnce({
status: 200,
headers: {},
data: 'ok',
});
await requestWithValidatedRedirects(
'https://example.com/submit',
{ data: { value: 1 }, method: 'POST' },
{ resolveHostname: publicResolver }
);
expect(axiosMock.mock.calls[1][0]).toMatchObject({
data: undefined,
method: 'GET',
url: 'https://example.com/result',
});
}
);
});
@@ -0,0 +1,215 @@
import axios, {
AxiosRequestConfig,
AxiosResponse,
RawAxiosRequestHeaders,
} from 'axios';
import type { LookupAddress } from 'node:dns';
import { Agent as HttpAgent } from 'node:http';
import { Agent as HttpsAgent } from 'node:https';
import { isIP, LookupFunction } from 'node:net';
import {
RemoteUrlPolicy,
UnsafeUrlError,
validateRemoteUrl,
} from '../events/url-safety';
const REDIRECT_STATUSES = new Set([301, 302, 303, 307, 308]);
const SENSITIVE_HEADERS = new Set([
'authorization',
'cookie',
'proxy-authorization',
]);
export interface ValidatedRequestAgentFactory {
createHttpAgent?(lookup?: LookupFunction): HttpAgent;
createHttpsAgent?(lookup?: LookupFunction): HttpsAgent;
}
export type ValidatedAxiosRequestConfig = Omit<
AxiosRequestConfig,
'httpAgent' | 'httpsAgent'
> & {
agentFactory?: ValidatedRequestAgentFactory;
};
function copyHeadersWithoutSensitiveValues(
headers: AxiosRequestConfig['headers']
): AxiosRequestConfig['headers'] {
if (!headers) {
return headers;
}
const source =
typeof (headers as { toJSON?: () => RawAxiosRequestHeaders }).toJSON ===
'function'
? (
headers as {
toJSON: () => RawAxiosRequestHeaders;
}
).toJSON()
: headers;
const sanitized: RawAxiosRequestHeaders = {};
for (const [name, value] of Object.entries(source)) {
if (!SENSITIVE_HEADERS.has(name.toLowerCase())) {
sanitized[name] = value;
}
}
return sanitized;
}
function createPinnedLookup(addresses: readonly string[]): LookupFunction {
const records: LookupAddress[] = addresses.map((address) => ({
address,
family: isIP(address),
}));
return (_hostname, options, callback) => {
const requestedFamily = options.family;
const eligibleRecords = requestedFamily
? records.filter((record) => record.family === requestedFamily)
: records;
if (eligibleRecords.length === 0) {
const error = new Error(
'Validated URL has no connectable address'
) as NodeJS.ErrnoException;
error.code = 'ENOTFOUND';
callback(error, []);
return;
}
if (options.all) {
callback(null, eligibleRecords);
return;
}
const selected = eligibleRecords[0];
callback(null, selected.address, selected.family);
};
}
function pinRequestToValidatedAddresses(
config: ValidatedAxiosRequestConfig,
url: URL,
addresses: readonly string[] | undefined
): AxiosRequestConfig {
const { agentFactory, ...axiosConfig } = config;
if (!addresses) {
if (url.protocol === 'https:' && agentFactory?.createHttpsAgent) {
return {
...axiosConfig,
httpsAgent: agentFactory.createHttpsAgent(),
};
}
if (url.protocol === 'http:' && agentFactory?.createHttpAgent) {
return {
...axiosConfig,
httpAgent: agentFactory.createHttpAgent(),
};
}
return axiosConfig;
}
const lookup = createPinnedLookup(addresses);
if (url.protocol === 'https:') {
return {
...axiosConfig,
httpsAgent:
agentFactory?.createHttpsAgent?.(lookup) ??
new HttpsAgent({ lookup }),
proxy: false,
};
}
return {
...axiosConfig,
httpAgent:
agentFactory?.createHttpAgent?.(lookup) ??
new HttpAgent({ lookup }),
proxy: false,
};
}
/**
* Runs an Axios request while validating the initial URL and every redirect.
* Redirects are followed manually so each target passes through the same
* private-network and protocol policy.
*/
export async function requestWithValidatedRedirects<T = unknown>(
rawUrl: string,
config: ValidatedAxiosRequestConfig = {},
policy: RemoteUrlPolicy = {},
maxRedirects = 5
): Promise<AxiosResponse<T>> {
const originalValidateStatus =
config.validateStatus ??
((status: number) => status >= 200 && status < 300);
let currentUrl = rawUrl;
let requestConfig = { ...config };
for (let redirectCount = 0; ; redirectCount += 1) {
const validatedTarget = await validateRemoteUrl(currentUrl, policy);
const validatedUrl = validatedTarget.url;
const pinnedConfig = pinRequestToValidatedAddresses(
requestConfig,
validatedUrl,
validatedTarget.addresses
);
const response = await axios<T>({
...pinnedConfig,
maxRedirects: 0,
url: validatedUrl.toString(),
validateStatus: (status) =>
REDIRECT_STATUSES.has(status) || originalValidateStatus(status),
});
if (!REDIRECT_STATUSES.has(response.status)) {
return response;
}
const location = response.headers?.location;
if (!location) {
throw new UnsafeUrlError(
'Redirect response did not include a location',
502
);
}
if (redirectCount >= maxRedirects) {
throw new UnsafeUrlError('Too many redirects', 502);
}
const nextUrl = new URL(location, validatedUrl);
const method = requestConfig.method?.toUpperCase();
const shouldRewriteToGet =
(response.status === 303 && method !== 'HEAD') ||
((response.status === 301 || response.status === 302) &&
method === 'POST');
if (shouldRewriteToGet) {
requestConfig = {
...requestConfig,
data: undefined,
method: 'GET',
};
}
if (nextUrl.origin !== validatedUrl.origin) {
if (requestConfig.data !== undefined) {
throw new UnsafeUrlError(
'Cross-origin redirects with request bodies are not supported',
502
);
}
requestConfig = {
...requestConfig,
auth: undefined,
headers: copyHeadersWithoutSensitiveValues(
requestConfig.headers
),
params: undefined,
};
}
currentUrl = nextUrl.toString();
}
}
@@ -0,0 +1,48 @@
import type BetterSqlite3 from 'better-sqlite3';
import { EpgDatabaseClearOperation } from './epg-database';
function createDatabaseMock(exec: jest.Mock) {
const database = {
close: jest.fn(),
exec,
};
const Database = jest.fn(() => database) as unknown as typeof BetterSqlite3;
return { Database, database };
}
describe('EpgDatabaseClearOperation', () => {
it('clears programs and channels in one transaction', () => {
const exec = jest.fn();
const { Database } = createDatabaseMock(exec);
new EpgDatabaseClearOperation(Database).run();
expect(exec.mock.calls.map(([statement]) => statement)).toEqual([
'BEGIN',
'DELETE FROM epg_programs',
'DELETE FROM epg_channels',
'COMMIT',
]);
});
it('rolls back when either delete fails', () => {
const failure = new Error('database is busy');
const exec = jest.fn((statement: string) => {
if (statement === 'DELETE FROM epg_channels') {
throw failure;
}
});
const { Database } = createDatabaseMock(exec);
expect(() => new EpgDatabaseClearOperation(Database).run()).toThrow(
failure
);
expect(exec.mock.calls.map(([statement]) => statement)).toEqual([
'BEGIN',
'DELETE FROM epg_programs',
'DELETE FROM epg_channels',
'ROLLBACK',
]);
});
});
@@ -0,0 +1,146 @@
import type BetterSqlite3 from 'better-sqlite3';
import { getIptvnatorDatabasePath } from '@iptvnator/shared/database/path-utils';
import type { ParsedChannel, ParsedProgram } from './epg-streaming-parser';
/**
* Database helper for worker-owned EPG operations.
* Creates its own connection to avoid blocking the main thread.
*/
export class EpgDatabase {
private readonly db: BetterSqlite3.Database;
private readonly knownChannelIds = new Set<string>();
private readonly insertChannelStmt: BetterSqlite3.Statement;
private readonly insertProgramStmt: BetterSqlite3.Statement;
private readonly deleteChannelsStmt: BetterSqlite3.Statement;
constructor(Database: typeof BetterSqlite3) {
this.db = new Database(getIptvnatorDatabasePath());
this.db.pragma('foreign_keys = ON');
this.db.pragma('journal_mode = WAL');
this.db.pragma('busy_timeout = 5000');
this.insertChannelStmt = this.db.prepare(`
INSERT INTO epg_channels (id, display_name, icon_url, url, source_url, updated_at)
VALUES (?, ?, ?, ?, ?, strftime('%Y-%m-%dT%H:%M:%SZ', 'now'))
ON CONFLICT(id) DO UPDATE SET
display_name = excluded.display_name,
icon_url = excluded.icon_url,
url = excluded.url,
source_url = excluded.source_url,
updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
`);
this.insertProgramStmt = this.db.prepare(`
INSERT INTO epg_programs (channel_id, start, stop, title, description, category, icon_url, rating, episode_num)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
`);
this.deleteChannelsStmt = this.db.prepare(`
DELETE FROM epg_channels WHERE source_url = ?
`);
}
/**
* Insert a batch of channels. When `clearFirst` is true, the existing rows
* for `sourceUrl` are deleted inside the same transaction as the insert so
* old data is preserved if the fetch/parse never produces any channels.
*/
insertChannels(
channels: ParsedChannel[],
sourceUrl: string,
clearFirst = false
): void {
const insertMany = this.db.transaction((channels: ParsedChannel[]) => {
if (clearFirst) {
this.deleteChannelsStmt.run(sourceUrl);
this.knownChannelIds.clear();
}
for (const channel of channels) {
const displayName =
channel.displayName?.[0]?.value || channel.id;
const iconUrl = channel.icon?.[0]?.src || null;
const url = channel.url?.[0] || null;
this.insertChannelStmt.run(
channel.id,
displayName,
iconUrl,
url,
sourceUrl
);
this.knownChannelIds.add(channel.id);
}
});
insertMany(channels);
}
/**
* Insert programs for channels already seen during the current parse.
*/
insertPrograms(programs: ParsedProgram[]): number {
let insertedCount = 0;
const insertMany = this.db.transaction((programs: ParsedProgram[]) => {
for (const program of programs) {
if (!this.knownChannelIds.has(program.channel)) continue;
const title = program.title?.[0]?.value || 'Unknown';
const description = program.desc?.[0]?.value || null;
const category = program.category?.[0]?.value || null;
const iconUrl = program.icon?.[0]?.src || null;
const rating = program.rating?.[0]?.value || null;
const episodeNum = program.episodeNum?.[0]?.value || null;
try {
this.insertProgramStmt.run(
program.channel,
program.start,
program.stop,
title,
description,
category,
iconUrl,
rating,
episodeNum
);
insertedCount++;
} catch {
// Skip individual failures such as FK constraint errors.
}
}
});
insertMany(programs);
return insertedCount;
}
close(): void {
this.db.close();
}
}
export class EpgDatabaseClearOperation {
private readonly db: BetterSqlite3.Database;
constructor(Database: typeof BetterSqlite3) {
this.db = new Database(getIptvnatorDatabasePath());
}
run(): void {
this.db.exec('BEGIN');
try {
this.db.exec('DELETE FROM epg_programs');
this.db.exec('DELETE FROM epg_channels');
this.db.exec('COMMIT');
} catch (error) {
this.db.exec('ROLLBACK');
throw error;
}
}
close(): void {
this.db.close();
}
}
@@ -1,15 +1,12 @@
import type BetterSqlite3 from 'better-sqlite3';
import { existsSync, mkdirSync } from 'fs';
import { getIptvnatorDatabasePath } from '@iptvnator/shared/database/path-utils';
import { Readable } from 'stream';
import { parentPort, workerData } from 'worker_threads';
import { createGunzip } from 'zlib';
import {
ParsedChannel,
ParsedProgram,
StreamingEpgParser,
} from './epg-streaming-parser';
import { EpgDatabase, EpgDatabaseClearOperation } from './epg-database';
import { StreamingEpgParser } from './epg-streaming-parser';
import { shouldGunzipEpgResponse } from './epg-response-utils';
import { isPrivateNetworkUrlAccessAllowed } from '../events/url-safety';
import { requestWithValidatedRedirects } from '../util/validated-axios';
import {
getNativeModuleSearchPaths,
getWorkerDataNativeModuleSearchPaths,
@@ -17,14 +14,11 @@ import {
registerNativeModuleSearchPaths,
} from './worker-runtime-paths';
let Database: typeof BetterSqlite3;
const nativeModuleSearchPaths = [
...getWorkerDataNativeModuleSearchPaths(workerData),
...getNativeModuleSearchPaths({
resourcesPath: (
process as NodeJS.Process & { resourcesPath?: string }
).resourcesPath,
resourcesPath: (process as NodeJS.Process & { resourcesPath?: string })
.resourcesPath,
}),
];
@@ -36,12 +30,11 @@ function loadBetterSqlite3(): typeof BetterSqlite3 {
loggerLabel: '[EPG Worker]',
searchPaths: nativeModuleSearchPaths,
fallbackRequire: () =>
// eslint-disable-next-line @typescript-eslint/no-require-imports
require('better-sqlite3') as typeof BetterSqlite3,
});
}
Database = loadBetterSqlite3();
const Database = loadBetterSqlite3();
/**
* Streaming EPG Parser Worker
@@ -76,143 +69,6 @@ const loggerLabel = '[EPG Worker]';
const CHANNEL_BATCH_SIZE = 100;
const PROGRAM_BATCH_SIZE = 1000;
/**
* Database helper class for EPG operations
* Creates its own connection to avoid blocking main thread
*/
class EpgDatabase {
private db: BetterSqlite3.Database;
private knownChannelIds: Set<string> = new Set();
// Prepared statements for better performance
private insertChannelStmt: BetterSqlite3.Statement;
private insertProgramStmt: BetterSqlite3.Statement;
private deleteChannelsStmt: BetterSqlite3.Statement;
constructor() {
const dbPath = getIptvnatorDatabasePath();
this.db = new Database(dbPath);
this.db.pragma('foreign_keys = ON');
this.db.pragma('journal_mode = WAL'); // Better concurrent write performance
this.db.pragma('busy_timeout = 5000');
// Prepare statements
// Use INSERT OR REPLACE to update existing channels and refresh updated_at
// Use strftime with ISO format for consistent date comparison
this.insertChannelStmt = this.db.prepare(`
INSERT INTO epg_channels (id, display_name, icon_url, url, source_url, updated_at)
VALUES (?, ?, ?, ?, ?, strftime('%Y-%m-%dT%H:%M:%SZ', 'now'))
ON CONFLICT(id) DO UPDATE SET
display_name = excluded.display_name,
icon_url = excluded.icon_url,
url = excluded.url,
source_url = excluded.source_url,
updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
`);
this.insertProgramStmt = this.db.prepare(`
INSERT INTO epg_programs (channel_id, start, stop, title, description, category, icon_url, rating, episode_num)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
`);
this.deleteChannelsStmt = this.db.prepare(`
DELETE FROM epg_channels WHERE source_url = ?
`);
}
/**
* Clear existing EPG data for a source URL
*/
clearSourceData(sourceUrl: string): void {
this.deleteChannelsStmt.run(sourceUrl);
this.knownChannelIds.clear();
}
/**
* Insert a batch of channels. When `clearFirst` is true, the existing rows
* for `sourceUrl` are deleted inside the same transaction as the insert so
* old data is preserved if the fetch/parse never produces any channels.
*/
insertChannels(
channels: ParsedChannel[],
sourceUrl: string,
clearFirst = false
): void {
const insertMany = this.db.transaction((channels: ParsedChannel[]) => {
if (clearFirst) {
this.deleteChannelsStmt.run(sourceUrl);
this.knownChannelIds.clear();
}
for (const channel of channels) {
const displayName =
channel.displayName?.[0]?.value || channel.id;
const iconUrl = channel.icon?.[0]?.src || null;
const url = channel.url?.[0] || null;
this.insertChannelStmt.run(
channel.id,
displayName,
iconUrl,
url,
sourceUrl
);
this.knownChannelIds.add(channel.id);
}
});
insertMany(channels);
}
/**
* Insert a batch of programs
* Only inserts programs for known channels to avoid FK constraint failures
*/
insertPrograms(programs: ParsedProgram[]): number {
let insertedCount = 0;
const insertMany = this.db.transaction((programs: ParsedProgram[]) => {
for (const prog of programs) {
// Skip if channel not known
if (!this.knownChannelIds.has(prog.channel)) continue;
const title = prog.title?.[0]?.value || 'Unknown';
const description = prog.desc?.[0]?.value || null;
const category = prog.category?.[0]?.value || null;
const iconUrl = prog.icon?.[0]?.src || null;
const rating = prog.rating?.[0]?.value || null;
const episodeNum = prog.episodeNum?.[0]?.value || null;
try {
this.insertProgramStmt.run(
prog.channel,
prog.start,
prog.stop,
title,
description,
category,
iconUrl,
rating,
episodeNum
);
insertedCount++;
} catch (err) {
// Skip individual failures (e.g., FK constraint)
}
}
});
insertMany(programs);
return insertedCount;
}
/**
* Close the database connection
*/
close(): void {
this.db.close();
}
}
/**
* Fetches and parses EPG data from URL using streaming
* Inserts directly into SQLite to avoid blocking main thread
@@ -221,7 +77,7 @@ async function fetchAndParseEpgStreaming(url: string): Promise<void> {
console.log(loggerLabel, `Fetching EPG from ${url}`);
// Create database connection in worker
const epgDb = new EpgDatabase();
const epgDb = new EpgDatabase(Database);
// Old rows for this source are retained until the first successful insert
// batch arrives — see `hasClearedSource` below. That way, a fetch or parse
@@ -230,13 +86,30 @@ async function fetchAndParseEpgStreaming(url: string): Promise<void> {
let hasClearedSource = false;
try {
const response = await fetch(url.trim());
const isGzipped = shouldGunzipEpgResponse(url, response);
// EPG URLs can originate from an untrusted M3U `url-tvg` attribute.
// Validate every redirect and require an explicit operator opt-in for
// private/LAN sources.
const response = await requestWithValidatedRedirects<Readable>(
url.trim(),
{
decompress: false,
method: 'GET',
responseType: 'stream',
},
{
allowPrivateNetworks: isPrivateNetworkUrlAccessAllowed(),
}
);
const responseUrl = response.config.url;
const isGzipped = shouldGunzipEpgResponse(url, {
headers: response.headers,
url: responseUrl,
});
if (response.url && response.url !== url) {
if (responseUrl && responseUrl !== url) {
console.log(
loggerLabel,
`Resolved EPG redirect: ${url} -> ${response.url}`
`Resolved EPG redirect: ${url} -> ${responseUrl}`
);
}
@@ -245,11 +118,11 @@ async function fetchAndParseEpgStreaming(url: string): Promise<void> {
`EPG response detected as gzipped: ${isGzipped}`
);
if (!response.ok) {
if (response.status < 200 || response.status >= 300) {
throw new Error(`HTTP error! status: ${response.status}`);
}
if (!response.body) {
if (!response.data) {
throw new Error('Response body is null');
}
@@ -276,15 +149,12 @@ async function fetchAndParseEpgStreaming(url: string): Promise<void> {
PROGRAM_BATCH_SIZE
);
// Convert web stream to Node.js stream
const nodeStream = Readable.fromWeb(response.body as any);
return new Promise((resolve, reject) => {
let dataStream: Readable = nodeStream;
let dataStream: Readable = response.data;
if (isGzipped) {
const gunzip = createGunzip();
dataStream = nodeStream.pipe(gunzip);
dataStream = response.data.pipe(gunzip);
gunzip.on('error', (err) => {
console.error(loggerLabel, 'Gunzip error:', err);
@@ -364,16 +234,12 @@ async function fetchAndParseEpgStreaming(url: string): Promise<void> {
* Runs in worker thread to avoid blocking main thread
*/
function clearAllEpgData(): void {
const dbPath = getIptvnatorDatabasePath();
const db = new Database(dbPath);
const clearOperation = new EpgDatabaseClearOperation(Database);
try {
console.log(loggerLabel, 'Clearing all EPG data...');
// Delete programs first (foreign key constraint)
db.exec('DELETE FROM epg_programs');
// Then delete channels
db.exec('DELETE FROM epg_channels');
clearOperation.run();
console.log(loggerLabel, 'All EPG data cleared');
@@ -387,7 +253,7 @@ function clearAllEpgData(): void {
};
parentPort?.postMessage(errorResponse);
} finally {
db.close();
clearOperation.close();
}
}
@@ -37,11 +37,25 @@ describe('shouldGunzipEpgResponse', () => {
).toBe(true);
});
it('supports plain Axios response header objects', () => {
expect(
shouldGunzipEpgResponse('https://example.com/guide', {
headers: {
'Content-Type': 'application/gzip',
},
url: 'https://example.com/guide',
})
).toBe(true);
});
it('returns true when content-disposition advertises a .gz filename', () => {
expect(
shouldGunzipEpgResponse('https://example.com/guide', {
headers: new Headers([
['content-disposition', 'attachment; filename="guide.xml.gz"'],
[
'content-disposition',
'attachment; filename="guide.xml.gz"',
],
]),
url: 'https://example.com/guide',
})
@@ -1,3 +1,20 @@
type HeaderReader =
| Record<string, unknown>
| {
get(name: string): unknown;
};
function getHeaderValue(headers: HeaderReader, name: string): string | null {
const value =
'get' in headers && typeof headers.get === 'function'
? headers.get(name)
: Object.entries(headers).find(
([headerName]) =>
headerName.toLowerCase() === name.toLowerCase()
)?.[1];
return value === null || value === undefined ? null : String(value);
}
function hasGzipPath(url: string | null | undefined): boolean {
if (!url) {
return false;
@@ -16,13 +33,13 @@ function hasGzipPath(url: string | null | undefined): boolean {
*/
export function shouldGunzipEpgResponse(
originalUrl: string,
response: { headers: Headers; url?: string }
response: { headers: HeaderReader; url?: string }
): boolean {
if (hasGzipPath(originalUrl) || hasGzipPath(response.url)) {
return true;
}
const contentType = response.headers.get('content-type');
const contentType = getHeaderValue(response.headers, 'content-type');
if (
contentType &&
/(application\/gzip|application\/x-gzip)/i.test(contentType)
@@ -30,7 +47,10 @@ export function shouldGunzipEpgResponse(
return true;
}
const contentDisposition = response.headers.get('content-disposition');
const contentDisposition = getHeaderValue(
response.headers,
'content-disposition'
);
if (contentDisposition?.toLowerCase().includes('.gz')) {
return true;
}
@@ -1,8 +1,11 @@
import axios from 'axios';
import { parentPort } from 'worker_threads';
import { parse } from 'iptv-playlist-parser';
import { readFile } from 'node:fs/promises';
import { createPlaylistObject, getFilenameFromUrl } from '@iptvnator/shared/m3u-utils';
import {
createPlaylistObject,
getFilenameFromUrl,
} from '@iptvnator/shared/m3u-utils';
import { createPlaylistAgentFactory } from '../util/secure-https';
import type {
Playlist,
PlaylistRefreshEvent,
@@ -12,8 +15,7 @@ import type {
PlaylistRefreshWorkerIncomingMessage,
PlaylistRefreshWorkerMessage,
} from './playlist-refresh.worker.types';
const https = require('https');
import { requestWithValidatedRedirects } from '../util/validated-axios';
type ActiveRefreshState = {
cancelled: boolean;
@@ -23,7 +25,9 @@ type ActiveRefreshState = {
const activeRefreshes = new Map<string, ActiveRefreshState>();
if (!parentPort) {
throw new Error('Playlist refresh worker must be started with a parent port');
throw new Error(
'Playlist refresh worker must be started with a parent port'
);
}
function postMessage(message: PlaylistRefreshWorkerMessage<Playlist>): void {
@@ -78,14 +82,16 @@ async function fetchPlaylistFromUrl(
emitEvent(payload, { status: 'started', phase: 'fetching' });
checkpoint(payload);
const agent = new https.Agent({
rejectUnauthorized: false,
});
const result = await axios.get(payload.url!, {
httpsAgent: agent,
signal: controller.signal,
timeout: 30000,
});
const result = await requestWithValidatedRedirects<string>(
payload.url!,
{
agentFactory: createPlaylistAgentFactory(),
method: 'GET',
signal: controller.signal,
timeout: 30000,
},
{ allowPrivateNetworks: true }
);
checkpoint(payload);
emitEvent(payload, { status: 'progress', phase: 'parsing' });
@@ -129,7 +135,9 @@ async function fetchPlaylistFromFile(
);
}
async function executeRefresh(payload: PlaylistRefreshPayload): Promise<Playlist> {
async function executeRefresh(
payload: PlaylistRefreshPayload
): Promise<Playlist> {
const controller = new AbortController();
activeRefreshes.set(payload.operationId, {
cancelled: false,
@@ -148,41 +156,45 @@ async function executeRefresh(payload: PlaylistRefreshPayload): Promise<Playlist
}
}
parentPort.on('message', async (message: PlaylistRefreshWorkerIncomingMessage) => {
if (message.type === 'cancel') {
const active = activeRefreshes.get(message.operationId);
if (active) {
active.cancelled = true;
active.controller.abort();
parentPort.on(
'message',
async (message: PlaylistRefreshWorkerIncomingMessage) => {
if (message.type === 'cancel') {
const active = activeRefreshes.get(message.operationId);
if (active) {
active.cancelled = true;
active.controller.abort();
}
return;
}
return;
}
try {
const result = await executeRefresh(message.payload);
postMessage({
type: 'response',
success: true,
result,
});
} catch (error) {
const payload = message.payload;
if (error instanceof Error && error.name === 'AbortError') {
emitEvent(payload, { status: 'cancelled', phase: 'parsing' });
} else {
emitEvent(payload, {
status: 'error',
phase: payload.url ? 'fetching' : 'reading-file',
error: error instanceof Error ? error.message : String(error),
try {
const result = await executeRefresh(message.payload);
postMessage({
type: 'response',
success: true,
result,
});
} catch (error) {
const payload = message.payload;
if (error instanceof Error && error.name === 'AbortError') {
emitEvent(payload, { status: 'cancelled', phase: 'parsing' });
} else {
emitEvent(payload, {
status: 'error',
phase: payload.url ? 'fetching' : 'reading-file',
error:
error instanceof Error ? error.message : String(error),
});
}
postMessage({
type: 'response',
success: false,
error: serializeError(error),
});
}
postMessage({
type: 'response',
success: false,
error: serializeError(error),
});
}
});
);
postMessage({ type: 'ready' });
+20 -9
View File
@@ -4,17 +4,25 @@ The download manager is a desktop-only feature that layers a curated queue, prog
## Backend responsibilities
- **Queue control (apps/electron-backend/src/app/events/downloads.events.ts)**
`DownloadTask` mirrors the shared `DownloadItem` table plus transient cancel/progress helpers. `enqueueDownload()` resolves a unique file path, persists a `queued` row in `downloads`, pushes the task onto `downloadQueue`, and triggers `processQueue()`. `processQueue()` keeps one active download, updates the row to `downloading`, and calls `startDownload()`.
- **electron-dl integration**
`startDownload()` now calls `electron-dl`’s `download()` helper. Headers (user agent, referer, origin) are attached, and the `onStarted`, `onProgress`, `onCompleted`, and `onCancel` callbacks translate the helper’s payload into Drizzle updates. The handler throttles progress broadcast, saves `filePath`/`fileName` from `electron-dl`, and marks failures/cancellations cleanly. Errors and cancellations delete partial files.
- **Queue control (`apps/electron-backend/src/app/events/database/download-runtime.ts`)**
`DownloadTask` mirrors the shared `DownloadItem` table plus transient cancel/progress helpers. Request validation and row creation live in `download-requests.ts`, while `downloads.events.ts` stays focused on IPC registration. `enqueueDownload()` pushes the task onto `downloadQueue` and triggers `processQueue()`. `processQueue()` keeps one active download, updates the row to `downloading`, and calls `startDownload()`.
- **electron-dl integration**
`startDownload()` calls `electron-dl`'s `download()` helper. Headers (user agent, referer, origin) are attached, and the `onStarted`, `onProgress`, `onCompleted`, and `onCancel` callbacks translate the helper's payload into Drizzle updates. A cancellation requested before `onStarted` is remembered and applied as soon as Electron supplies the `DownloadItem`, so the request cannot be lost in the startup race.
- **Destination collision policy**
Existing destination files are never overwritten. Before starting Electron's
download, the backend atomically reserves a free numbered filename with an
exclusive filesystem create. Electron may overwrite that empty reservation,
but cannot overwrite a file that existed before the reservation. The selected
`filePath` and `fileName` are persisted before transfer begins. Errors,
cancellations, and startup recovery remove that exact partial path and clear
it from the row; completed downloads replace it with Electron's final values.
- **IPC surface**
The backend exposes `DOWNLOADS_*` handlers for list retrieval, start/cancel/retry/remove operations, folder selection/reveal, and the `DOWNLOADS_UPDATE_EVENT` emitter that the renderer listens to in order to refresh its signal store.
## Renderer architecture
- **Downloads service** (`apps/web/src/app/services/downloads.service.ts`)
Signals back the current download list while `hasDownloads` and `isAvailable` gates UI rendering. Before each download the service resolves a download folder (stored in `SettingsStore` or fetched via `downloadsGetDefaultFolder`) and calls `downloadsStart`. The backend extracts the file extension from the URL or falls back to `mp4`. `onDownloadsUpdate` updates the signal, while helper methods `retryDownload`, `removeDownload`, `cancelDownload`, and `playDownload` talk to the corresponding IPC commands so retries reuse existing rows and completed items can open the recorded path.
- **Downloads service** (`libs/services/src/lib/downloads.service.ts`)
Signals back the current download list while `hasDownloads` and `isAvailable` gates UI rendering. Before each download the service asks the main process for the authorized folder and calls `downloadsStart`. The backend extracts the file extension from the URL or falls back to `mp4`. `onDownloadsUpdate` updates the signal, while helper methods `retryDownload`, `removeDownload`, `cancelDownload`, and `playDownload` talk to the corresponding IPC commands so retries reuse existing rows and completed items can open the recorded path.
- **Downloads view** (`libs/portal/downloads/feature`)
A standalone page exposes the queue, desktop-only messaging, folder picker, and action buttons. `downloads.component.html` now wraps the list inside a scrollable panel (`downloads__list-wrapper`) so long queues stay reachable, and `downloads.component.scss` drives a bold two-tone aesthetic inspired by the frontend-design mandate—gradient cards, floating avatars, and theme-aware variables triggered via `body.dark-theme`.
Failed/canceled cards now show retry/delete controls, queued/downloading cards show a cancel icon, and completed cards render inline play/open buttons with `mat-icon` cues. The header also shows the resolved download folder and a `CHANGE FOLDER` action.
@@ -33,9 +41,12 @@ The download manager is a desktop-only feature that layers a curated queue, prog
## Queuing, persistence, and UX notes
- Every download row writes to the shared `downloads` table with statuses (`queued`, `downloading`, `completed`, `failed`, `canceled`) plus metadata such as `bytesDownloaded`, `totalBytes`, `errorMessage`, and Xtream identifiers. Stale downloads reset to `failed` on startup.
- Queue cancellation removes the task or calls `downloadItem.cancel()` if the item is active; retries reuse the same database entry, preventing duplicate rows.
- Folder selection first checks stored preferences, falls back to the OS default downloads path, and finally prompts the user to pick a folder. The downloads service persists the chosen path via `SettingsStore`.
- Every download row writes to the shared `downloads` table with statuses (`queued`, `downloading`, `completed`, `failed`, `canceled`) plus metadata such as `bytesDownloaded`, `totalBytes`, `errorMessage`, and Xtream identifiers. On startup, `download-recovery.ts` deletes persisted partial reservations before stale queued/downloading rows become `failed`.
- Queue cancellation removes a queued task or records an active cancellation request and calls `downloadItem.cancel()` when the item is available; retries reuse the same database entry, preventing duplicate rows.
- The OS downloads path is always authorized. A custom folder becomes
authorized only after native folder selection, and the main process persists
that selection under Electron `userData`. Renderer settings may display the
path, but they are not trusted as authorization.
- The new UI leverages CSS variables for theme-specific backgrounds/borders, ensures `.downloads__list` can scroll inside its panel, and brings consistent badge/typography treatments to each card.
Keeping the backend queue, IPC handlers, shared schema, and renderer signals synchronized minimizes drift between platform rules and the UI. Future work might cover download list filters, cancel-all actions, or integration with upcoming playback analytics.
+46
View File
@@ -103,3 +103,49 @@ Rules:
When changing this flow, keep stale header cleanup covered. Switching from a
channel or playlist with custom headers to one without custom headers must clear
the previous override.
## Main-Process Remote Requests
Renderer-triggered HTTP requests must pass through the URL policy in
`apps/electron-backend/src/app/events/url-safety.ts`. The policy rejects
non-HTTP(S) URLs and embedded credentials, and strict callers also reject
loopback, private, reserved, and DNS-resolved private addresses. IPv4-mapped
IPv6 literals are decoded before classification, including hexadecimal forms
such as `::ffff:7f00:1`, so alternate IPv6 spelling cannot bypass IPv4 rules.
Remote request callers must use the validated Axios redirect helper so every
redirect target is checked before the main process follows it. Under the strict
policy, the helper pins the socket lookup to the IP addresses that passed
validation while retaining the original hostname for TLS SNI, certificate
validation, and virtual hosting. This prevents DNS rebinding between validation
and connection. Callers with custom TLS policy provide a typed agent factory;
the validated request layer supplies the pinned lookup instead of copying
private Node `Agent.options` state. Cross-origin redirects must not forward `Authorization`,
`Cookie`, `Proxy-Authorization`, Axios `params`, or request bodies.
EPG URLs are strict by default because an M3U playlist can supply them through
`url-tvg`. Operators who intentionally use a LAN-hosted EPG source can opt in
for that run with `IPTVNATOR_ALLOW_PRIVATE_NETWORK_URLS=1`. Directly configured
Xtream, Stalker, and playlist providers retain private-network support, but
still require HTTP(S), reject embedded credentials, and validate redirects.
Remote playlist TLS certificates are validated by default. The
`IPTVNATOR_ALLOW_INSECURE_TLS=1` escape hatch is only for explicitly trusted
providers with invalid or self-signed certificates.
## Filesystem Capabilities
Renderer IPC payloads are not filesystem authorization.
- `write-file` accepts only a path returned to the same renderer by the native
save dialog. The capability is single-use and is consumed before the write,
including when the filesystem operation fails.
- Download folders are owned by the Electron main process. The OS downloads
directory is always allowed; a custom directory is accepted only after the
native folder dialog selects it.
- The selected download directory is persisted under Electron `userData` and
returned by `DOWNLOADS_GET_DEFAULT_FOLDER`, so renderer-managed settings
cannot substitute an arbitrary host path.
- Downloads do not overwrite an existing destination file.
- Reveal and playback handlers accept only file paths recorded in IPTVnator's
downloads database.
@@ -16,16 +16,22 @@ import { SettingsStore } from './settings-store.service';
type TestDownloadsService = {
downloads: WritableSignal<DownloadItem[]>;
downloadFolder: WritableSignal<string>;
isAvailable: () => boolean;
isLoadingDownloads: Signal<boolean>;
hasLoadedDownloads: Signal<boolean>;
loadDownloads: DownloadsService['loadDownloads'];
loadDownloadFolder: DownloadsService['loadDownloadFolder'];
_isLoadingDownloads: WritableSignal<boolean>;
_hasLoadedDownloads: WritableSignal<boolean>;
loadDownloadsRequestId: number;
settingsStore: {
getDownloadFolder: () => string;
};
};
type DownloadsElectronStub = {
downloadsGetDefaultFolder?: jest.Mock<Promise<string>, []>;
downloadsGetList: jest.Mock<Promise<DownloadItem[]>, [string?]>;
};
@@ -67,18 +73,23 @@ describe('DownloadsService', () => {
const downloads = signal(initialDownloads);
const isLoadingDownloads = signal(false);
const hasLoadedDownloads = signal(false);
const downloadFolder = signal('');
const service = Object.create(
DownloadsService.prototype
) as TestDownloadsService;
Object.assign(service, {
downloads,
downloadFolder,
isAvailable: () => true,
_isLoadingDownloads: isLoadingDownloads,
isLoadingDownloads: isLoadingDownloads.asReadonly(),
_hasLoadedDownloads: hasLoadedDownloads,
hasLoadedDownloads: hasLoadedDownloads.asReadonly(),
loadDownloadsRequestId: 0,
settingsStore: {
getDownloadFolder: () => '/renderer-controlled',
},
});
return service;
@@ -132,6 +143,19 @@ describe('DownloadsService', () => {
expect(service.hasLoadedDownloads()).toBe(true);
});
it('uses the main-process authorized download folder instead of renderer storage', async () => {
const electron = {
downloadsGetDefaultFolder: jest.fn(async () => '/authorized'),
downloadsGetList: jest.fn(async () => []),
};
testWindow.electron = electron;
const service = createService();
await expect(service.loadDownloadFolder()).resolves.toBe('/authorized');
expect(service.downloadFolder()).toBe('/authorized');
expect(electron.downloadsGetDefaultFolder).toHaveBeenCalledTimes(1);
});
it('marks downloads as loaded after a failed request while preserving existing data', async () => {
const existing = createDownload(1);
const error = new Error('download query failed');
+3 -8
View File
@@ -137,14 +137,9 @@ export class DownloadsService implements OnDestroy {
async loadDownloadFolder(): Promise<string> {
if (!this.isAvailable()) return '';
// First check settings
const storedFolder = this.settingsStore.getDownloadFolder?.();
if (storedFolder) {
this.downloadFolder.set(storedFolder);
return storedFolder;
}
// Fall back to default
// The main process owns folder authorization. It returns either the OS
// default or a custom folder previously selected through a native
// dialog, rather than trusting renderer-managed settings.
try {
const defaultFolder =
await window.electron.downloadsGetDefaultFolder();
@@ -85,10 +85,7 @@
"
/>
</div>
<div
class="program-title"
[innerHTML]="program.title"
></div>
<div class="program-title">{{ program.title }}</div>
@if (isProgramPlaying(program)) {
<div class="program-progress">
<div class="epg-progress-track">
@@ -67,6 +67,31 @@ describe('HtmlVideoPlayerComponent', () => {
expect(component).toBeTruthy();
});
it('detaches volume/metadata/timeupdate listeners on destroy (no leak)', () => {
const el = component.videoPlayer.nativeElement;
const removeSpy = jest.spyOn(el, 'removeEventListener');
const handlers = component as unknown as {
handleVolumeChange: EventListener;
handleLoadedMetadata: EventListener;
handleTimeUpdate: EventListener;
};
fixture.destroy();
expect(removeSpy).toHaveBeenCalledWith(
'volumechange',
handlers.handleVolumeChange
);
expect(removeSpy).toHaveBeenCalledWith(
'loadedmetadata',
handlers.handleLoadedMetadata
);
expect(removeSpy).toHaveBeenCalledWith(
'timeupdate',
handlers.handleTimeUpdate
);
});
it('should call play channel function after input changes', () => {
jest.spyOn(component, 'playChannel');
jest.spyOn(global.console, 'error').mockImplementation(() => {
@@ -79,26 +79,38 @@ export class HtmlVideoPlayerComponent implements OnInit, OnChanges, OnDestroy {
this.playbackIssue.emit(null);
};
ngOnInit() {
this.videoPlayer.nativeElement.addEventListener('volumechange', () => {
this.onVolumeChange();
private readonly handleVolumeChange = (): void => {
this.onVolumeChange();
};
private readonly handleLoadedMetadata = (): void => {
if (this.startTime > 0) {
this.videoPlayer.nativeElement.currentTime = this.startTime;
}
};
private readonly handleTimeUpdate = (): void => {
this.timeUpdate.emit({
currentTime: this.videoPlayer.nativeElement.currentTime,
duration: this.videoPlayer.nativeElement.duration,
});
};
ngOnInit() {
this.videoPlayer.nativeElement.addEventListener(
'volumechange',
this.handleVolumeChange
);
this.videoPlayer.nativeElement.addEventListener(
'loadedmetadata',
() => {
if (this.startTime > 0) {
this.videoPlayer.nativeElement.currentTime = this.startTime;
}
}
this.handleLoadedMetadata
);
this.videoPlayer.nativeElement.addEventListener('timeupdate', () => {
this.timeUpdate.emit({
currentTime: this.videoPlayer.nativeElement.currentTime,
duration: this.videoPlayer.nativeElement.duration,
});
});
this.videoPlayer.nativeElement.addEventListener(
'timeupdate',
this.handleTimeUpdate
);
this.videoPlayer.nativeElement.addEventListener(
'error',
@@ -332,7 +344,15 @@ export class HtmlVideoPlayerComponent implements OnInit, OnChanges, OnDestroy {
ngOnDestroy(): void {
this.videoPlayer.nativeElement.removeEventListener(
'volumechange',
this.onVolumeChange
this.handleVolumeChange
);
this.videoPlayer.nativeElement.removeEventListener(
'loadedmetadata',
this.handleLoadedMetadata
);
this.videoPlayer.nativeElement.removeEventListener(
'timeupdate',
this.handleTimeUpdate
);
this.videoPlayer.nativeElement.removeEventListener(
'error',
@@ -1,29 +1,41 @@
@if (selectedPlayer() === 'videojs') {
<app-vjs-player
[options]="vjsOptions"
[volume]="volume()"
[startTime]="startTime()"
(timeUpdate)="timeUpdate.emit($event)"
(playbackIssue)="handlePlaybackIssue($event)"
/>
@defer (on immediate) {
<app-vjs-player
[options]="vjsOptions"
[volume]="volume()"
[startTime]="startTime()"
(timeUpdate)="timeUpdate.emit($event)"
(playbackIssue)="handlePlaybackIssue($event)"
/>
} @placeholder {
<div class="web-player-defer-placeholder" aria-hidden="true"></div>
}
} @else if (selectedPlayer() === 'html5') {
<app-html-video-player
[channel]="$any(channel)"
[volume]="volume()"
[showCaptions]="showCaptions()"
[startTime]="startTime()"
(timeUpdate)="timeUpdate.emit($event)"
(playbackIssue)="handlePlaybackIssue($event)"
/>
@defer (on immediate) {
<app-html-video-player
[channel]="$any(channel)"
[volume]="volume()"
[showCaptions]="showCaptions()"
[startTime]="startTime()"
(timeUpdate)="timeUpdate.emit($event)"
(playbackIssue)="handlePlaybackIssue($event)"
/>
} @placeholder {
<div class="web-player-defer-placeholder" aria-hidden="true"></div>
}
} @else if (selectedPlayer() === 'artplayer') {
<app-art-player
[channel]="$any(channel)"
[volume]="volume()"
[showCaptions]="showCaptions()"
[startTime]="startTime()"
(timeUpdate)="timeUpdate.emit($event)"
(playbackIssue)="handlePlaybackIssue($event)"
/>
@defer (on immediate) {
<app-art-player
[channel]="$any(channel)"
[volume]="volume()"
[showCaptions]="showCaptions()"
[startTime]="startTime()"
(timeUpdate)="timeUpdate.emit($event)"
(playbackIssue)="handlePlaybackIssue($event)"
/>
} @placeholder {
<div class="web-player-defer-placeholder" aria-hidden="true"></div>
}
} @else if (selectedPlayer() === 'embedded-mpv') {
<app-embedded-mpv-player
[playback]="resolvedPlayback()"
@@ -35,13 +47,17 @@
(nextEpisodeRequested)="nextEpisodeRequested.emit()"
/>
} @else {
<app-vjs-player
[options]="vjsOptions"
[volume]="volume()"
[startTime]="startTime()"
(timeUpdate)="timeUpdate.emit($event)"
(playbackIssue)="handlePlaybackIssue($event)"
/>
@defer (on immediate) {
<app-vjs-player
[options]="vjsOptions"
[volume]="volume()"
[startTime]="startTime()"
(timeUpdate)="timeUpdate.emit($event)"
(playbackIssue)="handlePlaybackIssue($event)"
/>
} @placeholder {
<div class="web-player-defer-placeholder" aria-hidden="true"></div>
}
}
@if (visiblePlaybackDiagnostic(); as issue) {
@@ -14,6 +14,15 @@ app-html-video-player {
height: 100%;
}
// Holds the player's space (black, no layout shift) for the one-frame gap
// before a @defer (on immediate) player chunk resolves.
.web-player-defer-placeholder {
display: block;
width: 100%;
height: 100%;
background: #000;
}
.web-player-diagnostic {
position: absolute;
inset: 0;
@@ -1,5 +1,9 @@
import { Component, input, output } from '@angular/core';
import { ComponentFixture, TestBed } from '@angular/core/testing';
import {
ComponentFixture,
DeferBlockBehavior,
TestBed,
} from '@angular/core/testing';
import { ClipboardModule } from '@angular/cdk/clipboard';
import { MatButtonModule } from '@angular/material/button';
import { MatIconModule } from '@angular/material/icon';
@@ -95,6 +99,8 @@ describe('WebPlayerViewComponent', () => {
runtimeCapabilities = { supportsManagedExternalPlayers: false };
await TestBed.configureTestingModule({
// @defer blocks render their main content synchronously in tests.
deferBlockBehavior: DeferBlockBehavior.Playthrough,
imports: [WebPlayerViewComponent, TranslateModule.forRoot()],
providers: [
{ provide: StorageMap, useValue: storageMap },
@@ -208,7 +214,7 @@ describe('WebPlayerViewComponent', () => {
]);
});
it('marks portal VOD playback as non-live for Video.js MPEG-TS playback', () => {
it('marks portal VOD playback as non-live for Video.js MPEG-TS playback', async () => {
const streamUrl = 'https://example.com/movie/123.ts';
fixture.componentRef.setInput('playback', {
streamUrl,
@@ -220,6 +226,8 @@ describe('WebPlayerViewComponent', () => {
},
});
fixture.detectChanges();
await fixture.whenStable();
fixture.detectChanges();
const player = fixture.debugElement.query(
@@ -238,7 +246,7 @@ describe('WebPlayerViewComponent', () => {
);
});
it('preserves playback HTTP metadata for channel-based players', () => {
it('preserves playback HTTP metadata for channel-based players', async () => {
const streamUrl = 'https://example.com/live/channel.m3u8';
fixture.componentRef.setInput(
'playerOverride',
@@ -257,6 +265,8 @@ describe('WebPlayerViewComponent', () => {
},
});
fixture.detectChanges();
await fixture.whenStable();
fixture.detectChanges();
const player = fixture.debugElement.query(
@@ -275,7 +285,7 @@ describe('WebPlayerViewComponent', () => {
);
});
it('falls back to playback headers when explicit HTTP metadata is absent', () => {
it('falls back to playback headers when explicit HTTP metadata is absent', async () => {
const streamUrl = 'https://example.com/live/channel.m3u8';
fixture.componentRef.setInput(
'playerOverride',
@@ -291,6 +301,8 @@ describe('WebPlayerViewComponent', () => {
},
});
fixture.detectChanges();
await fixture.whenStable();
fixture.detectChanges();
const player = fixture.debugElement.query(
@@ -376,18 +388,25 @@ describe('WebPlayerViewComponent', () => {
canNext: false,
autoplayEnabled: true,
};
fixture.componentRef.setInput('playerOverride', VideoPlayer.EmbeddedMpv);
fixture.componentRef.setInput(
'playerOverride',
VideoPlayer.EmbeddedMpv
);
fixture.componentRef.setInput('seriesNavigation', seriesNavigation);
(
component as unknown as {
playbackEnded: { subscribe: (fn: () => void) => void };
previousEpisodeRequested: { subscribe: (fn: () => void) => void };
previousEpisodeRequested: {
subscribe: (fn: () => void) => void;
};
nextEpisodeRequested: { subscribe: (fn: () => void) => void };
}
).playbackEnded.subscribe(() => events.push('ended'));
(
component as unknown as {
previousEpisodeRequested: { subscribe: (fn: () => void) => void };
previousEpisodeRequested: {
subscribe: (fn: () => void) => void;
};
}
).previousEpisodeRequested.subscribe(() => events.push('previous'));
(
@@ -494,7 +513,9 @@ describe('WebPlayerViewComponent', () => {
);
});
it('clears playback diagnostics when retrying inline playback', () => {
it('clears playback diagnostics when retrying inline playback', async () => {
fixture.detectChanges();
await fixture.whenStable();
fixture.detectChanges();
const player = fixture.debugElement.query(
By.directive(StubVjsPlayerComponent)
@@ -47,17 +47,23 @@ import {
toTimestamp,
} from './dashboard-mappers';
import {
buildStalkerDetailNavigationTarget,
buildStalkerStateItem,
buildXtreamNavigationTarget,
getGlobalFavoriteNavigation,
getRecentItemNavigation,
PORTAL_PLAYBACK_POSITIONS,
WorkspaceNavigationTarget,
} from '@iptvnator/portal/shared/util';
import type { PlaybackPositionData } from '@iptvnator/shared/interfaces';
import {
getGlobalFavoriteLink as getGlobalFavoriteLinkUtil,
getGlobalFavoriteNavigationState as getGlobalFavoriteNavigationStateUtil,
getPlaylistLink as getPlaylistLinkUtil,
getRecentItemLink as getRecentItemLinkUtil,
getRecentItemNavigationState as getRecentItemNavigationStateUtil,
getRecentlyAddedLink as getRecentlyAddedLinkUtil,
getRecentlyAddedNavigationState as getRecentlyAddedNavigationStateUtil,
isTypeInKind as isTypeInKindUtil,
type DashboardContentKind,
} from './dashboard-navigation.util';
export type DashboardContentKind = 'all' | 'channels' | 'vod' | 'series';
export type { DashboardContentKind };
// Compound key for looking up a playback position by recent item — a single
// playlist can contain the same xtream-id for a VOD and an episode (rare,
@@ -799,28 +805,11 @@ export class DashboardDataService {
type: PortalActivityType,
kind: DashboardContentKind
): boolean {
if (kind === 'all') {
return true;
}
if (kind === 'channels') {
return type === 'live';
}
if (kind === 'vod') {
return type === 'movie';
}
return type === 'series';
return isTypeInKindUtil(type, kind);
}
getPlaylistLink(playlist: PlaylistMeta): string[] {
if (playlist.serverUrl) {
return ['/workspace', 'xtreams', playlist._id, 'vod'];
}
if (playlist.macAddress) {
return ['/workspace', 'stalker', playlist._id, 'vod'];
}
return ['/workspace', 'playlists', playlist._id];
return getPlaylistLinkUtil(playlist);
}
getPlaylistProvider(playlist: PlaylistMeta): string {
@@ -858,13 +847,13 @@ export class DashboardDataService {
}
getRecentItemLink(item: GlobalRecentItem): string[] {
return getRecentItemNavigation(item).link;
return getRecentItemLinkUtil(item);
}
getRecentItemNavigationState(
item: GlobalRecentItem
): WorkspaceNavigationTarget['state'] {
return getRecentItemNavigation(item).state;
return getRecentItemNavigationStateUtil(item);
}
async removeGlobalRecentItem(item: GlobalRecentItem): Promise<void> {
@@ -965,67 +954,23 @@ export class DashboardDataService {
}
getGlobalFavoriteLink(item: DashboardFavoriteItem): string[] {
return getGlobalFavoriteNavigation(item).link;
return getGlobalFavoriteLinkUtil(item);
}
getGlobalFavoriteNavigationState(
item: DashboardFavoriteItem
): WorkspaceNavigationTarget['state'] {
return getGlobalFavoriteNavigation(item).state;
return getGlobalFavoriteNavigationStateUtil(item);
}
getRecentlyAddedLink(item: DashboardRecentlyAddedItem): string[] {
if (item.source === 'stalker' && item.type !== 'live') {
return buildStalkerDetailNavigationTarget({
playlistId: item.playlist_id,
type: item.type,
categoryId: item.category_id,
item: buildStalkerStateItem(item.stalker_item, {
id: item.id,
title: item.title,
type: item.type,
category_id: item.category_id,
poster_url: item.poster_url,
}),
}).link;
}
return buildXtreamNavigationTarget({
playlistId: item.playlist_id,
type: item.type,
categoryId: item.category_id,
itemId: item.xtream_id,
title: item.title,
imageUrl: item.poster_url,
}).link;
return getRecentlyAddedLinkUtil(item);
}
getRecentlyAddedNavigationState(
item: DashboardRecentlyAddedItem
): WorkspaceNavigationTarget['state'] {
if (item.source === 'stalker' && item.type !== 'live') {
return buildStalkerDetailNavigationTarget({
playlistId: item.playlist_id,
type: item.type,
categoryId: item.category_id,
item: buildStalkerStateItem(item.stalker_item, {
id: item.id,
title: item.title,
type: item.type,
category_id: item.category_id,
poster_url: item.poster_url,
}),
}).state;
}
return buildXtreamNavigationTarget({
playlistId: item.playlist_id,
type: item.type,
categoryId: item.category_id,
itemId: item.xtream_id,
title: item.title,
imageUrl: item.poster_url,
}).state;
return getRecentlyAddedNavigationStateUtil(item);
}
async removeGlobalFavorite(item: DashboardFavoriteItem): Promise<void> {
@@ -0,0 +1,128 @@
import {
PlaylistMeta,
PortalActivityType,
PortalAddedItem,
PortalFavoriteItem,
PortalRecentItem,
} from '@iptvnator/shared/interfaces';
import {
buildStalkerDetailNavigationTarget,
buildStalkerStateItem,
buildXtreamNavigationTarget,
getGlobalFavoriteNavigation,
getRecentItemNavigation,
WorkspaceNavigationTarget,
} from '@iptvnator/portal/shared/util';
/**
* Pure navigation/link helpers for dashboard items.
*
* Extracted from `DashboardDataService` so the routing logic can be unit-tested
* in isolation and the service stays a thin facade. None of these functions
* touch component/service state — they map a dashboard item to a router link
* (and optional navigation state) using the shared portal navigation builders.
*/
export type DashboardContentKind = 'all' | 'channels' | 'vod' | 'series';
export function isTypeInKind(
type: PortalActivityType,
kind: DashboardContentKind
): boolean {
if (kind === 'all') {
return true;
}
if (kind === 'channels') {
return type === 'live';
}
if (kind === 'vod') {
return type === 'movie';
}
return type === 'series';
}
export function getPlaylistLink(playlist: PlaylistMeta): string[] {
if (playlist.serverUrl) {
return ['/workspace', 'xtreams', playlist._id, 'vod'];
}
if (playlist.macAddress) {
return ['/workspace', 'stalker', playlist._id, 'vod'];
}
return ['/workspace', 'playlists', playlist._id];
}
export function getRecentItemLink(item: PortalRecentItem): string[] {
return getRecentItemNavigation(item).link;
}
export function getRecentItemNavigationState(
item: PortalRecentItem
): WorkspaceNavigationTarget['state'] {
return getRecentItemNavigation(item).state;
}
export function getGlobalFavoriteLink(item: PortalFavoriteItem): string[] {
return getGlobalFavoriteNavigation(item).link;
}
export function getGlobalFavoriteNavigationState(
item: PortalFavoriteItem
): WorkspaceNavigationTarget['state'] {
return getGlobalFavoriteNavigation(item).state;
}
export function getRecentlyAddedLink(item: PortalAddedItem): string[] {
if (item.source === 'stalker' && item.type !== 'live') {
return buildStalkerDetailNavigationTarget({
playlistId: item.playlist_id,
type: item.type,
categoryId: item.category_id,
item: buildStalkerStateItem(item.stalker_item, {
id: item.id,
title: item.title,
type: item.type,
category_id: item.category_id,
poster_url: item.poster_url,
}),
}).link;
}
return buildXtreamNavigationTarget({
playlistId: item.playlist_id,
type: item.type,
categoryId: item.category_id,
itemId: item.xtream_id,
title: item.title,
imageUrl: item.poster_url,
}).link;
}
export function getRecentlyAddedNavigationState(
item: PortalAddedItem
): WorkspaceNavigationTarget['state'] {
if (item.source === 'stalker' && item.type !== 'live') {
return buildStalkerDetailNavigationTarget({
playlistId: item.playlist_id,
type: item.type,
categoryId: item.category_id,
item: buildStalkerStateItem(item.stalker_item, {
id: item.id,
title: item.title,
type: item.type,
category_id: item.category_id,
poster_url: item.poster_url,
}),
}).state;
}
return buildXtreamNavigationTarget({
playlistId: item.playlist_id,
type: item.type,
categoryId: item.category_id,
itemId: item.xtream_id,
title: item.title,
imageUrl: item.poster_url,
}).state;
}