mirror of
https://github.com/4gray/iptvnator.git
synced 2026-10-08 17:06:15 -08:00
feat: implement playlist refresh functionality and cache management
This commit is contained in:
1 parent
4d047f89ca
commit
4e0afe8b20
16 files changed
+813
-45
No files matched your search
@@ -44,6 +44,17 @@ async function buildWorker() {
|
||||
'../../dist/apps/electron-backend/workers/database.worker.js'
|
||||
),
|
||||
},
|
||||
{
|
||||
label: 'playlist refresh worker',
|
||||
entry: path.join(
|
||||
__dirname,
|
||||
'src/app/workers/playlist-refresh.worker.ts'
|
||||
),
|
||||
outfile: path.join(
|
||||
__dirname,
|
||||
'../../dist/apps/electron-backend/workers/playlist-refresh.worker.js'
|
||||
),
|
||||
},
|
||||
];
|
||||
|
||||
for (const worker of workers) {
|
||||
@@ -70,6 +81,10 @@ async function buildWorker() {
|
||||
__dirname,
|
||||
'../../libs/shared/interfaces/src/index.ts'
|
||||
),
|
||||
'm3u-utils': path.join(
|
||||
__dirname,
|
||||
'../../libs/shared/m3u-utils/src/index.ts'
|
||||
),
|
||||
'database': path.join(
|
||||
__dirname,
|
||||
'../../libs/shared/database/src/index.ts'
|
||||
|
||||
@@ -1,9 +1,14 @@
|
||||
import { contextBridge, ipcRenderer } from 'electron';
|
||||
import type { ExternalPlayerSession } from 'shared-interfaces';
|
||||
import type {
|
||||
ExternalPlayerSession,
|
||||
PlaylistRefreshEvent,
|
||||
PlaylistRefreshPayload,
|
||||
} from 'shared-interfaces';
|
||||
|
||||
const PORTAL_DEBUG_EVENT = 'PORTAL_DEBUG_EVENT';
|
||||
const EXTERNAL_PLAYER_SESSION_UPDATE = 'EXTERNAL_PLAYER_SESSION_UPDATE';
|
||||
const DB_OPERATION_EVENT = 'DB_OPERATION_EVENT';
|
||||
const PLAYLIST_REFRESH_EVENT = 'PLAYLIST:REFRESH_EVENT';
|
||||
|
||||
type PortalDebugEvent = {
|
||||
requestId: string;
|
||||
@@ -112,6 +117,14 @@ contextBridge.exposeInMainWorld('electron', {
|
||||
ipcRenderer.on(DB_OPERATION_EVENT, handler);
|
||||
return () => ipcRenderer.off(DB_OPERATION_EVENT, handler);
|
||||
},
|
||||
onPlaylistRefreshEvent: (
|
||||
callback: (data: PlaylistRefreshEvent) => void
|
||||
) => {
|
||||
const handler = (_event: Electron.IpcRendererEvent, data: any) =>
|
||||
callback(data as PlaylistRefreshEvent);
|
||||
ipcRenderer.on(PLAYLIST_REFRESH_EVENT, handler);
|
||||
return () => ipcRenderer.off(PLAYLIST_REFRESH_EVENT, handler);
|
||||
},
|
||||
// DB save content progress listener
|
||||
onDbSaveContentProgress: (callback: (count: number) => void) => {
|
||||
const handler = (
|
||||
@@ -235,8 +248,15 @@ contextBridge.exposeInMainWorld('electron', {
|
||||
url: string;
|
||||
params: Record<string, string>;
|
||||
requestId?: string;
|
||||
sessionId?: string;
|
||||
suppressErrorLog?: boolean;
|
||||
}) => ipcRenderer.invoke('XTREAM_REQUEST', payload),
|
||||
xtreamCancelSession: (sessionId: string) =>
|
||||
ipcRenderer.invoke('XTREAM_CANCEL_SESSION', sessionId),
|
||||
refreshPlaylist: (payload: PlaylistRefreshPayload) =>
|
||||
ipcRenderer.invoke('PLAYLIST:REFRESH', payload),
|
||||
cancelPlaylistRefresh: (operationId: string) =>
|
||||
ipcRenderer.invoke('PLAYLIST:CANCEL_REFRESH', operationId),
|
||||
// Database operations
|
||||
dbCreatePlaylist: (playlist: any) =>
|
||||
ipcRenderer.invoke('DB_CREATE_PLAYLIST', playlist),
|
||||
@@ -314,6 +334,10 @@ contextBridge.exposeInMainWorld('electron', {
|
||||
type,
|
||||
operationId
|
||||
),
|
||||
dbClearXtreamImportCache: (
|
||||
playlistId: string,
|
||||
type: 'live' | 'movie' | 'series'
|
||||
) => ipcRenderer.invoke('DB_CLEAR_XTREAM_IMPORT_CACHE', playlistId, type),
|
||||
dbSearchContent: (
|
||||
playlistId: string,
|
||||
searchTerm: string,
|
||||
|
||||
@@ -3,6 +3,7 @@ import * as schema from 'database-schema';
|
||||
import type { AppDatabase } from '../database.types';
|
||||
import {
|
||||
checkpointOperation,
|
||||
chunkValues,
|
||||
type OperationControl,
|
||||
reportOperationProgress,
|
||||
} from './operation-control';
|
||||
@@ -292,6 +293,50 @@ export async function saveContent(
|
||||
return { success: true, count: totalInserted };
|
||||
}
|
||||
|
||||
export async function clearXtreamImportCache(
|
||||
db: AppDatabase,
|
||||
playlistId: string,
|
||||
type: 'live' | 'movie' | 'series'
|
||||
): Promise<{ success: boolean }> {
|
||||
const dbType =
|
||||
type === 'series' ? 'series' : type === 'movie' ? 'movies' : 'live';
|
||||
|
||||
const categoryRows = await db
|
||||
.select({ id: schema.categories.id })
|
||||
.from(schema.categories)
|
||||
.where(
|
||||
and(
|
||||
eq(schema.categories.playlistId, playlistId),
|
||||
eq(schema.categories.type, dbType)
|
||||
)
|
||||
);
|
||||
|
||||
const categoryIds = categoryRows.map((category) => category.id);
|
||||
if (categoryIds.length === 0) {
|
||||
return { success: true };
|
||||
}
|
||||
|
||||
const contentRows = await db
|
||||
.select({ id: schema.content.id })
|
||||
.from(schema.content)
|
||||
.where(inArray(schema.content.categoryId, categoryIds));
|
||||
|
||||
for (const chunk of chunkValues(
|
||||
contentRows.map((row) => row.id),
|
||||
100
|
||||
)) {
|
||||
await db.delete(schema.content).where(inArray(schema.content.id, chunk));
|
||||
}
|
||||
|
||||
for (const chunk of chunkValues(categoryIds, 100)) {
|
||||
await db
|
||||
.delete(schema.categories)
|
||||
.where(inArray(schema.categories.id, chunk));
|
||||
}
|
||||
|
||||
return { success: true };
|
||||
}
|
||||
|
||||
export async function getContentByXtreamId(
|
||||
db: AppDatabase,
|
||||
xtreamId: number,
|
||||
|
||||
@@ -0,0 +1,36 @@
|
||||
import { parseAppPlaylist } from './playlist.operations';
|
||||
|
||||
describe('playlist.operations', () => {
|
||||
it('hydrates updateDate from lastUpdated when payload is stale', async () => {
|
||||
const parsed = parseAppPlaylist({
|
||||
id: 'playlist-1',
|
||||
name: 'Refresh Xtream Source',
|
||||
serverUrl: 'http://localhost:8080',
|
||||
username: 'demo',
|
||||
password: 'secret',
|
||||
dateCreated: '2026-04-03T08:00:00.000Z',
|
||||
lastUpdated: '2026-04-03T11:15:00.000Z',
|
||||
type: 'xtream',
|
||||
autoRefresh: false,
|
||||
count: 0,
|
||||
importDate: '2026-04-03T08:00:00.000Z',
|
||||
payload: JSON.stringify({
|
||||
_id: 'playlist-1',
|
||||
title: 'Refresh Xtream Source',
|
||||
count: 0,
|
||||
importDate: '2026-04-03T08:00:00.000Z',
|
||||
autoRefresh: false,
|
||||
serverUrl: 'http://localhost:8080',
|
||||
username: 'demo',
|
||||
password: 'secret',
|
||||
}),
|
||||
} as any);
|
||||
|
||||
expect(parsed).toEqual(
|
||||
expect.objectContaining({
|
||||
_id: 'playlist-1',
|
||||
updateDate: new Date('2026-04-03T11:15:00.000Z').getTime(),
|
||||
})
|
||||
);
|
||||
});
|
||||
});
|
||||
@@ -139,26 +139,12 @@ function buildPlaylistRow(
|
||||
};
|
||||
}
|
||||
|
||||
function parseAppPlaylist(row: schema.Playlist): Record<string, unknown> {
|
||||
export function parseAppPlaylist(row: schema.Playlist): Record<string, unknown> {
|
||||
const payload = parseJsonValue<Record<string, unknown> | null>(
|
||||
row.payload,
|
||||
null
|
||||
);
|
||||
|
||||
if (payload && typeof payload === 'object') {
|
||||
return {
|
||||
...payload,
|
||||
_id:
|
||||
getStringValue(payload._id) ??
|
||||
getStringValue(payload.id) ??
|
||||
row.id,
|
||||
title:
|
||||
getStringValue(payload.title) ??
|
||||
getStringValue(payload.name) ??
|
||||
row.name,
|
||||
};
|
||||
}
|
||||
|
||||
const base = payload && typeof payload === 'object' ? payload : {};
|
||||
const favorites = parseJsonValue<unknown[]>(row.favorites, []);
|
||||
const recentlyViewed = parseJsonValue<unknown[]>(row.recentlyViewed, []);
|
||||
const importDate =
|
||||
@@ -166,30 +152,42 @@ function parseAppPlaylist(row: schema.Playlist): Record<string, unknown> {
|
||||
const portalUrl =
|
||||
row.portalUrl ??
|
||||
(row.type === PLAYLIST_TYPES.STALKER ? row.url : null);
|
||||
const updateDate =
|
||||
row.updateDate ??
|
||||
(row.lastUpdated ? new Date(row.lastUpdated).getTime() : undefined);
|
||||
|
||||
return {
|
||||
...base,
|
||||
_id: row.id,
|
||||
title: row.name,
|
||||
count: row.count ?? 0,
|
||||
importDate,
|
||||
lastUsage: row.lastUsage ?? importDate,
|
||||
title:
|
||||
getStringValue(base.title) ??
|
||||
getStringValue(base.name) ??
|
||||
row.name,
|
||||
count: row.count ?? getNumericValue(base.count) ?? 0,
|
||||
importDate: getStringValue(base.importDate) ?? importDate,
|
||||
lastUsage:
|
||||
row.lastUsage ??
|
||||
getStringValue(base.lastUsage) ??
|
||||
getStringValue(base.importDate) ??
|
||||
importDate,
|
||||
favorites,
|
||||
recentlyViewed,
|
||||
autoRefresh: row.autoRefresh ?? false,
|
||||
url: row.type === PLAYLIST_TYPES.M3U_URL ? row.url : undefined,
|
||||
filePath: row.filePath ?? undefined,
|
||||
userAgent: row.userAgent ?? undefined,
|
||||
referrer: row.referrer ?? undefined,
|
||||
origin: row.origin ?? undefined,
|
||||
updateDate:
|
||||
row.updateDate ??
|
||||
(row.lastUpdated ? new Date(row.lastUpdated).getTime() : undefined),
|
||||
position: row.position ?? undefined,
|
||||
serverUrl: row.serverUrl ?? undefined,
|
||||
username: row.username ?? undefined,
|
||||
password: row.password ?? undefined,
|
||||
macAddress: row.macAddress ?? undefined,
|
||||
portalUrl: portalUrl ?? undefined,
|
||||
autoRefresh: row.autoRefresh ?? Boolean(base.autoRefresh),
|
||||
url:
|
||||
row.type === PLAYLIST_TYPES.M3U_URL
|
||||
? row.url ?? getStringValue(base.url)
|
||||
: getStringValue(base.url),
|
||||
filePath: row.filePath ?? getStringValue(base.filePath),
|
||||
userAgent: row.userAgent ?? getStringValue(base.userAgent),
|
||||
referrer: row.referrer ?? getStringValue(base.referrer),
|
||||
origin: row.origin ?? getStringValue(base.origin),
|
||||
updateDate,
|
||||
position: row.position ?? getNumericValue(base.position),
|
||||
serverUrl: row.serverUrl ?? getStringValue(base.serverUrl),
|
||||
username: row.username ?? getStringValue(base.username),
|
||||
password: row.password ?? getStringValue(base.password),
|
||||
macAddress: row.macAddress ?? getStringValue(base.macAddress),
|
||||
portalUrl: portalUrl ?? getStringValue(base.portalUrl),
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@@ -60,6 +60,14 @@ ipcMain.handle(
|
||||
}
|
||||
);
|
||||
|
||||
handleWorkerRequest(
|
||||
'DB_CLEAR_XTREAM_IMPORT_CACHE',
|
||||
(playlistId: string, type: 'live' | 'movie' | 'series') => ({
|
||||
playlistId,
|
||||
type,
|
||||
})
|
||||
);
|
||||
|
||||
handleWorkerRequest(
|
||||
'DB_GET_CONTENT_BY_XTREAM_ID',
|
||||
(xtreamId: number, playlistId: string) => ({
|
||||
|
||||
@@ -4,12 +4,27 @@
|
||||
*/
|
||||
|
||||
import axios from 'axios';
|
||||
import { dialog, ipcMain } from 'electron';
|
||||
import { app, dialog, ipcMain, WebContents } from 'electron';
|
||||
import { parse } from 'iptv-playlist-parser';
|
||||
import { createPlaylistObject, getFilenameFromUrl } from 'm3u-utils';
|
||||
import { readFile, writeFile } from 'node:fs/promises';
|
||||
import { basename } from 'node:path';
|
||||
import { AUTO_UPDATE_PLAYLISTS } from 'shared-interfaces';
|
||||
import { pathToFileURL } from 'url';
|
||||
import { Worker } from 'worker_threads';
|
||||
import {
|
||||
AUTO_UPDATE_PLAYLISTS,
|
||||
PLAYLIST_CANCEL_REFRESH,
|
||||
PLAYLIST_REFRESH,
|
||||
PLAYLIST_REFRESH_EVENT,
|
||||
Playlist,
|
||||
PlaylistRefreshEvent,
|
||||
PlaylistRefreshPayload,
|
||||
} from 'shared-interfaces';
|
||||
import { resolveWorkerRuntimeBootstrap } from '../workers/worker-runtime-paths';
|
||||
import type {
|
||||
PlaylistRefreshWorkerMessage,
|
||||
PlaylistRefreshWorkerResponseMessage,
|
||||
} from '../workers/playlist-refresh.worker.types';
|
||||
|
||||
export default class PlaylistEvents {
|
||||
static bootstrapPlaylistEvents(): Electron.IpcMain {
|
||||
@@ -19,6 +34,15 @@ export default class PlaylistEvents {
|
||||
|
||||
const https = require('https');
|
||||
|
||||
type ActivePlaylistRefresh = {
|
||||
reject: (reason?: unknown) => void;
|
||||
resolve: (value: Playlist) => void;
|
||||
sender: WebContents;
|
||||
worker: Worker;
|
||||
};
|
||||
|
||||
const activePlaylistRefreshes = new Map<string, ActivePlaylistRefresh>();
|
||||
|
||||
/**
|
||||
* Fetches and parses a playlist from a URL
|
||||
* @param url - The URL to fetch the playlist from
|
||||
@@ -73,6 +97,46 @@ async function fetchPlaylistFromFile(
|
||||
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,
|
||||
appPath: app.getAppPath(),
|
||||
});
|
||||
|
||||
return new Worker(pathToFileURL(bootstrap.workerPath), {
|
||||
workerData: {
|
||||
nativeModuleSearchPaths: bootstrap.nativeModuleSearchPaths,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
function emitPlaylistRefreshEvent(
|
||||
sender: WebContents,
|
||||
event: PlaylistRefreshEvent
|
||||
): void {
|
||||
if (sender.isDestroyed()) {
|
||||
return;
|
||||
}
|
||||
|
||||
sender.send(PLAYLIST_REFRESH_EVENT, event);
|
||||
}
|
||||
|
||||
function createPlaylistRefreshError(error: {
|
||||
message: string;
|
||||
name?: string;
|
||||
stack?: string;
|
||||
}): Error {
|
||||
const workerError = new Error(error.message);
|
||||
workerError.name = error.name || 'PlaylistRefreshWorkerError';
|
||||
workerError.stack = error.stack || workerError.stack;
|
||||
return workerError;
|
||||
}
|
||||
|
||||
ipcMain.handle('fetch-playlist-by-url', async (event, url, title?: string) => {
|
||||
try {
|
||||
return await fetchPlaylistFromUrl(url, title);
|
||||
@@ -179,6 +243,93 @@ ipcMain.handle(AUTO_UPDATE_PLAYLISTS, async (event, playlists) => {
|
||||
return updatedPlaylists;
|
||||
});
|
||||
|
||||
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);
|
||||
};
|
||||
|
||||
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('error', async (error) => {
|
||||
await cleanup();
|
||||
reject(error);
|
||||
});
|
||||
|
||||
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}`
|
||||
)
|
||||
);
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
ipcMain.handle(
|
||||
PLAYLIST_CANCEL_REFRESH,
|
||||
async (_event, operationId: string): Promise<{ success: boolean }> => {
|
||||
const activeRefresh = activePlaylistRefreshes.get(operationId);
|
||||
if (!activeRefresh) {
|
||||
return { success: false };
|
||||
}
|
||||
|
||||
activeRefresh.worker.postMessage({
|
||||
type: 'cancel',
|
||||
operationId,
|
||||
});
|
||||
|
||||
return { success: true };
|
||||
}
|
||||
);
|
||||
|
||||
ipcMain.handle('save-file-dialog', async (event, defaultPath, filters) => {
|
||||
try {
|
||||
const { canceled, filePath } = await dialog.showSaveDialog({
|
||||
|
||||
@@ -0,0 +1,168 @@
|
||||
import { XTREAM_CANCEL_SESSION } from 'shared-interfaces';
|
||||
|
||||
const registeredHandlers = new Map<string, (...args: unknown[]) => unknown>();
|
||||
const axiosMock = Object.assign(jest.fn(), {
|
||||
isAxiosError: jest.fn(),
|
||||
});
|
||||
|
||||
function createDeferred<T>() {
|
||||
let resolve!: (value: T) => void;
|
||||
let reject!: (reason?: unknown) => void;
|
||||
const promise = new Promise<T>((res, rej) => {
|
||||
resolve = res;
|
||||
reject = rej;
|
||||
});
|
||||
|
||||
return { promise, resolve, reject };
|
||||
}
|
||||
|
||||
jest.mock('electron', () => ({
|
||||
ipcMain: {
|
||||
handle: jest.fn((channel: string, handler: (...args: unknown[]) => unknown) => {
|
||||
registeredHandlers.set(channel, handler);
|
||||
}),
|
||||
},
|
||||
}));
|
||||
|
||||
jest.mock('axios', () => ({
|
||||
__esModule: true,
|
||||
default: axiosMock,
|
||||
}));
|
||||
|
||||
jest.mock('./portal-debug.events', () => ({
|
||||
emitPortalDebugEvent: jest.fn(),
|
||||
}));
|
||||
|
||||
describe('XtreamEvents session cancellation', () => {
|
||||
let consoleErrorSpy: jest.SpyInstance;
|
||||
|
||||
beforeEach(async () => {
|
||||
jest.resetModules();
|
||||
registeredHandlers.clear();
|
||||
axiosMock.mockReset();
|
||||
axiosMock.isAxiosError.mockReset();
|
||||
consoleErrorSpy = jest.spyOn(console, 'error').mockImplementation();
|
||||
|
||||
await import('./xtream.events');
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
consoleErrorSpy.mockRestore();
|
||||
});
|
||||
|
||||
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 cancelError = Object.assign(new Error('cancelled'), {
|
||||
code: 'ERR_CANCELED',
|
||||
});
|
||||
let abortSignal: AbortSignal | undefined;
|
||||
|
||||
expect(requestHandler).toBeDefined();
|
||||
expect(cancelHandler).toBeDefined();
|
||||
|
||||
axiosMock.mockImplementation((config: { signal?: AbortSignal }) => {
|
||||
abortSignal = config.signal;
|
||||
return pendingRequest.promise;
|
||||
});
|
||||
axiosMock.isAxiosError.mockImplementation(
|
||||
(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>;
|
||||
|
||||
expect(abortSignal?.aborted).toBe(false);
|
||||
|
||||
const cancelResult = (await cancelHandler?.(
|
||||
{},
|
||||
'session-1'
|
||||
)) as { success: boolean; cancelled: number };
|
||||
|
||||
expect(cancelResult).toEqual({ success: true, cancelled: 1 });
|
||||
expect(abortSignal?.aborted).toBe(true);
|
||||
|
||||
pendingRequest.reject(cancelError);
|
||||
|
||||
await expect(requestPromise).rejects.toMatchObject({
|
||||
name: 'AbortError',
|
||||
status: 499,
|
||||
});
|
||||
});
|
||||
|
||||
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 cancelError = Object.assign(new Error('cancelled'), {
|
||||
code: 'ERR_CANCELED',
|
||||
});
|
||||
const abortSignals: AbortSignal[] = [];
|
||||
|
||||
expect(requestHandler).toBeDefined();
|
||||
expect(cancelHandler).toBeDefined();
|
||||
|
||||
axiosMock
|
||||
.mockImplementationOnce((config: { signal?: AbortSignal }) => {
|
||||
if (config.signal) {
|
||||
abortSignals.push(config.signal);
|
||||
}
|
||||
return firstRequest.promise;
|
||||
})
|
||||
.mockImplementationOnce((config: { signal?: AbortSignal }) => {
|
||||
if (config.signal) {
|
||||
abortSignals.push(config.signal);
|
||||
}
|
||||
return secondRequest.promise;
|
||||
});
|
||||
axiosMock.isAxiosError.mockImplementation(
|
||||
(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?.(
|
||||
{},
|
||||
'session-2'
|
||||
)) as { success: boolean; cancelled: number };
|
||||
|
||||
expect(cancelResult).toEqual({ success: true, cancelled: 2 });
|
||||
expect(abortSignals).toHaveLength(2);
|
||||
expect(abortSignals.every((signal) => signal.aborted)).toBe(true);
|
||||
|
||||
firstRequest.reject(cancelError);
|
||||
secondRequest.reject(cancelError);
|
||||
|
||||
await expect(firstPromise).rejects.toMatchObject({ name: 'AbortError' });
|
||||
await expect(secondPromise).rejects.toMatchObject({ name: 'AbortError' });
|
||||
});
|
||||
});
|
||||
@@ -5,7 +5,7 @@
|
||||
|
||||
import axios, { AxiosRequestConfig } from 'axios';
|
||||
import { ipcMain } from 'electron';
|
||||
import { PortalDebugEvent } from 'shared-interfaces';
|
||||
import { PortalDebugEvent, XTREAM_CANCEL_SESSION } from 'shared-interfaces';
|
||||
import { emitPortalDebugEvent } from './portal-debug.events';
|
||||
|
||||
export default class XtreamEvents {
|
||||
@@ -62,12 +62,14 @@ ipcMain.handle(
|
||||
url: string;
|
||||
params: Record<string, string>;
|
||||
requestId?: string;
|
||||
sessionId?: string;
|
||||
suppressErrorLog?: boolean;
|
||||
}
|
||||
) => {
|
||||
const startedAt = Date.now();
|
||||
let activeRequestKey: string | null = null;
|
||||
try {
|
||||
const { url, params, requestId } = payload;
|
||||
const { url, params, requestId, sessionId } = payload;
|
||||
|
||||
// Build URL with query parameters
|
||||
// Xtream API endpoint is always at /player_api.php
|
||||
@@ -76,6 +78,15 @@ ipcMain.handle(
|
||||
apiUrl.searchParams.append(key, value);
|
||||
});
|
||||
|
||||
const controller = new AbortController();
|
||||
if (requestId || sessionId) {
|
||||
activeRequestKey = requestId ?? crypto.randomUUID();
|
||||
activeXtreamRequests.set(activeRequestKey, {
|
||||
controller,
|
||||
sessionId,
|
||||
});
|
||||
}
|
||||
|
||||
// Configure axios request
|
||||
const config: AxiosRequestConfig = {
|
||||
method: 'GET',
|
||||
@@ -87,6 +98,7 @@ ipcMain.handle(
|
||||
},
|
||||
timeout: 30000, // 30 seconds timeout for Xtream API
|
||||
validateStatus: (status) => status < 500, // Don't throw on 4xx errors
|
||||
signal: controller.signal,
|
||||
};
|
||||
|
||||
const response = await axios(config);
|
||||
@@ -166,6 +178,14 @@ ipcMain.handle(
|
||||
|
||||
// Format error response
|
||||
if (axios.isAxiosError(error)) {
|
||||
if (error.code === 'ERR_CANCELED') {
|
||||
throw {
|
||||
type: 'ERROR',
|
||||
name: 'AbortError',
|
||||
message: 'Xtream request cancelled',
|
||||
status: 499,
|
||||
};
|
||||
}
|
||||
const errorResponse = {
|
||||
type: 'ERROR',
|
||||
message:
|
||||
@@ -188,6 +208,40 @@ ipcMain.handle(
|
||||
status: 500,
|
||||
};
|
||||
}
|
||||
} finally {
|
||||
if (activeRequestKey) {
|
||||
activeXtreamRequests.delete(activeRequestKey);
|
||||
}
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
ipcMain.handle(
|
||||
XTREAM_CANCEL_SESSION,
|
||||
async (_event, sessionId: string): Promise<{ success: boolean; cancelled: number }> => {
|
||||
if (!sessionId) {
|
||||
return { success: false, cancelled: 0 };
|
||||
}
|
||||
|
||||
let cancelled = 0;
|
||||
for (const activeRequest of activeXtreamRequests.values()) {
|
||||
if (activeRequest.sessionId !== sessionId) {
|
||||
continue;
|
||||
}
|
||||
|
||||
activeRequest.controller.abort();
|
||||
cancelled += 1;
|
||||
}
|
||||
|
||||
return {
|
||||
success: cancelled > 0,
|
||||
cancelled,
|
||||
};
|
||||
}
|
||||
);
|
||||
type ActiveXtreamRequest = {
|
||||
controller: AbortController;
|
||||
sessionId?: string;
|
||||
};
|
||||
|
||||
const activeXtreamRequests = new Map<string, ActiveXtreamRequest>();
|
||||
@@ -8,6 +8,7 @@ export const DB_WORKER_OPERATIONS = [
|
||||
'DB_GET_CONTENT',
|
||||
'DB_GET_GLOBAL_RECENTLY_ADDED',
|
||||
'DB_SAVE_CONTENT',
|
||||
'DB_CLEAR_XTREAM_IMPORT_CACHE',
|
||||
'DB_GET_CONTENT_BY_XTREAM_ID',
|
||||
'DB_SEARCH_CONTENT',
|
||||
'DB_GLOBAL_SEARCH',
|
||||
|
||||
@@ -30,6 +30,7 @@ import {
|
||||
reorderGlobalFavorites,
|
||||
} from '../database/operations/favorites.operations';
|
||||
import {
|
||||
clearXtreamImportCache,
|
||||
getContent,
|
||||
getContentByXtreamId,
|
||||
getGlobalRecentlyAdded,
|
||||
@@ -387,6 +388,15 @@ async function executeRequest(message: DbWorkerRequestMessage) {
|
||||
);
|
||||
}
|
||||
|
||||
case 'DB_CLEAR_XTREAM_IMPORT_CACHE': {
|
||||
const payload = message.payload as {
|
||||
playlistId: string;
|
||||
type: 'live' | 'movie' | 'series';
|
||||
};
|
||||
|
||||
return clearXtreamImportCache(db, payload.playlistId, payload.type);
|
||||
}
|
||||
|
||||
case 'DB_GET_CONTENT_BY_XTREAM_ID': {
|
||||
const payload = message.payload as {
|
||||
xtreamId: number;
|
||||
|
||||
@@ -0,0 +1,188 @@
|
||||
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 'm3u-utils';
|
||||
import type {
|
||||
Playlist,
|
||||
PlaylistRefreshEvent,
|
||||
PlaylistRefreshPayload,
|
||||
} from 'shared-interfaces';
|
||||
import type {
|
||||
PlaylistRefreshWorkerIncomingMessage,
|
||||
PlaylistRefreshWorkerMessage,
|
||||
} from './playlist-refresh.worker.types';
|
||||
|
||||
const https = require('https');
|
||||
|
||||
type ActiveRefreshState = {
|
||||
cancelled: boolean;
|
||||
controller: AbortController;
|
||||
};
|
||||
|
||||
const activeRefreshes = new Map<string, ActiveRefreshState>();
|
||||
|
||||
if (!parentPort) {
|
||||
throw new Error('Playlist refresh worker must be started with a parent port');
|
||||
}
|
||||
|
||||
function postMessage(message: PlaylistRefreshWorkerMessage<Playlist>): void {
|
||||
parentPort?.postMessage(message);
|
||||
}
|
||||
|
||||
function serializeError(error: unknown) {
|
||||
if (error instanceof Error) {
|
||||
return {
|
||||
name: error.name,
|
||||
message: error.message,
|
||||
stack: error.stack,
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
message: String(error),
|
||||
};
|
||||
}
|
||||
|
||||
function createAbortError(operationId: string): Error {
|
||||
const error = new Error(`Playlist refresh "${operationId}" was cancelled`);
|
||||
error.name = 'AbortError';
|
||||
return error;
|
||||
}
|
||||
|
||||
function emitEvent(
|
||||
payload: PlaylistRefreshPayload,
|
||||
partial: Omit<PlaylistRefreshEvent, 'operationId' | 'playlistId'>
|
||||
): void {
|
||||
postMessage({
|
||||
type: 'event',
|
||||
event: {
|
||||
operationId: payload.operationId,
|
||||
playlistId: payload.playlistId,
|
||||
...partial,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
function checkpoint(payload: PlaylistRefreshPayload): void {
|
||||
const active = activeRefreshes.get(payload.operationId);
|
||||
if (active?.cancelled) {
|
||||
throw createAbortError(payload.operationId);
|
||||
}
|
||||
}
|
||||
|
||||
async function fetchPlaylistFromUrl(
|
||||
payload: PlaylistRefreshPayload,
|
||||
controller: AbortController
|
||||
): Promise<Playlist> {
|
||||
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,
|
||||
});
|
||||
|
||||
checkpoint(payload);
|
||||
emitEvent(payload, { status: 'progress', phase: 'parsing' });
|
||||
const parsedPlaylist = parse(result.data);
|
||||
checkpoint(payload);
|
||||
|
||||
const extractedName =
|
||||
payload.url && payload.url.length > 1
|
||||
? getFilenameFromUrl(payload.url)
|
||||
: '';
|
||||
const playlistName =
|
||||
!extractedName || extractedName === 'Untitled playlist'
|
||||
? 'Imported from URL'
|
||||
: extractedName;
|
||||
|
||||
return createPlaylistObject(
|
||||
payload.title || playlistName,
|
||||
parsedPlaylist,
|
||||
payload.url,
|
||||
'URL'
|
||||
);
|
||||
}
|
||||
|
||||
async function fetchPlaylistFromFile(
|
||||
payload: PlaylistRefreshPayload
|
||||
): Promise<Playlist> {
|
||||
emitEvent(payload, { status: 'started', phase: 'reading-file' });
|
||||
checkpoint(payload);
|
||||
const fileContent = await readFile(payload.filePath!, 'utf-8');
|
||||
checkpoint(payload);
|
||||
|
||||
emitEvent(payload, { status: 'progress', phase: 'parsing' });
|
||||
const parsedPlaylist = parse(fileContent);
|
||||
checkpoint(payload);
|
||||
|
||||
return createPlaylistObject(
|
||||
payload.title,
|
||||
parsedPlaylist,
|
||||
payload.filePath,
|
||||
'FILE'
|
||||
);
|
||||
}
|
||||
|
||||
async function executeRefresh(payload: PlaylistRefreshPayload): Promise<Playlist> {
|
||||
const controller = new AbortController();
|
||||
activeRefreshes.set(payload.operationId, {
|
||||
cancelled: false,
|
||||
controller,
|
||||
});
|
||||
|
||||
try {
|
||||
const playlist = payload.url
|
||||
? await fetchPlaylistFromUrl(payload, controller)
|
||||
: await fetchPlaylistFromFile(payload);
|
||||
|
||||
checkpoint(payload);
|
||||
return playlist;
|
||||
} finally {
|
||||
activeRefreshes.delete(payload.operationId);
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
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: 'ready' });
|
||||
@@ -0,0 +1,40 @@
|
||||
import type { PlaylistRefreshEvent, PlaylistRefreshPayload } from 'shared-interfaces';
|
||||
|
||||
export interface PlaylistRefreshWorkerRequestMessage {
|
||||
type: 'request';
|
||||
payload: PlaylistRefreshPayload;
|
||||
}
|
||||
|
||||
export interface PlaylistRefreshWorkerCancelMessage {
|
||||
type: 'cancel';
|
||||
operationId: string;
|
||||
}
|
||||
|
||||
export interface PlaylistRefreshWorkerReadyMessage {
|
||||
type: 'ready';
|
||||
}
|
||||
|
||||
export interface PlaylistRefreshWorkerEventMessage {
|
||||
type: 'event';
|
||||
event: PlaylistRefreshEvent;
|
||||
}
|
||||
|
||||
export interface PlaylistRefreshWorkerResponseMessage<TResult = unknown> {
|
||||
type: 'response';
|
||||
success: boolean;
|
||||
result?: TResult;
|
||||
error?: {
|
||||
name?: string;
|
||||
message: string;
|
||||
stack?: string;
|
||||
};
|
||||
}
|
||||
|
||||
export type PlaylistRefreshWorkerIncomingMessage =
|
||||
| PlaylistRefreshWorkerRequestMessage
|
||||
| PlaylistRefreshWorkerCancelMessage;
|
||||
|
||||
export type PlaylistRefreshWorkerMessage<TResult = unknown> =
|
||||
| PlaylistRefreshWorkerReadyMessage
|
||||
| PlaylistRefreshWorkerEventMessage
|
||||
| PlaylistRefreshWorkerResponseMessage<TResult>;
|
||||
@@ -50,6 +50,7 @@ export class ElectronService extends DataService {
|
||||
XtreamCodeActions.GetLiveCategories,
|
||||
XtreamCodeActions.GetVodCategories,
|
||||
XtreamCodeActions.GetSeriesCategories,
|
||||
XtreamCodeActions.GetShortEpg,
|
||||
]);
|
||||
|
||||
constructor() {
|
||||
@@ -432,6 +433,8 @@ export class ElectronService extends DataService {
|
||||
url: string;
|
||||
params: Record<string, string>;
|
||||
requestId?: string;
|
||||
sessionId?: string;
|
||||
suppressErrorLog?: boolean;
|
||||
}) {
|
||||
const context = createPortalDebugRequestContext({
|
||||
provider: 'xtream',
|
||||
@@ -456,9 +459,9 @@ export class ElectronService extends DataService {
|
||||
return result;
|
||||
} catch (error: unknown) {
|
||||
const action = payload.params?.action;
|
||||
const isSilentAction = action
|
||||
? this.silentXtreamActions.has(action)
|
||||
: false;
|
||||
const isSilentAction =
|
||||
payload.suppressErrorLog === true ||
|
||||
(action ? this.silentXtreamActions.has(action) : false);
|
||||
const normalizedMessage = this.getReadableXtreamErrorMessage(error);
|
||||
const errorInfo = this.getErrorDetails(error);
|
||||
|
||||
|
||||
@@ -62,6 +62,7 @@ export class PwaService extends DataService {
|
||||
XtreamCodeActions.GetLiveCategories,
|
||||
XtreamCodeActions.GetVodCategories,
|
||||
XtreamCodeActions.GetSeriesCategories,
|
||||
XtreamCodeActions.GetShortEpg,
|
||||
]);
|
||||
|
||||
/** Proxy URL to avoid CORS issues */
|
||||
@@ -244,6 +245,9 @@ export class PwaService extends DataService {
|
||||
url: string;
|
||||
params: Record<string, string>;
|
||||
macAddress?: string;
|
||||
requestId?: string;
|
||||
sessionId?: string;
|
||||
suppressErrorLog?: boolean;
|
||||
}) {
|
||||
const headers = payload.macAddress
|
||||
? {
|
||||
@@ -292,7 +296,9 @@ export class PwaService extends DataService {
|
||||
|
||||
if (!response.payload) {
|
||||
const action = payload.params.action;
|
||||
const isSilentAction = this.silentXtreamActions.has(action);
|
||||
const isSilentAction =
|
||||
payload.suppressErrorLog === true ||
|
||||
this.silentXtreamActions.has(action);
|
||||
const normalizedMessage =
|
||||
this.getReadableXtreamErrorMessage(response);
|
||||
logPortalDebugEvent(
|
||||
@@ -332,7 +338,9 @@ export class PwaService extends DataService {
|
||||
} catch (error: unknown) {
|
||||
logPortalDebugEvent(createPortalDebugErrorEvent(context, error));
|
||||
const action = payload.params.action;
|
||||
const isSilentAction = this.silentXtreamActions.has(action);
|
||||
const isSilentAction =
|
||||
payload.suppressErrorLog === true ||
|
||||
this.silentXtreamActions.has(action);
|
||||
const normalizedMessage = this.getReadableXtreamErrorMessage(error);
|
||||
const errorInfo = this.getErrorDetails(error);
|
||||
|
||||
|
||||
Vendored
+19
@@ -10,6 +10,8 @@ import {
|
||||
ExternalPlayerSession,
|
||||
PlaybackPositionData,
|
||||
Playlist,
|
||||
PlaylistRefreshEvent,
|
||||
PlaylistRefreshPayload,
|
||||
PortalDebugEvent,
|
||||
} from 'shared-interfaces';
|
||||
import {
|
||||
@@ -31,6 +33,9 @@ declare global {
|
||||
onPortalDebugEvent?: (
|
||||
callback: (data: PortalDebugEvent) => void
|
||||
) => () => void;
|
||||
onPlaylistRefreshEvent?: (
|
||||
callback: (data: PlaylistRefreshEvent) => void
|
||||
) => () => void;
|
||||
getAppVersion: () => Promise<string>;
|
||||
platform: string;
|
||||
fetchPlaylistByUrl: (
|
||||
@@ -119,8 +124,18 @@ declare global {
|
||||
url: string;
|
||||
params: Record<string, string>;
|
||||
requestId?: string;
|
||||
sessionId?: string;
|
||||
suppressErrorLog?: boolean;
|
||||
}) => Promise<{ payload: unknown; action: string }>;
|
||||
xtreamCancelSession: (
|
||||
sessionId: string
|
||||
) => Promise<{ success: boolean; cancelled: number }>;
|
||||
refreshPlaylist: (
|
||||
payload: PlaylistRefreshPayload
|
||||
) => Promise<Playlist>;
|
||||
cancelPlaylistRefresh: (
|
||||
operationId: string
|
||||
) => Promise<{ success: boolean }>;
|
||||
// Database operations
|
||||
dbCreatePlaylist: (
|
||||
playlist: Playlist
|
||||
@@ -202,6 +217,10 @@ declare global {
|
||||
type: string,
|
||||
operationId?: string
|
||||
) => Promise<{ success: boolean; count: number }>;
|
||||
dbClearXtreamImportCache: (
|
||||
playlistId: string,
|
||||
type: 'live' | 'movie' | 'series'
|
||||
) => Promise<{ success: boolean }>;
|
||||
dbSearchContent: (
|
||||
playlistId: string,
|
||||
searchTerm: string,
|
||||
|
||||
Reference in new issue
Block a user