From 4e0afe8b209f8abe2445c1d92a7979f411e4f991 Mon Sep 17 00:00:00 2001 From: 4gray Date: Sat, 4 Apr 2026 13:39:45 +0200 Subject: [PATCH] feat: implement playlist refresh functionality and cache management --- apps/electron-backend/build-worker.js | 15 ++ .../src/app/api/main.preload.ts | 26 ++- .../database/operations/content.operations.ts | 45 +++++ .../operations/playlist.operations.spec.ts | 36 ++++ .../operations/playlist.operations.ts | 68 +++---- .../src/app/events/database/content.events.ts | 8 + .../src/app/events/playlist.events.ts | 155 ++++++++++++++- .../src/app/events/xtream.events.spec.ts | 168 ++++++++++++++++ .../src/app/events/xtream.events.ts | 58 +++++- .../src/app/workers/database-worker.types.ts | 1 + .../src/app/workers/database.worker.ts | 10 + .../app/workers/playlist-refresh.worker.ts | 188 ++++++++++++++++++ .../workers/playlist-refresh.worker.types.ts | 40 ++++ apps/web/src/app/services/electron.service.ts | 9 +- apps/web/src/app/services/pwa.service.ts | 12 +- apps/web/src/typings.d.ts | 19 ++ 16 files changed, 813 insertions(+), 45 deletions(-) create mode 100644 apps/electron-backend/src/app/database/operations/playlist.operations.spec.ts create mode 100644 apps/electron-backend/src/app/events/xtream.events.spec.ts create mode 100644 apps/electron-backend/src/app/workers/playlist-refresh.worker.ts create mode 100644 apps/electron-backend/src/app/workers/playlist-refresh.worker.types.ts diff --git a/apps/electron-backend/build-worker.js b/apps/electron-backend/build-worker.js index c35b03aec..7d7d7c48f 100644 --- a/apps/electron-backend/build-worker.js +++ b/apps/electron-backend/build-worker.js @@ -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' diff --git a/apps/electron-backend/src/app/api/main.preload.ts b/apps/electron-backend/src/app/api/main.preload.ts index a552a2606..d99aed546 100644 --- a/apps/electron-backend/src/app/api/main.preload.ts +++ b/apps/electron-backend/src/app/api/main.preload.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; 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, diff --git a/apps/electron-backend/src/app/database/operations/content.operations.ts b/apps/electron-backend/src/app/database/operations/content.operations.ts index e513f938f..efc8e59dc 100644 --- a/apps/electron-backend/src/app/database/operations/content.operations.ts +++ b/apps/electron-backend/src/app/database/operations/content.operations.ts @@ -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, diff --git a/apps/electron-backend/src/app/database/operations/playlist.operations.spec.ts b/apps/electron-backend/src/app/database/operations/playlist.operations.spec.ts new file mode 100644 index 000000000..6b3416117 --- /dev/null +++ b/apps/electron-backend/src/app/database/operations/playlist.operations.spec.ts @@ -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(), + }) + ); + }); +}); diff --git a/apps/electron-backend/src/app/database/operations/playlist.operations.ts b/apps/electron-backend/src/app/database/operations/playlist.operations.ts index 96e66b627..bd9e4cdef 100644 --- a/apps/electron-backend/src/app/database/operations/playlist.operations.ts +++ b/apps/electron-backend/src/app/database/operations/playlist.operations.ts @@ -139,26 +139,12 @@ function buildPlaylistRow( }; } -function parseAppPlaylist(row: schema.Playlist): Record { +export function parseAppPlaylist(row: schema.Playlist): Record { const payload = parseJsonValue | 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(row.favorites, []); const recentlyViewed = parseJsonValue(row.recentlyViewed, []); const importDate = @@ -166,30 +152,42 @@ function parseAppPlaylist(row: schema.Playlist): Record { 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), }; } diff --git a/apps/electron-backend/src/app/events/database/content.events.ts b/apps/electron-backend/src/app/events/database/content.events.ts index 7e0a2bc56..655c8b65d 100644 --- a/apps/electron-backend/src/app/events/database/content.events.ts +++ b/apps/electron-backend/src/app/events/database/content.events.ts @@ -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) => ({ diff --git a/apps/electron-backend/src/app/events/playlist.events.ts b/apps/electron-backend/src/app/events/playlist.events.ts index 171ef5a84..c0035dfaf 100644 --- a/apps/electron-backend/src/app/events/playlist.events.ts +++ b/apps/electron-backend/src/app/events/playlist.events.ts @@ -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(); + /** * 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((resolve, reject) => { + const cleanup = async (): Promise => { + 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) => { + 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; + 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({ diff --git a/apps/electron-backend/src/app/events/xtream.events.spec.ts b/apps/electron-backend/src/app/events/xtream.events.spec.ts new file mode 100644 index 000000000..743338f48 --- /dev/null +++ b/apps/electron-backend/src/app/events/xtream.events.spec.ts @@ -0,0 +1,168 @@ +import { XTREAM_CANCEL_SESSION } from 'shared-interfaces'; + +const registeredHandlers = new Map unknown>(); +const axiosMock = Object.assign(jest.fn(), { + isAxiosError: jest.fn(), +}); + +function createDeferred() { + let resolve!: (value: T) => void; + let reject!: (reason?: unknown) => void; + const promise = new Promise((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; + + 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; + const secondPromise = requestHandler?.({}, { + url: 'http://localhost:3211', + params: { + action: 'get_vod_streams', + password: 'secret', + username: 'user1', + }, + sessionId: 'session-2', + suppressErrorLog: true, + }) as Promise; + + 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' }); + }); +}); diff --git a/apps/electron-backend/src/app/events/xtream.events.ts b/apps/electron-backend/src/app/events/xtream.events.ts index 9bf4026d5..73f21485f 100644 --- a/apps/electron-backend/src/app/events/xtream.events.ts +++ b/apps/electron-backend/src/app/events/xtream.events.ts @@ -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; 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(); diff --git a/apps/electron-backend/src/app/workers/database-worker.types.ts b/apps/electron-backend/src/app/workers/database-worker.types.ts index 797a4bd6c..5967e5495 100644 --- a/apps/electron-backend/src/app/workers/database-worker.types.ts +++ b/apps/electron-backend/src/app/workers/database-worker.types.ts @@ -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', diff --git a/apps/electron-backend/src/app/workers/database.worker.ts b/apps/electron-backend/src/app/workers/database.worker.ts index e8e48a051..98240ea68 100644 --- a/apps/electron-backend/src/app/workers/database.worker.ts +++ b/apps/electron-backend/src/app/workers/database.worker.ts @@ -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; diff --git a/apps/electron-backend/src/app/workers/playlist-refresh.worker.ts b/apps/electron-backend/src/app/workers/playlist-refresh.worker.ts new file mode 100644 index 000000000..959ce6bed --- /dev/null +++ b/apps/electron-backend/src/app/workers/playlist-refresh.worker.ts @@ -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(); + +if (!parentPort) { + throw new Error('Playlist refresh worker must be started with a parent port'); +} + +function postMessage(message: PlaylistRefreshWorkerMessage): 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 +): 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 { + 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 { + 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 { + 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' }); diff --git a/apps/electron-backend/src/app/workers/playlist-refresh.worker.types.ts b/apps/electron-backend/src/app/workers/playlist-refresh.worker.types.ts new file mode 100644 index 000000000..ed6607311 --- /dev/null +++ b/apps/electron-backend/src/app/workers/playlist-refresh.worker.types.ts @@ -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 { + type: 'response'; + success: boolean; + result?: TResult; + error?: { + name?: string; + message: string; + stack?: string; + }; +} + +export type PlaylistRefreshWorkerIncomingMessage = + | PlaylistRefreshWorkerRequestMessage + | PlaylistRefreshWorkerCancelMessage; + +export type PlaylistRefreshWorkerMessage = + | PlaylistRefreshWorkerReadyMessage + | PlaylistRefreshWorkerEventMessage + | PlaylistRefreshWorkerResponseMessage; diff --git a/apps/web/src/app/services/electron.service.ts b/apps/web/src/app/services/electron.service.ts index cadc61797..b5b23de6e 100644 --- a/apps/web/src/app/services/electron.service.ts +++ b/apps/web/src/app/services/electron.service.ts @@ -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; 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); diff --git a/apps/web/src/app/services/pwa.service.ts b/apps/web/src/app/services/pwa.service.ts index 3d090343e..b3d20c234 100644 --- a/apps/web/src/app/services/pwa.service.ts +++ b/apps/web/src/app/services/pwa.service.ts @@ -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; 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); diff --git a/apps/web/src/typings.d.ts b/apps/web/src/typings.d.ts index c0bd3b027..5bca10657 100644 --- a/apps/web/src/typings.d.ts +++ b/apps/web/src/typings.d.ts @@ -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; platform: string; fetchPlaylistByUrl: ( @@ -119,8 +124,18 @@ declare global { url: string; params: Record; requestId?: string; + sessionId?: string; suppressErrorLog?: boolean; }) => Promise<{ payload: unknown; action: string }>; + xtreamCancelSession: ( + sessionId: string + ) => Promise<{ success: boolean; cancelled: number }>; + refreshPlaylist: ( + payload: PlaylistRefreshPayload + ) => Promise; + 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,