Files
iptvnator/apps/electron-backend/src/app/workers/database.worker.ts
T
1d9a563d1a test(performance): count startup phases and SQL statements for the J1 launch journey (#1715)
* test(performance): count startup phases and SQL statements for the J1 launch journey

Implements plan item A2. With IPTVNATOR_PERF_CAPTURE=1 the main process
keeps named counters and registers a main-only performance:read-counters
IPC handler; without the flag nothing is counted and the handler does not
exist.

- debug-trace.ts owns the registry; traceStartupPhase replaces the
  trace('startup', ...) sites and counts main.startupPhases.
- The database worker counts executed statements through better-sqlite3's
  Statement prototype (the verbose callback expands every statement and
  made bulk inserts 2-4x slower) and posts the count over its message
  port, flushed before every other worker message. The main-thread shared
  connection is counted through a new connection observer in the shared
  database library.
- The first main window freezes main.modulesRegisteredBeforeWindow at
  creation and main.sqlStatementsBeforeReadyToShow at ready-to-show.
- The journey gate drops the ready-to-show that Electron emits for the
  about:blank detour, so the app sees the real document's first paint,
  and taps the counters handler; the J1 record reads both counters after
  the renderer probe completes.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* test(database): require one SQL statement per exec during initialization

The performance capture counts one exec call as one statement, because
SQL cannot be split reliably in the counter (trigger bodies contain
semicolons). The historical-upgrade driver now wraps exec on every
connection initDatabase opens and fails on a batch, so that counting
assumption holds for the fresh profile and all historical schemas.
Documents the definition in the counter and the architecture docs.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* test(performance): count SQL statements only for the launch journey

Codex review: the M3U import, refresh-cancellation and Xtream benchmarks
also run with IPTVNATOR_PERF_CAPTURE=1, so the statement hook wrapped
every row of their bulk inserts and changed what they measure.

SQL counting now also needs IPTVNATOR_PERF_COUNT_SQL=1, which only the
launch journey sets; a harness test fails if another source sets it.
Startup phases, the window snapshot and the read handler stay on the
capture flag. Without SQL counting no ready-to-show listener is attached,
so a zero is never reported for statements nobody counted.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

---------

Co-authored-by: 4gray <fourgray@proton.me>
Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 21:01:26 +02:00

1364 lines
43 KiB
TypeScript

import { migrateAppPlaylists } from '../database/operations/playlist-migration.operations';
import {
closeWorkerDatabase,
flushWorkerSqlStatementCount,
getWorkerDatabase,
} from './database.worker-connection';
import { parentPort, workerData } from 'worker_threads';
import type {
ContentMetadataPatch,
VodSourcePin,
XtreamBackupFavoriteItem,
XtreamBackupRecentlyViewedItem,
} from '@iptvnator/shared/interfaces';
import type {
DbOperationEvent,
DbWorkerIncomingMessage,
DbWorkerMessage,
DbWorkerRequestMessage,
} from './database-worker.types';
import {
DB_OPERATION_NAMES,
DB_OPERATION_PHASES,
} from './database-worker.types';
import {
getAllCategories,
getCategories,
hasCategories,
saveCategories,
setCategoryLocks,
updateCategoryVisibility,
} from '../database/operations/category.operations';
import { setParentalLockActive } from '../database/parental-lock-state';
import {
addFavorite,
getAllGlobalFavorites,
getFavorites,
getGlobalFavorites,
isFavorite,
removeFavorite,
reorderGlobalFavorites,
} from '../database/operations/favorites.operations';
import {
clearXtreamImportCache,
getContent,
getContentByXtreamId,
getGlobalRecentlyAdded,
globalSearch,
hasContent,
saveContent,
searchContent,
} from '../database/operations/content.operations';
import { setContentMetadataIfMissing } from '../database/operations/content-metadata.operations';
import {
clearAllPlaybackPositions,
clearPlaybackPosition,
clearPlaybackPositionsBatch,
getAllPlaybackPositions,
getPlaybackPosition,
getRecentPlaybackPositions,
getSeriesPlaybackPositions,
savePlaybackPosition,
savePlaybackPositionsBatch,
} from '../database/operations/playback-position.operations';
import {
createPlaylist,
deleteAllPlaylists,
deletePlaylist,
getAppPlaylist,
getAppPlaylistFavoriteChannels,
getAppPlaylistMetas,
getAppPlaylists,
getAppState,
getPlaylist,
setAppState,
type AppPlaylistGetPhaseCapture,
type AppPlaylistUpsertPhaseCapture,
updatePlaylist,
upsertAppPlaylist,
upsertAppPlaylists,
} from '../database/operations/playlist.operations';
import { setPlaylistServerTimezone } from '../database/operations/playlist-server-timezone.operations';
import {
addRecentItem,
clearPlaylistRecentItems,
clearRecentlyViewed,
getRecentItems,
getRecentlyViewed,
removeRecentItem,
removeRecentItemsBatch,
} from '../database/operations/recently-viewed.operations';
import { matchTitles } from '../database/operations/title-match.operations';
import {
findTitleSources,
type FindTitleSourcesRequest,
} from '../database/operations/title-sources.operations';
import {
clearTmdbMetadata,
getTmdbCacheStats,
getTmdbMetadata,
setTmdbMetadata,
} from '../database/operations/tmdb.operations';
import {
clearVodSourcePin,
getVodSourcePin,
clearVodSourcePinsForPlaylist,
replaceVodSourcePinsForPlaylist,
listVodSourcePinsForPlaylist,
setVodSourcePin,
} from '../database/operations/vod-source-pin.operations';
import {
deleteXtreamContent,
restoreXtreamUserData,
} from '../database/operations/xtream.operations';
import {
armWorkerPerformanceCapture,
executeWithWorkerPerformanceCapture,
registerDatabaseWorkerPerformanceCapture,
releaseDatabaseWorkerPerformanceCapture,
stampWorkerPerformanceResponsePostedEpoch,
startWorkerPerformanceCapture,
type WorkerPerformanceCapture,
} from './worker-performance-capture';
import {
captureWorkerPerformancePhase,
captureWorkerPerformancePhaseAsync,
createWorkerPerformancePhaseAdapter,
} from './worker-performance-phase';
import {
handleDatabaseWorkerPostGcHeapRequest,
isDatabaseWorkerPostGcHeapRequest,
} from './database-worker-post-gc-heap';
import { publishDatabaseWorkerCancelReceipt } from './database-worker-cancel-receipt';
import {
createProgressEventThrottle,
type ProgressEventThrottle,
} from './operation-progress-throttle';
const loggerLabel = '[DB Worker]';
// Seeded by the main process so a (re)started worker is locked before its
// first read; DB_SET_PARENTAL_LOCK_STATE updates it afterwards.
setParentalLockActive(
(workerData as { parentalLockActive?: unknown } | undefined)
?.parentalLockActive === true
);
const batchDelayMs = Number.parseInt(
process.env['IPTVNATOR_DB_WORKER_BATCH_DELAY_MS'] ?? '0',
10
);
/**
* Shortest gap between two progress events of one operation. Phase starts,
* phase completions and terminal events are never held back; see
* `operation-progress-throttle.ts`.
*/
const progressEventMinIntervalMs = 100;
type ActiveOperationState = {
cancelled: boolean;
performanceRequestId?: string;
};
type PreRegisteredActiveOperation = {
operationId: string;
state: ActiveOperationState;
};
type OperationController = {
control: {
checkpoint: () => Promise<void>;
onProgress: (progress: {
phase: string;
current?: number;
total?: number;
increment?: number;
}) => Promise<void>;
};
emitStarted: (event: Partial<DbOperationEvent> & { phase: string }) => void;
emitCompleted: (
event?: Partial<DbOperationEvent> & { phase?: string }
) => void;
emitCancelled: (
event?: Partial<DbOperationEvent> & { phase?: string }
) => void;
emitError: (
error: unknown,
event?: Partial<DbOperationEvent> & { phase?: string }
) => void;
cleanup: () => void;
};
const activeOperations = new Map<string, ActiveOperationState>();
const activePerformanceCaptures = new Set<WorkerPerformanceCapture>();
if (!parentPort) {
throw new Error('Database worker must be started with a parent port');
}
function createAbortError(message: string): Error {
const error = new Error(message);
error.name = 'AbortError';
return error;
}
function isAbortError(error: unknown): boolean {
return error instanceof Error && error.name === 'AbortError';
}
function serializeError(error: unknown) {
if (error instanceof Error) {
return {
name: error.name,
message: error.message,
stack: error.stack,
};
}
return {
message: String(error),
};
}
function postMessage(message: DbWorkerMessage): void {
flushWorkerSqlStatementCount();
parentPort?.postMessage(message);
}
function postEvent(requestId: string, event: DbOperationEvent): void {
postMessage({
type: 'event',
requestId,
event,
});
}
async function pauseBetweenBatches(): Promise<void> {
if (batchDelayMs <= 0) {
await new Promise<void>((resolve) => setImmediate(resolve));
return;
}
await new Promise((resolve) => setTimeout(resolve, batchDelayMs));
}
function createOperationController(
config: {
operationId?: string;
operation: string;
playlistId?: string;
requestId: string;
cancellable?: boolean;
},
preRegisteredState?: ActiveOperationState
): OperationController {
const { operationId, operation, playlistId, requestId } = config;
const cancellable = config.cancellable ?? true;
const activeState =
operationId && cancellable
? (preRegisteredState ?? { cancelled: false })
: null;
if (operationId && activeState) {
activeOperations.set(operationId, activeState);
}
let lastEvent: Partial<DbOperationEvent> = {};
const progressThrottle = createProgressEventThrottle({
minIntervalMs: progressEventMinIntervalMs,
});
const send = (
status: DbOperationEvent['status'],
event: Partial<DbOperationEvent> = {}
): void => {
const mergedEvent: DbOperationEvent = {
operationId,
operation,
playlistId,
status,
phase: event.phase ?? lastEvent.phase,
current: event.current ?? lastEvent.current,
total: event.total ?? lastEvent.total,
increment: event.increment,
error: event.error,
};
lastEvent = {
phase: mergedEvent.phase,
current: mergedEvent.current,
total: mergedEvent.total,
};
postEvent(requestId, mergedEvent);
};
const checkpoint = async (): Promise<void> => {
if (activeState?.cancelled) {
throw createAbortError(`Operation "${operation}" was cancelled`);
}
await pauseBetweenBatches();
if (activeState?.cancelled) {
throw createAbortError(`Operation "${operation}" was cancelled`);
}
};
const sendProgress = (
updates: ReturnType<ProgressEventThrottle['push']>
): void => {
for (const update of updates) {
send('progress', update);
}
};
return {
control: {
checkpoint,
onProgress: async (progress) => {
sendProgress(progressThrottle.push(progress));
},
},
emitStarted: (event) => {
send('started', event);
},
emitCompleted: (event) => {
sendProgress(progressThrottle.flush());
send('completed', event);
},
emitCancelled: (event) => {
sendProgress(progressThrottle.flush());
send('cancelled', event);
},
emitError: (error, event) => {
sendProgress(progressThrottle.flush());
send('error', {
...event,
error: error instanceof Error ? error.message : String(error),
});
},
cleanup: () => {
if (
operationId &&
activeOperations.get(operationId) === activeState
) {
activeOperations.delete(operationId);
}
},
};
}
async function executeTrackedOperation<TResult>(
config: Parameters<typeof createOperationController>[0],
handler: (controller: OperationController) => Promise<TResult>,
preRegisteredState?: ActiveOperationState
): Promise<TResult> {
const controller = createOperationController(config, preRegisteredState);
try {
return await handler(controller);
} catch (error) {
if (isAbortError(error)) {
controller.emitCancelled();
} else {
controller.emitError(error);
}
throw error;
} finally {
controller.cleanup();
}
}
function preRegisterCancellableOperation(
message: DbWorkerRequestMessage
): PreRegisteredActiveOperation | null {
switch (message.operation) {
case 'DB_SAVE_CONTENT':
case 'DB_DELETE_PLAYLIST':
case 'DB_DELETE_XTREAM_CONTENT':
case 'DB_RESTORE_XTREAM_USER_DATA':
break;
default:
return null;
}
if (
typeof message.payload !== 'object' ||
message.payload === null ||
Array.isArray(message.payload)
) {
return null;
}
const operationId = (message.payload as Record<string, unknown>)[
'operationId'
];
if (typeof operationId !== 'string' || operationId.length === 0) {
return null;
}
const state: ActiveOperationState = { cancelled: false };
activeOperations.set(operationId, state);
return { operationId, state };
}
function releasePreRegisteredOperation(
operation: PreRegisteredActiveOperation | null
): void {
if (
operation &&
activeOperations.get(operation.operationId) === operation.state
) {
activeOperations.delete(operation.operationId);
}
}
async function executeRequest(
message: DbWorkerRequestMessage,
preRegisteredState?: ActiveOperationState,
performanceCapture: WorkerPerformanceCapture | null = null
) {
const db = await getWorkerDatabase();
switch (message.operation) {
case 'DB_HAS_CATEGORIES': {
const payload = message.payload as {
playlistId: string;
type: 'live' | 'movies' | 'series';
};
return hasCategories(db, payload.playlistId, payload.type);
}
case 'DB_GET_CATEGORIES': {
const payload = message.payload as {
playlistId: string;
type: 'live' | 'movies' | 'series';
};
const capturePhase =
createWorkerPerformancePhaseAdapter(performanceCapture);
return getCategories(
db,
payload.playlistId,
payload.type,
capturePhase
);
}
case 'DB_SAVE_CATEGORIES': {
const payload = message.payload as {
playlistId: string;
categories: Array<{
category_name: string;
category_id: string | number;
}>;
type: 'live' | 'movies' | 'series';
hiddenCategoryXtreamIds?: number[];
lockedCategoryXtreamIds?: number[];
};
const capturePhase =
createWorkerPerformancePhaseAdapter(performanceCapture);
return saveCategories(
db,
payload.playlistId,
payload.categories,
payload.type,
payload.hiddenCategoryXtreamIds,
payload.lockedCategoryXtreamIds,
capturePhase
);
}
case 'DB_GET_ALL_CATEGORIES': {
const payload = message.payload as {
playlistId: string;
type: 'live' | 'movies' | 'series';
};
return getAllCategories(db, payload.playlistId, payload.type);
}
case 'DB_UPDATE_CATEGORY_VISIBILITY': {
const payload = message.payload as {
categoryIds: number[];
hidden: boolean;
};
return updateCategoryVisibility(
db,
payload.categoryIds,
payload.hidden
);
}
case 'DB_SET_CATEGORY_LOCKS': {
const payload = message.payload as {
playlistId: string;
type: 'live' | 'movies' | 'series';
lockedXtreamIds: number[];
};
return setCategoryLocks(
db,
payload.playlistId,
payload.type,
payload.lockedXtreamIds
);
}
case 'DB_HAS_CONTENT': {
const payload = message.payload as {
playlistId: string;
type: 'live' | 'movie' | 'series';
};
return hasContent(db, payload.playlistId, payload.type);
}
case 'DB_GET_CONTENT': {
const payload = message.payload as {
playlistId: string;
type: 'live' | 'movie' | 'series';
};
const capturePhase =
createWorkerPerformancePhaseAdapter(performanceCapture);
return getContent(
db,
payload.playlistId,
payload.type,
capturePhase
);
}
case 'DB_GET_GLOBAL_RECENTLY_ADDED': {
const payload = message.payload as {
kind?: 'all' | 'vod' | 'series';
limit?: number;
playlistType?:
'xtream' | 'stalker' | 'm3u-file' | 'm3u-text' | 'm3u-url';
};
return getGlobalRecentlyAdded(
db,
payload.kind,
payload.limit,
payload.playlistType
);
}
case 'DB_SAVE_CONTENT': {
const payload = message.payload as {
playlistId: string;
streams: Array<Record<string, unknown>>;
type: 'live' | 'movie' | 'series';
operationId?: string;
};
const capturePhase =
createWorkerPerformancePhaseAdapter(performanceCapture);
return executeTrackedOperation(
{
requestId: message.requestId,
operation: DB_OPERATION_NAMES.SAVE_CONTENT,
operationId: payload.operationId,
playlistId: payload.playlistId,
},
async (controller) => {
controller.emitStarted({
phase: DB_OPERATION_PHASES.PREPARING_CONTENT,
current: 0,
total: payload.streams.length,
});
const result = await saveContent(
db,
payload.playlistId,
payload.streams,
payload.type,
controller.control,
capturePhase
);
controller.emitCompleted({
phase: DB_OPERATION_PHASES.SAVING_CONTENT,
current: result.count,
total: result.count,
});
return result;
},
preRegisteredState
);
}
case 'DB_CLEAR_XTREAM_IMPORT_CACHE': {
const payload = message.payload as {
playlistId: string;
type: 'live' | 'movie' | 'series';
};
const capturePhase =
createWorkerPerformancePhaseAdapter(performanceCapture);
return clearXtreamImportCache(
db,
payload.playlistId,
payload.type,
capturePhase
);
}
case 'DB_GET_CONTENT_BY_XTREAM_ID': {
const payload = message.payload as {
xtreamId: number;
playlistId: string;
contentType?: 'live' | 'movie' | 'series';
};
return getContentByXtreamId(
db,
payload.xtreamId,
payload.playlistId,
payload.contentType
);
}
case 'DB_SET_CONTENT_METADATA_IF_MISSING': {
const payload = message.payload as {
contentId: number;
patch?: ContentMetadataPatch;
};
return setContentMetadataIfMissing(
db,
payload.contentId,
payload.patch
);
}
case 'DB_SEARCH_CONTENT': {
const payload = message.payload as {
playlistId: string;
searchTerm: string;
types: string[];
excludeHidden?: boolean;
};
const capturePhase =
createWorkerPerformancePhaseAdapter(performanceCapture);
return searchContent(
db,
payload.playlistId,
payload.searchTerm,
payload.types,
payload.excludeHidden,
capturePhase
);
}
case 'DB_GLOBAL_SEARCH': {
const payload = message.payload as {
searchTerm: string;
types: string[];
excludeHidden?: boolean;
sources?: Array<'xtream' | 'm3u'>;
options?: {
limit?: number;
offset?: number;
};
};
return globalSearch(
db,
payload.searchTerm,
payload.types,
payload.excludeHidden,
payload.sources,
payload.options
);
}
case 'DB_CREATE_PLAYLIST': {
return createPlaylist(
db,
message.payload as {
id: string;
name: string;
serverUrl?: string;
username?: string;
password?: string;
macAddress?: string;
url?: string;
type: string;
}
);
}
case 'DB_UPSERT_APP_PLAYLIST': {
const capturePhase: AppPlaylistUpsertPhaseCapture | undefined =
performanceCapture
? {
captureAsync: (phase, execute, metadata) =>
captureWorkerPerformancePhaseAsync(
performanceCapture,
phase,
execute,
metadata
),
captureSync: (phase, execute, metadata) =>
captureWorkerPerformancePhase(
performanceCapture,
phase,
execute,
metadata
),
}
: undefined;
return upsertAppPlaylist(
db,
message.payload as Record<string, unknown>,
capturePhase
);
}
case 'DB_MIGRATE_APP_PLAYLISTS': {
const payload = message.payload as {
playlists: Record<string, unknown>[];
key?: string;
};
return migrateAppPlaylists(db, payload.playlists, payload.key);
}
case 'DB_UPSERT_APP_PLAYLISTS': {
return upsertAppPlaylists(
db,
message.payload as Record<string, unknown>[]
);
}
case 'DB_GET_APP_PLAYLISTS':
return getAppPlaylists(db);
case 'DB_GET_APP_PLAYLIST_METAS':
return getAppPlaylistMetas(db);
case 'DB_GET_APP_PLAYLIST': {
const payload = message.payload as { playlistId: string };
const capturePhase: AppPlaylistGetPhaseCapture | undefined =
performanceCapture
? {
captureAsync: (phase, execute, metadata) =>
captureWorkerPerformancePhaseAsync(
performanceCapture,
phase,
execute,
metadata
),
captureSync: (phase, execute, metadata) =>
captureWorkerPerformancePhase(
performanceCapture,
phase,
execute,
metadata
),
}
: undefined;
return getAppPlaylist(db, payload.playlistId, capturePhase);
}
case 'DB_GET_APP_PLAYLIST_FAVORITE_CHANNELS': {
const payload = message.payload as { playlistId: string };
return getAppPlaylistFavoriteChannels(db, payload.playlistId);
}
case 'DB_GET_PLAYLIST': {
const payload = message.payload as { playlistId: string };
return getPlaylist(db, payload.playlistId);
}
case 'DB_SET_PLAYLIST_SERVER_TIMEZONE': {
const payload = message.payload as {
playlistId: string;
connection: {
serverUrl: string;
username: string;
password: string;
};
serverTimezone: string;
};
return setPlaylistServerTimezone(
db,
payload.playlistId,
payload.connection,
payload.serverTimezone
);
}
case 'DB_UPDATE_PLAYLIST': {
const payload = message.payload as {
playlistId: string;
updates: {
name?: string;
username?: string;
password?: string;
serverUrl?: string;
lastUpdated?: string;
};
};
return updatePlaylist(db, payload.playlistId, payload.updates);
}
case 'DB_DELETE_PLAYLIST': {
const payload = message.payload as {
playlistId: string;
operationId?: string;
};
const capturePhase =
createWorkerPerformancePhaseAdapter(performanceCapture);
return executeTrackedOperation(
{
requestId: message.requestId,
operation: DB_OPERATION_NAMES.DELETE_PLAYLIST,
operationId: payload.operationId,
playlistId: payload.playlistId,
},
async (controller) => {
controller.emitStarted({
phase: DB_OPERATION_PHASES.DELETING_FAVORITES,
current: 0,
});
const result = await deletePlaylist(
db,
payload.playlistId,
controller.control,
capturePhase
);
controller.emitCompleted({
phase: DB_OPERATION_PHASES.DELETING_PLAYLIST,
current: 1,
total: 1,
});
return result;
},
preRegisteredState
);
}
case 'DB_GET_APP_STATE': {
const payload = message.payload as { key: string };
return getAppState(db, payload.key);
}
case 'DB_SET_APP_STATE': {
const payload = message.payload as { key: string; value: string };
return setAppState(db, payload.key, payload.value);
}
case 'DB_GET_TMDB_METADATA': {
const payload = message.payload as {
mediaType: 'movie' | 'tv';
lookupKey: string;
language: string;
};
return getTmdbMetadata(
db,
payload.mediaType,
payload.lookupKey,
payload.language
);
}
case 'DB_SET_TMDB_METADATA': {
const payload = message.payload as {
entry: Parameters<typeof setTmdbMetadata>[1];
};
return setTmdbMetadata(db, payload.entry);
}
case 'DB_GET_TMDB_CACHE_STATS':
return getTmdbCacheStats(db);
case 'DB_CLEAR_TMDB_METADATA':
return clearTmdbMetadata(db);
case 'DB_MATCH_TITLES': {
const payload = message.payload as { titles: string[] };
return matchTitles(db, payload.titles);
}
case 'DB_FIND_TITLE_SOURCES': {
const payload = message.payload as {
request: FindTitleSourcesRequest;
};
return findTitleSources(db, payload.request);
}
case 'DB_GET_VOD_SOURCE_PIN': {
const payload = message.payload as { matchKeys: string[] };
return getVodSourcePin(db, payload.matchKeys);
}
case 'DB_CLEAR_VOD_SOURCE_PINS_FOR_PLAYLIST': {
const payload = message.payload as { playlistId: string };
return clearVodSourcePinsForPlaylist(db, payload.playlistId);
}
case 'DB_LIST_VOD_SOURCE_PINS': {
const payload = message.payload as { playlistId: string };
return listVodSourcePinsForPlaylist(db, payload.playlistId);
}
case 'DB_SET_VOD_SOURCE_PIN': {
const payload = message.payload as {
pin: VodSourcePin;
retireKeys?: string[];
aliasKeys?: string[];
};
return setVodSourcePin(
db,
payload.pin,
payload.retireKeys ?? [],
payload.aliasKeys ?? []
);
}
case 'DB_REPLACE_VOD_SOURCE_PINS': {
const payload = message.payload as {
playlistId: string;
pins: VodSourcePin[];
};
return replaceVodSourcePinsForPlaylist(
db,
payload.playlistId,
payload.pins ?? []
);
}
case 'DB_CLEAR_VOD_SOURCE_PIN': {
const payload = message.payload as { matchKeys: string[] };
return clearVodSourcePin(db, payload.matchKeys);
}
case 'DB_DELETE_ALL_PLAYLISTS': {
const payload = message.payload as { operationId?: string };
return executeTrackedOperation(
{
requestId: message.requestId,
operation: DB_OPERATION_NAMES.DELETE_ALL_PLAYLISTS,
operationId: payload.operationId,
cancellable: false,
},
async (controller) => {
controller.emitStarted({
phase: DB_OPERATION_PHASES.DELETING_FAVORITES,
current: 0,
total: 7,
});
const result = await deleteAllPlaylists(
db,
controller.control
);
controller.emitCompleted({
phase: DB_OPERATION_PHASES.DELETING_PLAYLISTS,
current: 7,
total: 7,
});
return result;
}
);
}
case 'DB_DELETE_XTREAM_CONTENT': {
const payload = message.payload as {
playlistId: string;
operationId?: string;
};
const capturePhase =
createWorkerPerformancePhaseAdapter(performanceCapture);
return executeTrackedOperation(
{
requestId: message.requestId,
operation: DB_OPERATION_NAMES.DELETE_XTREAM_CONTENT,
operationId: payload.operationId,
playlistId: payload.playlistId,
},
async (controller) => {
controller.emitStarted({
phase: DB_OPERATION_PHASES.COLLECTING_USER_DATA,
current: 0,
});
const result = await deleteXtreamContent(
db,
payload.playlistId,
controller.control,
capturePhase
);
controller.emitCompleted({
phase: DB_OPERATION_PHASES.DELETING_CATEGORIES,
});
return result;
},
preRegisteredState
);
}
case 'DB_RESTORE_XTREAM_USER_DATA': {
const payload = message.payload as {
playlistId: string;
favorites: XtreamBackupFavoriteItem[];
recentlyViewed: XtreamBackupRecentlyViewedItem[];
operationId?: string;
};
return executeTrackedOperation(
{
requestId: message.requestId,
operation: DB_OPERATION_NAMES.RESTORE_XTREAM_USER_DATA,
operationId: payload.operationId,
playlistId: payload.playlistId,
},
async (controller) => {
const totalItems =
payload.favorites.length +
payload.recentlyViewed.length;
controller.emitStarted({
phase: DB_OPERATION_PHASES.RESTORING_FAVORITES,
current: 0,
total: totalItems,
});
const result = await restoreXtreamUserData(
db,
payload.playlistId,
payload.favorites,
payload.recentlyViewed,
controller.control
);
controller.emitCompleted({
phase: DB_OPERATION_PHASES.RESTORING_RECENTLY_VIEWED,
current: totalItems,
total: totalItems,
});
return result;
},
preRegisteredState
);
}
case 'DB_ADD_FAVORITE': {
const payload = message.payload as {
contentId: number;
playlistId: string;
backdropUrl?: string;
};
return addFavorite(db, payload.contentId, payload.playlistId, {
backdropUrl: payload.backdropUrl,
});
}
case 'DB_REMOVE_FAVORITE': {
const payload = message.payload as {
contentId: number;
playlistId: string;
};
return removeFavorite(db, payload.contentId, payload.playlistId);
}
case 'DB_IS_FAVORITE': {
const payload = message.payload as {
contentId: number;
playlistId: string;
};
return isFavorite(db, payload.contentId, payload.playlistId);
}
case 'DB_GET_FAVORITES': {
const payload = message.payload as { playlistId: string };
return getFavorites(db, payload.playlistId);
}
case 'DB_GET_GLOBAL_FAVORITES':
return getGlobalFavorites(db);
case 'DB_GET_ALL_GLOBAL_FAVORITES':
return getAllGlobalFavorites(db);
case 'DB_REORDER_GLOBAL_FAVORITES': {
const payload = message.payload as {
updates: {
content_id: number;
playlist_id: string;
position: number;
}[];
};
return reorderGlobalFavorites(db, payload.updates);
}
case 'DB_GET_RECENTLY_VIEWED':
return getRecentlyViewed(db);
case 'DB_CLEAR_RECENTLY_VIEWED':
return clearRecentlyViewed(db);
case 'DB_GET_RECENT_ITEMS': {
const payload = message.payload as { playlistId: string };
return getRecentItems(db, payload.playlistId);
}
case 'DB_ADD_RECENT_ITEM': {
const payload = message.payload as {
contentId: number;
playlistId: string;
backdropUrl?: string;
};
return addRecentItem(db, payload.contentId, payload.playlistId, {
backdropUrl: payload.backdropUrl,
});
}
case 'DB_CLEAR_PLAYLIST_RECENT_ITEMS': {
const payload = message.payload as { playlistId: string };
return clearPlaylistRecentItems(db, payload.playlistId);
}
case 'DB_REMOVE_RECENT_ITEM': {
const payload = message.payload as {
contentId: number;
playlistId: string;
};
return removeRecentItem(db, payload.contentId, payload.playlistId);
}
case 'DB_REMOVE_RECENT_ITEMS_BATCH': {
const payload = message.payload as {
items: { contentId: number; playlistId: string }[];
};
return removeRecentItemsBatch(db, payload.items);
}
case 'DB_SAVE_PLAYBACK_POSITION': {
const payload = message.payload as {
playlistId: string;
data: {
contentXtreamId: number;
contentType: 'vod' | 'episode';
seriesXtreamId?: number;
seasonNumber?: number;
episodeNumber?: number;
positionSeconds: number;
durationSeconds?: number;
playlistType?:
| 'xtream'
| 'stalker'
| 'm3u-file'
| 'm3u-text'
| 'm3u-url';
};
};
return savePlaybackPosition(db, payload.playlistId, payload.data);
}
case 'DB_GET_PLAYBACK_POSITION': {
const payload = message.payload as {
playlistId: string;
contentXtreamId: number;
contentType: 'vod' | 'episode';
};
return getPlaybackPosition(
db,
payload.playlistId,
payload.contentXtreamId,
payload.contentType
);
}
case 'DB_GET_SERIES_PLAYBACK_POSITIONS': {
const payload = message.payload as {
playlistId: string;
seriesXtreamId: number;
};
return getSeriesPlaybackPositions(
db,
payload.playlistId,
payload.seriesXtreamId
);
}
case 'DB_GET_RECENT_PLAYBACK_POSITIONS': {
const payload = message.payload as {
playlistId: string;
limit?: number;
};
return getRecentPlaybackPositions(
db,
payload.playlistId,
payload.limit
);
}
case 'DB_GET_ALL_PLAYBACK_POSITIONS': {
const payload = message.payload as { playlistId: string };
return getAllPlaybackPositions(db, payload.playlistId);
}
case 'DB_CLEAR_ALL_PLAYBACK_POSITIONS': {
const payload = message.payload as { playlistId: string };
return clearAllPlaybackPositions(db, payload.playlistId);
}
case 'DB_CLEAR_PLAYBACK_POSITION': {
const payload = message.payload as {
playlistId: string;
contentXtreamId: number;
contentType: 'vod' | 'episode';
};
return clearPlaybackPosition(
db,
payload.playlistId,
payload.contentXtreamId,
payload.contentType
);
}
case 'DB_SAVE_PLAYBACK_POSITIONS_BATCH': {
const payload = message.payload as {
playlistId: string;
items: {
contentXtreamId: number;
contentType: 'vod' | 'episode';
seriesXtreamId?: number;
seasonNumber?: number;
episodeNumber?: number;
positionSeconds: number;
durationSeconds?: number;
playlistType?:
| 'xtream'
| 'stalker'
| 'm3u-file'
| 'm3u-text'
| 'm3u-url';
}[];
};
return savePlaybackPositionsBatch(
db,
payload.playlistId,
payload.items
);
}
case 'DB_CLEAR_PLAYBACK_POSITIONS_BATCH': {
const payload = message.payload as {
playlistId: string;
items: {
contentXtreamId: number;
contentType: 'vod' | 'episode';
}[];
};
return clearPlaybackPositionsBatch(
db,
payload.playlistId,
payload.items
);
}
}
}
parentPort.on('message', async (message: DbWorkerIncomingMessage) => {
if (isDatabaseWorkerPostGcHeapRequest(message)) {
handleDatabaseWorkerPostGcHeapRequest(
message,
activePerformanceCaptures.size === 0
);
return;
}
if (message.type === 'parental-lock') {
setParentalLockActive(message.active === true);
return;
}
if (message.type === 'cancel') {
const activeOperation = activeOperations.get(message.operationId);
if (activeOperation) {
activeOperation.cancelled = true;
if (activeOperation.performanceRequestId !== undefined) {
publishDatabaseWorkerCancelReceipt(
message.operationId,
activeOperation.performanceRequestId,
postMessage
);
}
}
return;
}
const preRegisteredOperation = preRegisterCancellableOperation(message);
let performanceCapture: WorkerPerformanceCapture | null = null;
const execution = await (async () => {
try {
performanceCapture = startWorkerPerformanceCapture({
requestId: message.requestId,
});
if (performanceCapture && preRegisteredOperation) {
preRegisteredOperation.state.performanceRequestId =
message.requestId;
}
registerDatabaseWorkerPerformanceCapture(
activePerformanceCaptures,
performanceCapture
);
await armWorkerPerformanceCapture(performanceCapture);
return await executeWithWorkerPerformanceCapture(
performanceCapture,
() =>
executeRequest(
message,
preRegisteredOperation?.state,
performanceCapture
)
);
} finally {
releaseDatabaseWorkerPerformanceCapture(
activePerformanceCaptures,
performanceCapture
);
releasePreRegisteredOperation(preRegisteredOperation);
}
})();
if (execution.success) {
postMessage({
type: 'response',
requestId: message.requestId,
success: true,
result: execution.result,
performance: stampWorkerPerformanceResponsePostedEpoch(
performanceCapture,
execution.performance
),
});
} else {
console.error(
loggerLabel,
`Error handling ${message.operation}:`,
execution.error
);
postMessage({
type: 'response',
requestId: message.requestId,
success: false,
error: serializeError(execution.error),
performance: stampWorkerPerformanceResponsePostedEpoch(
performanceCapture,
execution.performance
),
});
}
});
process.on('exit', () => {
closeWorkerDatabase();
});
postMessage({ type: 'ready' });