diff --git a/AGENTS.md b/AGENTS.md index f3cd7a990..c1b428902 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -85,7 +85,7 @@ IPTVNATOR_TRACE_STARTUP=1 nx serve electron-backend - `IPTVNATOR_TRACE_PLAYER=1` traces external-player activity and bounded Embedded MPV runtime-probe stderr - `IPTVNATOR_TRACE_RENDERER_CONSOLE=1` mirrors renderer console output into the Electron terminal - `IPTVNATOR_PERF_CAPTURE=1` enables development/test-only, redacted preload IPC phase markers for refresh/DB benchmark correlation; benchmark tooling sets it explicitly, and production launches must leave it unset - - `IPTVNATOR_PERF_WORKER_PROFILING=1` enables development/test-only event-loop metrics in database and playlist-refresh worker responses; the performance benchmark sets it automatically, and production launches must leave it unset + - `IPTVNATOR_PERF_WORKER_PROFILING=1` enables development/test-only, request-scoped worker timestamps, thread CPU, event-loop utilization, and event-loop delay metrics in database and playlist-refresh responses; overlapping database requests are explicitly invalidated instead of misattributed, the performance benchmark sets the flag automatically, and production launches must leave it unset - Settings, portal request/response, and trace payloads must use `@iptvnator/shared/logging` or the redacting portal logger before reaching diff --git a/CLAUDE.md b/CLAUDE.md index 72eed72c7..bd872ab79 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -142,7 +142,7 @@ Useful narrower flags: - `IPTVNATOR_TRACE_PLAYER=1` traces external-player activity and bounded Embedded MPV runtime-probe stderr - `IPTVNATOR_TRACE_RENDERER_CONSOLE=1` mirrors renderer console logs into the Electron terminal - `IPTVNATOR_PERF_CAPTURE=1` enables development/test-only, redacted preload IPC phase markers for refresh/DB benchmark correlation; benchmark tooling sets it explicitly, and production launches must leave it unset -- `IPTVNATOR_PERF_WORKER_PROFILING=1` enables development/test-only event-loop metrics in database and playlist-refresh worker responses; the performance benchmark sets it automatically, and production launches must leave it unset +- `IPTVNATOR_PERF_WORKER_PROFILING=1` enables development/test-only, request-scoped worker timestamps, thread CPU, event-loop utilization, and event-loop delay metrics in database and playlist-refresh responses; overlapping database requests are explicitly invalidated instead of misattributed, the performance benchmark sets the flag automatically, and production launches must leave it unset Settings, portal request/response, and trace payloads must use `@iptvnator/shared/logging` or the redacting portal logger before reaching diff --git a/apps/electron-backend-e2e/src/performance/database-request-identity-capture.spec.ts b/apps/electron-backend-e2e/src/performance/database-request-identity-capture.spec.ts new file mode 100644 index 000000000..2e9672718 --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/database-request-identity-capture.spec.ts @@ -0,0 +1,229 @@ +/* eslint-disable playwright/expect-expect -- These are Node assertion-based performance contract tests. */ +import assert from 'node:assert/strict'; +import { EventEmitter } from 'node:events'; +import test from 'node:test'; + +import { + installDatabaseRequestIdentityCaptureInMain, + type DatabaseRequestIdentityCaptureApi, +} from './database-request-identity-capture'; + +const CHANNEL = 'IPTVNATOR_PERF_CAPTURE_MARKER'; +let nextStateId = 1; + +function createCapture(): { + api: DatabaseRequestIdentityCaptureApi; + ipcMain: EventEmitter; +} { + const ipcMain = new EventEmitter(); + const stateKey = `__identityCapture${nextStateId}`; + nextStateId += 1; + installDatabaseRequestIdentityCaptureInMain({ ipcMain } as never, stateKey); + const api = (globalThis as unknown as Record)[ + stateKey + ] as DatabaseRequestIdentityCaptureApi; + api.start(); + return { api, ipcMain }; +} + +function marker(overrides: Record = {}) { + return { + correlationState: 'correlated', + invalidReason: null, + ipcCallId: 2, + method: 'dbGetAppPlaylist', + operationId: 'refresh-operation-1', + phase: 'start', + playlistId: 'playlist-1', + sourceEpochMs: 100, + ...overrides, + }; +} + +test('correlates the real preload marker to a DB request whose payload omits operationId', () => { + const { api, ipcMain } = createCapture(); + ipcMain.emit(CHANNEL, {}, marker()); + + const identity = api.matchDatabaseRequest({ + operation: 'DB_GET_APP_PLAYLIST', + payload: { playlistId: 'playlist-1' }, + requestId: 'request-1', + type: 'request', + }); + + assert.deepEqual(identity, { + operationId: 'refresh-operation-1', + operationIdUnavailableReason: null, + }); +}); + +test('resolves a response-retained request identity when the marker arrives later and resets between captures', () => { + const { api, ipcMain } = createCapture(); + const identity = api.matchDatabaseRequest({ + operation: 'DB_GET_APP_PLAYLIST', + payload: { playlistId: 'playlist-1' }, + requestId: 'request-before-marker', + type: 'request', + }); + assert.equal( + identity.operationIdUnavailableReason, + 'preload-performance-marker-missing' + ); + const responseOutcome = { identity }; + + ipcMain.emit(CHANNEL, {}, marker()); + assert.deepEqual(responseOutcome.identity, { + operationId: 'refresh-operation-1', + operationIdUnavailableReason: null, + }); + + api.stop(); + ipcMain.emit(CHANNEL, {}, marker({ operationId: 'marker-while-stopped' })); + api.start(); + const afterReset = api.matchDatabaseRequest({ + operation: 'DB_GET_APP_PLAYLIST', + payload: { playlistId: 'playlist-1' }, + requestId: 'request-after-reset', + type: 'request', + }); + assert.deepEqual(afterReset, { + operationId: null, + operationIdUnavailableReason: 'preload-performance-marker-missing', + }); +}); + +test('fails closed for missing and ambiguous preload markers', () => { + const missingCapture = createCapture(); + const missing = missingCapture.api.matchDatabaseRequest({ + operation: 'DB_GET_APP_PLAYLIST', + payload: { playlistId: 'playlist-1' }, + requestId: 'request-missing', + type: 'request', + }); + assert.deepEqual(missing, { + operationId: null, + operationIdUnavailableReason: 'preload-performance-marker-missing', + }); + + const ambiguousCapture = createCapture(); + ambiguousCapture.ipcMain.emit(CHANNEL, {}, marker()); + ambiguousCapture.ipcMain.emit( + CHANNEL, + {}, + marker({ + ipcCallId: 3, + operationId: 'refresh-operation-2', + }) + ); + const ambiguous = ambiguousCapture.api.matchDatabaseRequest({ + operation: 'DB_GET_APP_PLAYLIST', + payload: { playlistId: 'playlist-1' }, + requestId: 'request-ambiguous', + type: 'request', + }); + assert.deepEqual(ambiguous, { + operationId: null, + operationIdUnavailableReason: 'preload-performance-marker-ambiguous', + }); +}); + +test('quarantines a key after response-before-marker ambiguity', () => { + const { api, ipcMain } = createCapture(); + const requestA = api.matchDatabaseRequest({ + operation: 'DB_GET_APP_PLAYLIST', + payload: { playlistId: 'playlist-1' }, + requestId: 'request-a', + type: 'request', + }); + const requestB = api.matchDatabaseRequest({ + operation: 'DB_GET_APP_PLAYLIST', + payload: { playlistId: 'playlist-1' }, + requestId: 'request-b', + type: 'request', + }); + + ipcMain.emit( + CHANNEL, + {}, + marker({ ipcCallId: 2, operationId: 'refresh-operation-a' }) + ); + ipcMain.emit( + CHANNEL, + {}, + marker({ ipcCallId: 3, operationId: 'refresh-operation-b' }) + ); + + const requestC = api.matchDatabaseRequest({ + operation: 'DB_GET_APP_PLAYLIST', + payload: { playlistId: 'playlist-1' }, + requestId: 'request-c', + type: 'request', + }); + ipcMain.emit( + CHANNEL, + {}, + marker({ + ipcCallId: 4, + operationId: 'other-playlist-operation', + playlistId: 'playlist-2', + }) + ); + const otherPlaylist = api.matchDatabaseRequest({ + operation: 'DB_GET_APP_PLAYLIST', + payload: { playlistId: 'playlist-2' }, + requestId: 'other-playlist-request', + type: 'request', + }); + ipcMain.emit( + CHANNEL, + {}, + marker({ + ipcCallId: 5, + method: 'dbUpsertAppPlaylist', + operationId: 'other-db-operation', + }) + ); + const otherOperation = api.matchDatabaseRequest({ + operation: 'DB_UPSERT_APP_PLAYLIST', + payload: { _id: 'playlist-1' }, + requestId: 'other-operation-request', + type: 'request', + }); + + const ambiguousIdentity = { + operationId: null, + operationIdUnavailableReason: 'preload-performance-marker-ambiguous', + }; + assert.deepEqual(requestA, ambiguousIdentity); + assert.deepEqual(requestB, ambiguousIdentity); + assert.deepEqual(requestC, ambiguousIdentity); + assert.deepEqual(otherPlaylist, { + operationId: 'other-playlist-operation', + operationIdUnavailableReason: null, + }); + assert.deepEqual(otherOperation, { + operationId: 'other-db-operation', + operationIdUnavailableReason: null, + }); + + api.stop(); + api.start(); + ipcMain.emit( + CHANNEL, + {}, + marker({ + ipcCallId: 6, + operationId: 'fresh-capture-operation', + }) + ); + const freshCaptureIdentity = api.matchDatabaseRequest({ + operation: 'DB_GET_APP_PLAYLIST', + payload: { playlistId: 'playlist-1' }, + requestId: 'fresh-capture-request', + type: 'request', + }); + assert.deepEqual(freshCaptureIdentity, { + operationId: 'fresh-capture-operation', + operationIdUnavailableReason: null, + }); +}); diff --git a/apps/electron-backend-e2e/src/performance/database-request-identity-capture.ts b/apps/electron-backend-e2e/src/performance/database-request-identity-capture.ts new file mode 100644 index 000000000..a0912d3e3 --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/database-request-identity-capture.ts @@ -0,0 +1,270 @@ +import type { ElectronApplication } from '@playwright/test'; + +export const DATABASE_REQUEST_IDENTITY_CAPTURE_STATE_KEY = + '__iptvnatorM3uRefreshDatabaseRequestIdentityCapture'; + +export interface DatabaseRequestIdentity { + operationId: string | null; + operationIdUnavailableReason: string | null; +} + +export interface DatabaseRequestIdentityCaptureApi { + matchDatabaseRequest(message: unknown): DatabaseRequestIdentity; + start(): void; + stop(): void; +} + +interface DatabaseRequestIdentityElectron { + readonly ipcMain: { + on( + channel: string, + listener: (event: unknown, marker: unknown) => void + ): void; + }; +} + +export function installDatabaseRequestIdentityCaptureInMain( + { ipcMain }: DatabaseRequestIdentityElectron, + stateKey: string +): void { + type JsonRecord = Record; + interface PendingMarker { + readonly ipcCallId: number; + readonly operation: string; + readonly operationId: string; + readonly playlistId: string; + } + interface PendingRequest { + readonly identity: DatabaseRequestIdentity; + readonly operation: string; + readonly playlistId: string; + } + interface QuarantinedKey { + readonly operation: string; + readonly playlistId: string; + } + + const target = globalThis as unknown as Record; + if (target[stateKey] !== undefined) { + return; + } + + const markerChannel = 'IPTVNATOR_PERF_CAPTURE_MARKER'; + const markerMethodToOperation: Readonly> = { + dbGetAppPlaylist: 'DB_GET_APP_PLAYLIST', + dbUpsertAppPlaylist: 'DB_UPSERT_APP_PLAYLIST', + }; + let active = false; + let pendingMarkers: PendingMarker[] = []; + let pendingRequests: PendingRequest[] = []; + let quarantinedKeys: QuarantinedKey[] = []; + + const readRecord = (value: unknown): JsonRecord | null => + typeof value === 'object' && value !== null + ? (value as JsonRecord) + : null; + const readPlaylistId = ( + operation: string, + payload: unknown + ): string | null => { + const value = readRecord(payload); + if (!value) { + return null; + } + const playlistId = + operation === 'DB_UPSERT_APP_PLAYLIST' + ? value['_id'] + : value['playlistId']; + return typeof playlistId === 'string' && playlistId.length > 0 + ? playlistId + : null; + }; + const matches = ( + candidate: { operation: string; playlistId: string }, + operation: string, + playlistId: string + ): boolean => + candidate.operation === operation && + candidate.playlistId === playlistId; + const markAmbiguous = (requests: PendingRequest[]): void => { + for (const request of requests) { + request.identity.operationId = null; + request.identity.operationIdUnavailableReason = + 'preload-performance-marker-ambiguous'; + } + }; + const isQuarantined = (operation: string, playlistId: string): boolean => + quarantinedKeys.some((key) => matches(key, operation, playlistId)); + const quarantine = (operation: string, playlistId: string): void => { + if (!isQuarantined(operation, playlistId)) { + quarantinedKeys.push({ operation, playlistId }); + } + const requestCandidates = pendingRequests.filter((request) => + matches(request, operation, playlistId) + ); + markAmbiguous(requestCandidates); + pendingRequests = pendingRequests.filter( + (request) => !matches(request, operation, playlistId) + ); + pendingMarkers = pendingMarkers.filter( + (marker) => !matches(marker, operation, playlistId) + ); + }; + + ipcMain.on(markerChannel, (_event, input) => { + if (!active) { + return; + } + const marker = readRecord(input); + if ( + !marker || + marker['phase'] !== 'start' || + marker['correlationState'] !== 'correlated' || + marker['invalidReason'] !== null + ) { + return; + } + const operation = + typeof marker['method'] === 'string' + ? markerMethodToOperation[marker['method']] + : undefined; + const operationId = marker['operationId']; + const playlistId = marker['playlistId']; + const ipcCallId = marker['ipcCallId']; + if ( + operation === undefined || + typeof operationId !== 'string' || + operationId.length === 0 || + typeof playlistId !== 'string' || + playlistId.length === 0 || + !Number.isSafeInteger(ipcCallId) || + Number(ipcCallId) < 1 + ) { + return; + } + + if (isQuarantined(operation, playlistId)) { + return; + } + const requestCandidates = pendingRequests.filter((request) => + matches(request, operation, playlistId) + ); + if (requestCandidates.length === 1) { + const request = requestCandidates[0]; + if (!request) { + return; + } + request.identity.operationId = operationId; + request.identity.operationIdUnavailableReason = null; + pendingRequests = pendingRequests.filter( + (candidate) => candidate !== request + ); + return; + } + if (requestCandidates.length > 1) { + quarantine(operation, playlistId); + return; + } + + pendingMarkers.push({ + ipcCallId: Number(ipcCallId), + operation, + operationId, + playlistId, + }); + pendingMarkers.sort((left, right) => left.ipcCallId - right.ipcCallId); + }); + + const api: DatabaseRequestIdentityCaptureApi = { + matchDatabaseRequest: (input) => { + const message = readRecord(input); + const operation = message?.['operation']; + const playlistId = + typeof operation === 'string' + ? readPlaylistId(operation, message?.['payload']) + : null; + if ( + !active || + (operation !== 'DB_GET_APP_PLAYLIST' && + operation !== 'DB_UPSERT_APP_PLAYLIST') || + playlistId === null + ) { + return { + operationId: null, + operationIdUnavailableReason: + 'preload-performance-marker-not-applicable', + }; + } + + if (isQuarantined(operation, playlistId)) { + return { + operationId: null, + operationIdUnavailableReason: + 'preload-performance-marker-ambiguous', + }; + } + const markerCandidates = pendingMarkers.filter((marker) => + matches(marker, operation, playlistId) + ); + if (markerCandidates.length === 1) { + const marker = markerCandidates[0]; + if (!marker) { + return { + operationId: null, + operationIdUnavailableReason: + 'preload-performance-marker-missing', + }; + } + pendingMarkers = pendingMarkers.filter( + (candidate) => candidate !== marker + ); + return { + operationId: marker.operationId, + operationIdUnavailableReason: null, + }; + } + if (markerCandidates.length > 1) { + quarantine(operation, playlistId); + return { + operationId: null, + operationIdUnavailableReason: + 'preload-performance-marker-ambiguous', + }; + } + + const identity: DatabaseRequestIdentity = { + operationId: null, + operationIdUnavailableReason: + 'preload-performance-marker-missing', + }; + pendingRequests.push({ + identity, + operation, + playlistId, + }); + return identity; + }, + start: () => { + pendingMarkers = []; + pendingRequests = []; + quarantinedKeys = []; + active = true; + }, + stop: () => { + active = false; + pendingMarkers = []; + pendingRequests = []; + quarantinedKeys = []; + }, + }; + target[stateKey] = api; +} + +export async function installDatabaseRequestIdentityCapture( + electronApp: ElectronApplication +): Promise { + await electronApp.evaluate( + installDatabaseRequestIdentityCaptureInMain, + DATABASE_REQUEST_IDENTITY_CAPTURE_STATE_KEY + ); +} diff --git a/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation-contract.ts b/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation-contract.ts index f0abcd22f..344767515 100644 --- a/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation-contract.ts +++ b/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation-contract.ts @@ -48,6 +48,29 @@ export interface MainTimelineRecord { readonly type: string; } +export interface WorkerRequestPerformanceMetrics { + readonly eventLoopDelay: EventLoopDelayMetrics | null; + readonly eventLoopDelayUnavailableReason: string | null; + readonly eventLoopUtilization: number | null; + readonly eventLoopUtilizationUnavailableReason: string | null; + readonly histogramFlushedEpochMs: number | null; + readonly invalidReason: string | null; + readonly operation: string | null; + readonly operationId: string | null; + readonly operationIdUnavailableReason: string | null; + readonly performanceCaptureUnavailableReason: string | null; + readonly playlistId: string | null; + readonly requestId: string | null; + readonly requestReceivedEpochMs: number | null; + readonly responseEpochMs: number; + readonly success: boolean; + readonly threadCpuSystemMicros: number | null; + readonly threadCpuUnavailableReason: string | null; + readonly threadCpuUserMicros: number | null; + readonly workEndedEpochMs: number | null; + readonly workStartedEpochMs: number | null; +} + export interface WorkerCaptureMetrics { readonly cancelPostedEpochMs: number | null; readonly cpuSystemMicros: number | null; @@ -62,6 +85,7 @@ export interface WorkerCaptureMetrics { readonly playlistId: string | null; readonly postGcHeapUsedBytes: number | null; readonly profilePath: string | null; + readonly requests: readonly WorkerRequestPerformanceMetrics[]; readonly responseEpochMs: number | null; readonly snapshotPath: string | null; readonly terminatedEpochMs: number | null; @@ -176,13 +200,15 @@ export interface CancellationBenchmarkSummary { readonly cancelToDurableTerminalMs: NumericDistribution; readonly cancelToWorkerTerminatedMs: NumericDistribution; readonly cancelTransportLatencyMs: NumericDistribution; - readonly databaseWorkerEventLoopDelayMaxMs: NumericDistribution; - readonly databaseWorkerEventLoopDelayP95Ms: NumericDistribution; - readonly databaseWorkerEventLoopDelayP99Ms: NumericDistribution; - readonly databaseWorkerEventLoopUtilization: NumericDistribution; readonly databaseWorkerExternalPeakBytes: NumericDistribution; readonly databaseWorkerHeapPeakBytes: NumericDistribution; readonly databaseWorkerPostGcHeapBytes: NumericDistribution; + readonly databaseWorkerRequestEventLoopDelayMaxMs: NumericDistribution; + readonly databaseWorkerRequestEventLoopDelayP95Ms: NumericDistribution; + readonly databaseWorkerRequestEventLoopDelayP99Ms: NumericDistribution; + readonly databaseWorkerRequestEventLoopUtilization: NumericDistribution; + readonly databaseWorkerRequestThreadCpuSystemMicros: NumericDistribution; + readonly databaseWorkerRequestThreadCpuUserMicros: NumericDistribution; readonly dataFetchMs: NumericDistribution; readonly dbCommitProxyMs: NumericDistribution; readonly ipcStoreDispatchProxyMs: NumericDistribution; @@ -196,13 +222,15 @@ export interface CancellationBenchmarkSummary { readonly mainRssPostGcBytes: NumericDistribution; readonly parseNormalizeCloneProxyMs: NumericDistribution; readonly persistencePreparationProxyMs: NumericDistribution; - readonly playlistWorkerEventLoopUtilization: NumericDistribution; - readonly playlistWorkerEventLoopDelayMaxMs: NumericDistribution; - readonly playlistWorkerEventLoopDelayP95Ms: NumericDistribution; - readonly playlistWorkerEventLoopDelayP99Ms: NumericDistribution; readonly playlistWorkerExternalPeakBytes: NumericDistribution; readonly playlistWorkerHeapPeakBytes: NumericDistribution; readonly playlistWorkerPostGcHeapBytes: NumericDistribution; + readonly playlistWorkerRequestEventLoopDelayMaxMs: NumericDistribution; + readonly playlistWorkerRequestEventLoopDelayP95Ms: NumericDistribution; + readonly playlistWorkerRequestEventLoopDelayP99Ms: NumericDistribution; + readonly playlistWorkerRequestEventLoopUtilization: NumericDistribution; + readonly playlistWorkerRequestThreadCpuSystemMicros: NumericDistribution; + readonly playlistWorkerRequestThreadCpuUserMicros: NumericDistribution; readonly rendererFrameGapMs: NumericDistribution; readonly rendererHeartbeatDelayMs: NumericDistribution; readonly rendererHeapPeakBytes: NumericDistribution; diff --git a/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation-report.ts b/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation-report.ts index 2f063cd0a..7a8e4c220 100644 --- a/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation-report.ts +++ b/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation-report.ts @@ -9,6 +9,7 @@ import type { NumericDistribution, PerformanceWorkerKind, RendererCaptureMetrics, + WorkerRequestPerformanceMetrics, } from './m3u-refresh-cancellation-contract'; import { PERFORMANCE_ITERATION_KIND, @@ -63,15 +64,28 @@ export function createCancellationBenchmarkSummary( return worker ? select(worker) : null; }) ); - const workerEventLoopDelayMetric = ( + const workerRequestMetric = ( + kind: PerformanceWorkerKind, + select: (request: WorkerRequestPerformanceMetrics) => number | null + ): NumericDistribution => + summarizeNumbers( + measured.flatMap((iteration) => + iteration.main.workers + .filter((worker) => worker.kind === kind) + .flatMap((worker) => + worker.requests.map((request) => select(request)) + ) + ) + ); + const workerRequestEventLoopDelayMetric = ( kind: PerformanceWorkerKind, percentile: keyof NonNullable< - MainCaptureMetrics['workers'][number]['eventLoopDelay'] + WorkerRequestPerformanceMetrics['eventLoopDelay'] > ): NumericDistribution => - workerMetric( + workerRequestMetric( kind, - (worker) => worker.eventLoopDelay?.[percentile] ?? null + (request) => request.eventLoopDelay?.[percentile] ?? null ); return Object.freeze({ @@ -105,22 +119,6 @@ export function createCancellationBenchmarkSummary( (iteration) => iteration.phases.cancelTransportLatencyMs ) ), - databaseWorkerEventLoopDelayMaxMs: workerEventLoopDelayMetric( - PERFORMANCE_WORKER_KIND.DATABASE, - 'maxMs' - ), - databaseWorkerEventLoopDelayP95Ms: workerEventLoopDelayMetric( - PERFORMANCE_WORKER_KIND.DATABASE, - 'p95Ms' - ), - databaseWorkerEventLoopDelayP99Ms: workerEventLoopDelayMetric( - PERFORMANCE_WORKER_KIND.DATABASE, - 'p99Ms' - ), - databaseWorkerEventLoopUtilization: workerMetric( - PERFORMANCE_WORKER_KIND.DATABASE, - (worker) => worker.eventLoopUtilization - ), databaseWorkerExternalPeakBytes: workerMetric( PERFORMANCE_WORKER_KIND.DATABASE, (worker) => worker.peakExternalBytes @@ -133,6 +131,33 @@ export function createCancellationBenchmarkSummary( PERFORMANCE_WORKER_KIND.DATABASE, (worker) => worker.postGcHeapUsedBytes ), + databaseWorkerRequestEventLoopDelayMaxMs: + workerRequestEventLoopDelayMetric( + PERFORMANCE_WORKER_KIND.DATABASE, + 'maxMs' + ), + databaseWorkerRequestEventLoopDelayP95Ms: + workerRequestEventLoopDelayMetric( + PERFORMANCE_WORKER_KIND.DATABASE, + 'p95Ms' + ), + databaseWorkerRequestEventLoopDelayP99Ms: + workerRequestEventLoopDelayMetric( + PERFORMANCE_WORKER_KIND.DATABASE, + 'p99Ms' + ), + databaseWorkerRequestEventLoopUtilization: workerRequestMetric( + PERFORMANCE_WORKER_KIND.DATABASE, + (request) => request.eventLoopUtilization + ), + databaseWorkerRequestThreadCpuSystemMicros: workerRequestMetric( + PERFORMANCE_WORKER_KIND.DATABASE, + (request) => request.threadCpuSystemMicros + ), + databaseWorkerRequestThreadCpuUserMicros: workerRequestMetric( + PERFORMANCE_WORKER_KIND.DATABASE, + (request) => request.threadCpuUserMicros + ), dataFetchMs: summarizeNumbers( measured.map((iteration) => iteration.phases.dataFetchMs) ), @@ -185,22 +210,6 @@ export function createCancellationBenchmarkSummary( iteration.phases.persistencePreparationProxyMs ) ), - playlistWorkerEventLoopUtilization: workerMetric( - PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, - (worker) => worker.eventLoopUtilization - ), - playlistWorkerEventLoopDelayMaxMs: workerEventLoopDelayMetric( - PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, - 'maxMs' - ), - playlistWorkerEventLoopDelayP95Ms: workerEventLoopDelayMetric( - PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, - 'p95Ms' - ), - playlistWorkerEventLoopDelayP99Ms: workerEventLoopDelayMetric( - PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, - 'p99Ms' - ), playlistWorkerExternalPeakBytes: workerMetric( PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, (worker) => worker.peakExternalBytes @@ -213,6 +222,33 @@ export function createCancellationBenchmarkSummary( PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, (worker) => worker.postGcHeapUsedBytes ), + playlistWorkerRequestEventLoopDelayMaxMs: + workerRequestEventLoopDelayMetric( + PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, + 'maxMs' + ), + playlistWorkerRequestEventLoopDelayP95Ms: + workerRequestEventLoopDelayMetric( + PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, + 'p95Ms' + ), + playlistWorkerRequestEventLoopDelayP99Ms: + workerRequestEventLoopDelayMetric( + PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, + 'p99Ms' + ), + playlistWorkerRequestEventLoopUtilization: workerRequestMetric( + PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, + (request) => request.eventLoopUtilization + ), + playlistWorkerRequestThreadCpuSystemMicros: workerRequestMetric( + PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, + (request) => request.threadCpuSystemMicros + ), + playlistWorkerRequestThreadCpuUserMicros: workerRequestMetric( + PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, + (request) => request.threadCpuUserMicros + ), rendererFrameGapMs: summarizeNumbers( flattenRenderer( (iteration) => iteration.renderer.probe.frameGapsMs diff --git a/apps/electron-backend-e2e/src/performance/m3u-refresh-main-capture.ts b/apps/electron-backend-e2e/src/performance/m3u-refresh-main-capture.ts index fd9435445..f61fa12b4 100644 --- a/apps/electron-backend-e2e/src/performance/m3u-refresh-main-capture.ts +++ b/apps/electron-backend-e2e/src/performance/m3u-refresh-main-capture.ts @@ -1,7 +1,17 @@ /* eslint-disable max-lines -- The injected main-process protocol must remain self-contained for Playwright serialization. */ import type { ElectronApplication } from '@playwright/test'; +import { + DATABASE_REQUEST_IDENTITY_CAPTURE_STATE_KEY, + installDatabaseRequestIdentityCapture, + type DatabaseRequestIdentity, + type DatabaseRequestIdentityCaptureApi, +} from './database-request-identity-capture'; import type { MainCaptureMetrics } from './m3u-refresh-cancellation-contract'; +import { + selectMainCaptureGeneration, + type MainCaptureGenerationTransport, +} from './worker-request-performance'; const MAIN_CAPTURE_STATE_KEY = '__iptvnatorM3uRefreshMainCapture'; @@ -23,7 +33,13 @@ export interface MainCaptureStatus { export async function installMainCapture( electronApp: ElectronApplication ): Promise { - await electronApp.evaluate(async ({ app, BrowserWindow }, stateKey) => { + await installDatabaseRequestIdentityCapture(electronApp); + const captureStateKeys = { + databaseRequestIdentityStateKey: + DATABASE_REQUEST_IDENTITY_CAPTURE_STATE_KEY, + stateKey: MAIN_CAPTURE_STATE_KEY, + }; + await electronApp.evaluate(async ({ app, BrowserWindow }, input) => { type JsonRecord = Record; type WorkerTransferable = import('node:worker_threads').Transferable; interface CpuProfileHandle { @@ -42,6 +58,15 @@ export async function installMainCapture( idle: number; utilization: number; } + interface WorkerRequestPerformance { + identity: DatabaseRequestIdentity; + operation: string | null; + performanceCapture: unknown; + playlistId: string | null; + requestId: string | null; + responseEpochMs: number; + success: boolean; + } interface InstrumentedWorker { cpuUsage?(): Promise; getHeapSnapshot?(): Promise; @@ -62,15 +87,11 @@ export async function installMainCapture( } interface WorkerRecord { cancelPostedEpochMs: number | null; + captureGeneration: number | null; cpuFirst: WorkerCpuUsage | null; cpuLast: WorkerCpuUsage | null; elu: number | null; eluStart: WorkerElu | null; - eventLoopDelay: { - maxMs: number; - p95Ms: number; - p99Ms: number; - } | null; externalPeak: number; finalized: boolean; finalizing: Promise | null; @@ -81,6 +102,7 @@ export async function installMainCapture( postGcHeapUsed: number | null; profileHandle: Promise | null; profilePath: string | null; + requestPerformance: WorkerRequestPerformance[]; responseEpochMs: number | null; sampleBusy: boolean; sampleTimer: NodeJS.Timeout | null; @@ -99,9 +121,12 @@ export async function installMainCapture( } const target = globalThis as unknown as Record; - if (target[stateKey] !== undefined) { + if (target[input.stateKey] !== undefined) { return; } + const databaseRequestIdentityCapture = target[ + input.databaseRequestIdentityStateKey + ] as DatabaseRequestIdentityCaptureApi; const runtimeProcess = process as typeof process & { getBuiltinModule(id: string): unknown; @@ -132,6 +157,7 @@ export async function installMainCapture( const dbRequests = new Map< string, { + identity: DatabaseRequestIdentity; operation: string; playlistId: string | null; record: WorkerRecord; @@ -140,6 +166,7 @@ export async function installMainCapture( const state = { active: false, + captureGeneration: 0, cpuStart: null as NodeJS.CpuUsage | null, diagnostic: false, eventLoopDelay: null as ReturnType< @@ -200,6 +227,9 @@ export async function installMainCapture( const playlistId = value['playlistId'] ?? value['_id']; return typeof playlistId === 'string' ? playlistId : null; }; + const isCurrentCaptureRecord = (record: WorkerRecord): boolean => + state.active && + record.captureGeneration === state.captureGeneration; const createWorkerRecord = ( worker: InstrumentedWorker, kind: WorkerRecord['kind'] @@ -210,11 +240,11 @@ export async function installMainCapture( } const record: WorkerRecord = { cancelPostedEpochMs: null, + captureGeneration: null, cpuFirst: null, cpuLast: null, elu: null, eluStart: null, - eventLoopDelay: null, externalPeak: 0, finalized: false, finalizing: null, @@ -225,6 +255,7 @@ export async function installMainCapture( postGcHeapUsed: null, profileHandle: null, profilePath: null, + requestPerformance: [], responseEpochMs: null, sampleBusy: false, sampleTimer: null, @@ -238,30 +269,13 @@ export async function installMainCapture( return; } const message = incoming as JsonRecord; - const performanceCapture = readWorkerPerformance(message); - if (performanceCapture) { - record.eventLoopDelay = record.eventLoopDelay - ? { - maxMs: Math.max( - record.eventLoopDelay.maxMs, - performanceCapture.eventLoopDelay.maxMs - ), - p95Ms: Math.max( - record.eventLoopDelay.p95Ms, - performanceCapture.eventLoopDelay.p95Ms - ), - p99Ms: Math.max( - record.eventLoopDelay.p99Ms, - performanceCapture.eventLoopDelay.p99Ms - ), - } - : performanceCapture.eventLoopDelay; + if (!isCurrentCaptureRecord(record)) { + return; } if (record.kind === 'playlist-refresh.worker') { if (message['type'] === 'event') { const event = message['event'] as - | JsonRecord - | undefined; + JsonRecord | undefined; recordTimeline({ operationId: record.operationId ?? undefined, playlistId: record.playlistId ?? undefined, @@ -270,7 +284,20 @@ export async function installMainCapture( )}:${String(event?.['phase'] ?? 'none')}`, }); } else if (message['type'] === 'response') { - record.responseEpochMs = nowEpochMs(); + const responseEpochMs = nowEpochMs(); + record.responseEpochMs = responseEpochMs; + record.requestPerformance.push({ + identity: { + operationId: record.operationId, + operationIdUnavailableReason: null, + }, + operation: null, + performanceCapture: message['performance'] ?? null, + playlistId: record.playlistId, + requestId: null, + responseEpochMs, + success: message['success'] === true, + }); recordTimeline({ operationId: record.operationId ?? undefined, playlistId: record.playlistId ?? undefined, @@ -285,10 +312,23 @@ export async function installMainCapture( typeof message['requestId'] === 'string' ) { const request = dbRequests.get(message['requestId']); + if (request) { + record.requestPerformance.push({ + identity: request.identity, + operation: request.operation, + performanceCapture: message['performance'] ?? null, + playlistId: request.playlistId, + requestId: message['requestId'], + responseEpochMs: nowEpochMs(), + success: message['success'] === true, + }); + } if (request) { dbRequests.delete(message['requestId']); recordTimeline({ operation: request.operation, + operationId: + request.identity.operationId ?? undefined, playlistId: request.playlistId ?? undefined, requestId: message['requestId'], success: message['success'] === true, @@ -299,46 +339,6 @@ export async function installMainCapture( }); return record; }; - const readWorkerPerformance = ( - message: JsonRecord - ): { - eventLoopDelay: { - maxMs: number; - p95Ms: number; - p99Ms: number; - }; - } | null => { - const performanceCapture = message['performance']; - if ( - typeof performanceCapture !== 'object' || - performanceCapture === null - ) { - return null; - } - const eventLoopDelay = (performanceCapture as JsonRecord)[ - 'eventLoopDelay' - ]; - if (typeof eventLoopDelay !== 'object' || eventLoopDelay === null) { - return null; - } - const value = eventLoopDelay as JsonRecord; - const maxMs = value['maxMs']; - const p95Ms = value['p95Ms']; - const p99Ms = value['p99Ms']; - if ( - typeof maxMs !== 'number' || - !Number.isFinite(maxMs) || - typeof p95Ms !== 'number' || - !Number.isFinite(p95Ms) || - typeof p99Ms !== 'number' || - !Number.isFinite(p99Ms) - ) { - return null; - } - return { - eventLoopDelay: { maxMs, p95Ms, p99Ms }, - }; - }; const sampleWorker = async (record: WorkerRecord): Promise => { if (record.sampleBusy || record.finalized) { return; @@ -373,14 +373,45 @@ export async function installMainCapture( record.sampleBusy = false; } }; + const resetWorkerForCapture = (record: WorkerRecord): void => { + if (record.sampleTimer) { + clearInterval(record.sampleTimer); + } + record.cancelPostedEpochMs = null; + record.captureGeneration = state.captureGeneration; + record.cpuFirst = null; + record.cpuLast = null; + record.elu = null; + record.eluStart = null; + record.externalPeak = 0; + record.finalized = false; + record.finalizing = null; + record.heapPeak = 0; + record.operationId = null; + record.playlistId = null; + record.postGcHeapUsed = null; + record.profileHandle = null; + record.profilePath = null; + record.requestPerformance = []; + record.responseEpochMs = null; + record.sampleBusy = false; + record.sampleTimer = null; + record.snapshotPath = null; + record.terminatedEpochMs = null; + }; const startWorker = (record: WorkerRecord): void => { - if (!state.active || record.sampleTimer !== null) { + if (!state.active) { + return; + } + if (record.captureGeneration !== state.captureGeneration) { + resetWorkerForCapture(record); + } + if (record.sampleTimer !== null) { return; } record.elu = null; record.eluStart = record.worker.performance?.eventLoopUtilization() ?? null; - record.eventLoopDelay = null; void sampleWorker(record); record.sampleTimer = setInterval( () => void sampleWorker(record), @@ -456,7 +487,11 @@ export async function installMainCapture( const kind = classifyRequest(value); if (kind) { const record = createWorkerRecord(this, kind); - if (kind === 'playlist-refresh.worker') { + startWorker(record); + if ( + isCurrentCaptureRecord(record) && + kind === 'playlist-refresh.worker' + ) { const payload = value['payload'] as JsonRecord; record.operationId = String(payload['operationId']); record.playlistId = String(payload['playlistId']); @@ -467,23 +502,47 @@ export async function installMainCapture( type: 'playlist-request', }); } else if ( + isCurrentCaptureRecord(record) && typeof value['requestId'] === 'string' && typeof value['operation'] === 'string' ) { + const payload = value['payload']; + const payloadOperationId = + typeof payload === 'object' && + payload !== null && + typeof (payload as JsonRecord)[ + 'operationId' + ] === 'string' && + String((payload as JsonRecord)['operationId']) + .length > 0 + ? String( + (payload as JsonRecord)['operationId'] + ) + : null; + const identity = + payloadOperationId === null + ? databaseRequestIdentityCapture.matchDatabaseRequest( + value + ) + : { + operationId: payloadOperationId, + operationIdUnavailableReason: null, + }; dbRequests.set(value['requestId'], { + identity, operation: value['operation'], playlistId: readDatabasePlaylistId(value), record, }); recordTimeline({ operation: value['operation'], + operationId: identity.operationId ?? undefined, playlistId: readDatabasePlaylistId(value) ?? undefined, requestId: value['requestId'], type: 'db-request', }); } - startWorker(record); } } else if ( value['type'] === 'cancel' && @@ -505,7 +564,7 @@ export async function installMainCapture( this: InstrumentedWorker ): Promise { const record = records.get(this); - if (!record || !state.active) { + if (!record || !isCurrentCaptureRecord(record)) { return originalTerminate.call(this); } const markTerminated = (code: number): number => { @@ -627,6 +686,7 @@ export async function installMainCapture( ).length, playlistTerminated: [...records.values()].filter( (record) => + isCurrentCaptureRecord(record) && record.kind === 'playlist-refresh.worker' && record.terminatedEpochMs !== null ).length, @@ -645,7 +705,11 @@ export async function installMainCapture( ).length, }), start: async (options: MainCaptureStartOptions): Promise => { + state.captureGeneration += 1; state.active = true; + databaseRequestIdentityCapture.start(); + dbRequests.clear(); + operationWorkers.clear(); state.diagnostic = options.diagnostic; state.outputDirectory = options.outputDirectory; state.timeline = []; @@ -682,7 +746,7 @@ export async function installMainCapture( ); } }, - stop: async (): Promise => { + stop: async (): Promise => { if (state.sampleTimer) { clearInterval(state.sampleTimer); state.sampleTimer = null; @@ -705,6 +769,7 @@ export async function installMainCapture( const delay = state.eventLoopDelay; const databaseRecords = [...records.values()].filter( (record) => + record.captureGeneration === state.captureGeneration && record.kind === 'database.worker' && record.sampleTimer !== null ); @@ -756,78 +821,91 @@ export async function installMainCapture( } state.windowListeners = []; state.active = false; + databaseRequestIdentityCapture.stop(); const workers = [...records.values()] .filter( (record) => record.kind === 'playlist-refresh.worker' || - record.heapPeak > 0 + record.heapPeak > 0 || + record.requestPerformance.length > 0 ) .map((record) => ({ - cancelPostedEpochMs: record.cancelPostedEpochMs, - cpuSystemMicros: - record.cpuFirst && record.cpuLast - ? record.cpuLast.system - record.cpuFirst.system - : null, - cpuUserMicros: - record.cpuFirst && record.cpuLast - ? record.cpuLast.user - record.cpuFirst.user - : null, - eventLoopDelay: record.eventLoopDelay, - eventLoopDelayUnavailableReason: - record.eventLoopDelay === null - ? record.terminatedEpochMs !== null - ? 'worker-terminated-before-profile-flush' - : 'worker-self-profile-result-unavailable' - : null, - eventLoopUtilization: record.elu, - kind: record.kind, - operationId: record.operationId, - peakExternalBytes: record.externalPeak, - peakHeapUsedBytes: record.heapPeak, - playlistId: record.playlistId, - postGcHeapUsedBytes: record.postGcHeapUsed, - profilePath: record.profilePath, - responseEpochMs: record.responseEpochMs, - snapshotPath: record.snapshotPath, - terminatedEpochMs: record.terminatedEpochMs, + captureGeneration: record.captureGeneration, + metrics: { + cancelPostedEpochMs: record.cancelPostedEpochMs, + cpuSystemMicros: + record.cpuFirst && record.cpuLast + ? record.cpuLast.system - + record.cpuFirst.system + : null, + cpuUserMicros: + record.cpuFirst && record.cpuLast + ? record.cpuLast.user - record.cpuFirst.user + : null, + eventLoopUtilization: record.elu, + kind: record.kind, + operationId: record.operationId, + peakExternalBytes: record.externalPeak, + peakHeapUsedBytes: record.heapPeak, + playlistId: record.playlistId, + postGcHeapUsedBytes: record.postGcHeapUsed, + profilePath: record.profilePath, + responseEpochMs: record.responseEpochMs, + snapshotPath: record.snapshotPath, + terminatedEpochMs: record.terminatedEpochMs, + }, + requests: record.requestPerformance.map((request) => ({ + operation: request.operation, + operationId: request.identity.operationId, + operationIdUnavailableReason: + request.identity.operationIdUnavailableReason, + performanceCapture: request.performanceCapture, + playlistId: request.playlistId, + requestId: request.requestId, + responseEpochMs: request.responseEpochMs, + success: request.success, + })), })); return { - cpuProfilePath: state.mainProfilePath, - cpuSystemMicros: cpu.system, - cpuUserMicros: cpu.user, - eventLoopDelay: { - maxMs: Number(delay?.max ?? 0) / 1e6, - p95Ms: Number(delay?.percentile(95) ?? 0) / 1e6, - p99Ms: Number(delay?.percentile(99) ?? 0) / 1e6, + captureGeneration: state.captureGeneration, + metrics: { + cpuProfilePath: state.mainProfilePath, + cpuSystemMicros: cpu.system, + cpuUserMicros: cpu.user, + eventLoopDelay: { + maxMs: Number(delay?.max ?? 0) / 1e6, + p95Ms: Number(delay?.percentile(95) ?? 0) / 1e6, + p99Ms: Number(delay?.percentile(99) ?? 0) / 1e6, + }, + // Electron's Chromium-owned main message pump currently + // leaves Node's libuv ELU counters at zero. Preserve that + // as unavailable instead of reporting a misleading 0%. + eventLoopUtilization, + eventLoopUtilizationUnavailableReason: + eventLoopUtilization === null + ? 'electron-main-embedded-event-loop' + : null, + heapSnapshotPath: state.mainSnapshotPath, + memory: { + peakHeapUsedBytes: state.mainPeakHeap, + peakRssBytes: state.mainPeakRss, + postGcHeapUsedBytes: state.postGcHeap, + postGcRssBytes: state.postGcRss, + }, + rendererPeakRssBytes: state.rendererPeakRss, + responsiveEvents: state.responsiveEvents, + rssScope: + 'electron-main-process-including-worker-threads-and-native-memory', + timeline: state.timeline, + unresponsiveEvents: state.unresponsiveEvents, }, - // Electron's Chromium-owned main message pump currently - // leaves Node's libuv ELU counters at zero. Preserve that - // as unavailable instead of reporting a misleading 0%. - eventLoopUtilization, - eventLoopUtilizationUnavailableReason: - eventLoopUtilization === null - ? 'electron-main-embedded-event-loop' - : null, - heapSnapshotPath: state.mainSnapshotPath, - memory: { - peakHeapUsedBytes: state.mainPeakHeap, - peakRssBytes: state.mainPeakRss, - postGcHeapUsedBytes: state.postGcHeap, - postGcRssBytes: state.postGcRss, - }, - rendererPeakRssBytes: state.rendererPeakRss, - responsiveEvents: state.responsiveEvents, - rssScope: - 'electron-main-process-including-worker-threads-and-native-memory', - timeline: state.timeline, - unresponsiveEvents: state.unresponsiveEvents, workers, }; }, }; - target[stateKey] = api; - }, MAIN_CAPTURE_STATE_KEY); + target[input.stateKey] = api; + }, captureStateKeys); } export async function startMainCapture( @@ -861,11 +939,15 @@ export async function readMainCaptureStatus( export async function stopMainCapture( electronApp: ElectronApplication ): Promise { - return electronApp.evaluate(async (_electron, stateKey) => { - const target = globalThis as unknown as Record; - const api = target[stateKey] as { - stop(): Promise; - }; - return api.stop(); - }, MAIN_CAPTURE_STATE_KEY); + const transport = await electronApp.evaluate( + async (_electron, stateKey) => { + const target = globalThis as unknown as Record; + const api = target[stateKey] as { + stop(): Promise; + }; + return api.stop(); + }, + MAIN_CAPTURE_STATE_KEY + ); + return selectMainCaptureGeneration(transport); } diff --git a/apps/electron-backend-e2e/src/performance/worker-request-performance-envelope.spec.ts b/apps/electron-backend-e2e/src/performance/worker-request-performance-envelope.spec.ts new file mode 100644 index 000000000..e87e3c345 --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/worker-request-performance-envelope.spec.ts @@ -0,0 +1,172 @@ +/* eslint-disable playwright/expect-expect -- These are Node assertion-based performance contract tests. */ +import assert from 'node:assert/strict'; +import test from 'node:test'; + +import type { WorkerRequestPerformanceMetrics } from './m3u-refresh-cancellation-contract'; +import { normalizeWorkerRequestPerformanceOutcome } from './worker-request-performance'; + +const validCapture = { + eventLoopDelay: { maxMs: 6, p95Ms: 5, p99Ms: 5.5 }, + eventLoopDelayUnavailableReason: null, + eventLoopUtilization: 0.5, + eventLoopUtilizationUnavailableReason: null, + histogramFlushedEpochMs: 1_050, + invalidReason: null, + requestReceivedEpochMs: 1_000, + threadCpuSystemMicros: 50, + threadCpuUnavailableReason: null, + threadCpuUserMicros: 100, + workEndedEpochMs: 1_040, + workStartedEpochMs: 1_010, +} as const; + +function normalize(performanceCapture: unknown) { + return normalizeWorkerRequestPerformanceOutcome({ + operation: 'DB_GET_APP_PLAYLIST', + operationId: 'operation-42', + operationIdUnavailableReason: null, + performanceCapture, + playlistId: 'playlist-1', + requestId: 'request-1', + responseEpochMs: 1_060, + success: true, + }); +} + +function assertClosedCapture( + capture: WorkerRequestPerformanceMetrics, + reason: string +): void { + assert.equal(capture.performanceCaptureUnavailableReason, reason); + assert.equal(capture.operationId, 'operation-42'); + for (const value of [ + capture.eventLoopDelay, + capture.eventLoopDelayUnavailableReason, + capture.eventLoopUtilization, + capture.eventLoopUtilizationUnavailableReason, + capture.histogramFlushedEpochMs, + capture.invalidReason, + capture.requestReceivedEpochMs, + capture.threadCpuSystemMicros, + capture.threadCpuUnavailableReason, + capture.threadCpuUserMicros, + capture.workEndedEpochMs, + capture.workStartedEpochMs, + ]) { + assert.equal(value, null); + } +} + +test('absent captures are missing while present non-record captures are invalid', () => { + for (const missing of [undefined, null]) { + assertClosedCapture( + normalize(missing), + 'worker-performance-capture-missing' + ); + } + for (const malformed of [42, [], 'capture']) { + assertClosedCapture( + normalize(malformed), + 'worker-performance-capture-invalid' + ); + } +}); + +test('malformed nested fields and cross-field invariants fail the whole capture closed', () => { + const malformed = [ + { ...validCapture, eventLoopDelay: [] }, + { + ...validCapture, + eventLoopDelay: { maxMs: 6, p95Ms: 'bad', p99Ms: 5.5 }, + }, + { + ...validCapture, + eventLoopDelay: { maxMs: 5, p95Ms: 6, p99Ms: 5.5 }, + }, + { ...validCapture, eventLoopDelayUnavailableReason: 42 }, + { + ...validCapture, + eventLoopDelay: null, + eventLoopDelayUnavailableReason: null, + histogramFlushedEpochMs: null, + }, + { + ...validCapture, + eventLoopDelay: null, + eventLoopDelayUnavailableReason: + 'event-loop-delay-capture-unavailable', + }, + { ...validCapture, eventLoopUtilization: 1.01 }, + { + ...validCapture, + eventLoopUtilization: null, + eventLoopUtilizationUnavailableReason: null, + }, + { ...validCapture, threadCpuSystemMicros: null }, + { + ...validCapture, + threadCpuSystemMicros: null, + threadCpuUnavailableReason: null, + threadCpuUserMicros: null, + }, + { ...validCapture, threadCpuUnavailableReason: 'unavailable' }, + { ...validCapture, histogramFlushedEpochMs: null }, + { + ...validCapture, + eventLoopDelay: null, + eventLoopDelayUnavailableReason: 'event-loop-delay-invalid', + histogramFlushedEpochMs: null, + }, + { ...validCapture, histogramFlushedEpochMs: 1_039 }, + { ...validCapture, workStartedEpochMs: 999 }, + { + ...validCapture, + invalidReason: 'overlapping-database-worker-requests', + }, + { ...validCapture, threadCpuUserMicros: undefined }, + ]; + for (const performanceCapture of malformed) { + assertClosedCapture( + normalize(performanceCapture), + 'worker-performance-capture-invalid' + ); + } +}); + +test('coherent metric-unavailable and invalid worker results remain valid raw captures', () => { + const unavailable = { + ...validCapture, + eventLoopDelay: null, + eventLoopDelayUnavailableReason: 'event-loop-delay-capture-unavailable', + eventLoopUtilization: null, + eventLoopUtilizationUnavailableReason: + 'event-loop-utilization-unavailable', + histogramFlushedEpochMs: null, + threadCpuSystemMicros: null, + threadCpuUnavailableReason: 'thread-cpu-usage-unavailable', + threadCpuUserMicros: null, + }; + for (const capture of [ + unavailable, + { + ...validCapture, + eventLoopDelay: null, + eventLoopDelayUnavailableReason: 'event-loop-delay-invalid', + }, + ]) { + const accepted = normalize(capture); + assert.equal(accepted.performanceCaptureUnavailableReason, null); + assert.equal(accepted.requestReceivedEpochMs, 1_000); + } + + const invalidReason = 'overlapping-database-worker-requests'; + const legitimateInvalid = normalize({ + ...unavailable, + eventLoopDelayUnavailableReason: invalidReason, + eventLoopUtilizationUnavailableReason: invalidReason, + invalidReason, + threadCpuUnavailableReason: invalidReason, + }); + assert.equal(legitimateInvalid.performanceCaptureUnavailableReason, null); + assert.equal(legitimateInvalid.invalidReason, invalidReason); +}); diff --git a/apps/electron-backend-e2e/src/performance/worker-request-performance.spec.ts b/apps/electron-backend-e2e/src/performance/worker-request-performance.spec.ts new file mode 100644 index 000000000..72a3ae522 --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/worker-request-performance.spec.ts @@ -0,0 +1,302 @@ +/* eslint-disable playwright/expect-expect -- These are Node assertion-based performance contract tests. */ +import assert from 'node:assert/strict'; +import { readFileSync } from 'node:fs'; +import test from 'node:test'; + +import { + type CancellationBenchmarkManifest, + type CancellationIterationResult, + PERFORMANCE_ITERATION_KIND, + PERFORMANCE_WORKER_KIND, + type WorkerCaptureMetrics, + type WorkerRequestPerformanceMetrics, +} from './m3u-refresh-cancellation-contract'; +import { + normalizeWorkerRequestPerformanceOutcome, + selectMainCaptureGeneration, +} from './worker-request-performance'; +import { createCancellationBenchmarkSummary } from './m3u-refresh-cancellation-report'; + +type RequestPerformanceFixture = WorkerRequestPerformanceMetrics; +function requestPerformance( + requestId: string, + p95Ms: number, + threadCpuUserMicros: number +): RequestPerformanceFixture { + return normalizeWorkerRequestPerformanceOutcome( + requestTransport( + { + eventLoopDelay: { + maxMs: p95Ms + 1, + p95Ms, + p99Ms: p95Ms + 0.5, + }, + eventLoopDelayUnavailableReason: null, + eventLoopUtilization: p95Ms / 100, + eventLoopUtilizationUnavailableReason: null, + histogramFlushedEpochMs: 1_050 + p95Ms, + invalidReason: null, + requestReceivedEpochMs: 1_000, + threadCpuSystemMicros: threadCpuUserMicros / 2, + threadCpuUnavailableReason: null, + threadCpuUserMicros, + workEndedEpochMs: 1_040, + workStartedEpochMs: 1_010, + }, + requestId + ) + ); +} +function requestTransport( + performanceCapture: unknown, + requestId = 'request-1' +) { + return { + operation: 'DB_GET_APP_PLAYLIST', + operationId: 'operation-42', + operationIdUnavailableReason: null, + performanceCapture, + playlistId: 'playlist-1', + requestId, + responseEpochMs: 1_060, + success: true, + } as const; +} +function assertClosedCapture( + capture: WorkerRequestPerformanceMetrics, + reason: string +): void { + assert.equal(capture.performanceCaptureUnavailableReason, reason); + assert.equal(capture.operationId, 'operation-42'); + for (const value of [ + capture.eventLoopDelay, + capture.eventLoopDelayUnavailableReason, + capture.eventLoopUtilization, + capture.eventLoopUtilizationUnavailableReason, + capture.histogramFlushedEpochMs, + capture.invalidReason, + capture.requestReceivedEpochMs, + capture.threadCpuSystemMicros, + capture.threadCpuUnavailableReason, + capture.threadCpuUserMicros, + capture.workEndedEpochMs, + capture.workStartedEpochMs, + ]) { + assert.equal(value, null); + } +} +function databaseWorker( + requests: readonly RequestPerformanceFixture[] +): WorkerCaptureMetrics { + return { + kind: PERFORMANCE_WORKER_KIND.DATABASE, + operationId: null, + peakExternalBytes: 0, + peakHeapUsedBytes: 0, + playlistId: null, + postGcHeapUsedBytes: null, + requests, + responseEpochMs: null, + terminatedEpochMs: null, + } as WorkerCaptureMetrics; +} +function playlistWorker( + request: RequestPerformanceFixture, + operationId: string, + terminatedEpochMs: number +): WorkerCaptureMetrics { + return { + ...databaseWorker([request]), + kind: PERFORMANCE_WORKER_KIND.PLAYLIST_REFRESH, + operationId, + playlistId: request.playlistId, + responseEpochMs: request.responseEpochMs, + terminatedEpochMs, + }; +} + +function measuredIteration( + workers: readonly WorkerCaptureMetrics[] +): CancellationIterationResult { + return { + cancellationEffectObserved: true, + kind: PERFORMANCE_ITERATION_KIND.MEASURED, + main: { + eventLoopDelay: { maxMs: 0, p95Ms: 0, p99Ms: 0 }, + eventLoopUtilization: null, + memory: { + peakHeapUsedBytes: 0, + peakRssBytes: 0, + postGcHeapUsedBytes: null, + postGcRssBytes: null, + }, + rendererPeakRssBytes: 0, + responsiveEvents: 0, + unresponsiveEvents: 0, + workers, + } as CancellationIterationResult['main'], + phases: {} as CancellationIterationResult['phases'], + renderer: { + peakHeapUsedBytes: 0, + postGcHeapUsedBytes: null, + probe: { + frameGapsMs: [], + heartbeatDelaysMs: [], + longTasksMs: [], + }, + } as CancellationIterationResult['renderer'], + runId: 'measured-1', + }; +} + +test('summary distributions use every request capture instead of a worker-level percentile maximum', () => { + const requests = [ + requestPerformance('request-1', 5, 100), + requestPerformance('request-2', 15, 300), + ]; + const summary = createCancellationBenchmarkSummary( + {} as CancellationBenchmarkManifest, + [measuredIteration([databaseWorker(requests)])] + ); + const measured = summary.measured as unknown as Record< + string, + { count: number; max: number | null; min: number | null } + >; + + assert.deepEqual( + summary.measured.databaseWorkerRequestEventLoopDelayP95Ms, + { + count: 2, + max: 15, + mean: 10, + median: 10, + min: 5, + p95: 14.5, + p99: 14.9, + } + ); + assert.deepEqual(measured['databaseWorkerRequestThreadCpuUserMicros'], { + count: 2, + max: 300, + mean: 200, + median: 200, + min: 100, + p95: 290, + p99: 298, + }); + assert.equal( + summary.iterations[0]?.main.workers[0]?.requests[0]?.operationId, + 'operation-42' + ); + assert.equal( + 'databaseWorkerEventLoopDelayP95Ms' in summary.measured, + false, + 'request-percentile distributions must not look like an operation-wide percentile' + ); +}); + +test('main capture retains raw request identity, timestamps, metrics, and reasons', () => { + const source = readFileSync( + new URL('./m3u-refresh-main-capture.ts', import.meta.url), + 'utf8' + ); + + assert.match(source, /requestPerformance: \[\]/); + assert.match(source, /record\.requestPerformance\.length > 0/); + assert.match(source, /record\.requestPerformance\.map/); + assert.match(source, /operationId: request\.identity\.operationId/); + assert.match(source, /identity: request\.identity/); + assert.match(source, /performanceCapture:\s+message\['performance'\]/); + assert.match(source, /captureGeneration: state\.captureGeneration/); + assert.doesNotMatch( + source, + /record\.eventLoopDelay = record\.eventLoopDelay/ + ); +}); + +test('capture-generation selection excludes a pre-start seed worker', () => { + const seed = playlistWorker( + requestPerformance('seed-request', 999, 999), + 'seed-operation', + 1_100 + ); + const measured = playlistWorker( + { + ...requestPerformance('measured-request', 8, 80), + success: false, + }, + 'measured-cancel-operation', + 2_100 + ); + + const mainMetrics = measuredIteration([]).main; + const toTransportWorker = ( + captureGeneration: number | null, + worker: WorkerCaptureMetrics + ) => ({ + captureGeneration, + metrics: worker, + requests: worker.requests.map((request) => ({ + ...requestTransport(request, request.requestId ?? undefined), + success: request.success, + })), + }); + const selected = selectMainCaptureGeneration({ + captureGeneration: 2, + metrics: mainMetrics, + workers: [ + toTransportWorker(null, seed), + toTransportWorker(2, measured), + ], + }); + + assert.equal(selected.workers.length, 1); + assert.equal( + selected.workers[0]?.requests[0]?.requestId, + 'measured-request' + ); + assert.equal(selected.workers[0]?.operationId, 'measured-cancel-operation'); + assert.equal(selected.workers[0]?.terminatedEpochMs, 2_100); +}); + +test('raw outcomes retain missing and malformed captures while summaries exclude their null metrics', () => { + const valid = normalizeWorkerRequestPerformanceOutcome( + requestTransport(requestPerformance('request-1', 5, 100)) + ); + const missing = normalizeWorkerRequestPerformanceOutcome( + requestTransport(null, 'request-2') + ); + const malformed = normalizeWorkerRequestPerformanceOutcome( + requestTransport( + { + requestReceivedEpochMs: 'invalid', + workEndedEpochMs: 1_080, + workStartedEpochMs: 1_075, + }, + 'request-3' + ) + ); + const worker = databaseWorker([valid, missing, malformed]); + const summary = createCancellationBenchmarkSummary( + {} as CancellationBenchmarkManifest, + [measuredIteration([worker])] + ); + + assert.equal(worker.requests.length, 3); + assertClosedCapture( + worker.requests[1] as WorkerRequestPerformanceMetrics, + 'worker-performance-capture-missing' + ); + assertClosedCapture( + worker.requests[2] as WorkerRequestPerformanceMetrics, + 'worker-performance-capture-invalid' + ); + assert.equal( + summary.measured.databaseWorkerRequestEventLoopDelayP95Ms.count, + 1 + ); + assert.equal( + summary.measured.databaseWorkerRequestThreadCpuUserMicros.count, + 1 + ); +}); diff --git a/apps/electron-backend-e2e/src/performance/worker-request-performance.ts b/apps/electron-backend-e2e/src/performance/worker-request-performance.ts new file mode 100644 index 000000000..edf157282 --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/worker-request-performance.ts @@ -0,0 +1,302 @@ +import type { + EventLoopDelayMetrics, + MainCaptureMetrics, + WorkerCaptureMetrics, + WorkerRequestPerformanceMetrics, +} from './m3u-refresh-cancellation-contract'; + +export interface WorkerRequestPerformanceOutcomeTransport { + readonly operation: string | null; + readonly operationId: string | null; + readonly operationIdUnavailableReason: string | null; + readonly performanceCapture: unknown; + readonly playlistId: string | null; + readonly requestId: string | null; + readonly responseEpochMs: number; + readonly success: boolean; +} + +export interface MainCaptureGenerationTransport { + readonly captureGeneration: number; + readonly metrics: Omit; + readonly workers: readonly { + readonly captureGeneration: number | null; + readonly metrics: Omit< + WorkerCaptureMetrics, + 'eventLoopDelay' | 'eventLoopDelayUnavailableReason' | 'requests' + >; + readonly requests: readonly WorkerRequestPerformanceOutcomeTransport[]; + }[]; +} + +type ParsedWorkerPerformanceCapture = Pick< + WorkerRequestPerformanceMetrics, + | 'eventLoopDelay' + | 'eventLoopDelayUnavailableReason' + | 'eventLoopUtilization' + | 'eventLoopUtilizationUnavailableReason' + | 'histogramFlushedEpochMs' + | 'invalidReason' + | 'requestReceivedEpochMs' + | 'threadCpuSystemMicros' + | 'threadCpuUnavailableReason' + | 'threadCpuUserMicros' + | 'workEndedEpochMs' + | 'workStartedEpochMs' +>; + +const INVALID_REASONS = new Set(['overlapping-database-worker-requests']); +const EVENT_LOOP_DELAY_REASONS = new Set([ + ...INVALID_REASONS, + 'event-loop-delay-arm-timeout', + 'event-loop-delay-capture-unavailable', + 'event-loop-delay-flush-timeout', + 'event-loop-delay-invalid', +]); +const EVENT_LOOP_UTILIZATION_REASONS = new Set([ + ...INVALID_REASONS, + 'event-loop-utilization-unavailable', +]); +const THREAD_CPU_REASONS = new Set([ + ...INVALID_REASONS, + 'thread-cpu-usage-invalid', + 'thread-cpu-usage-unavailable', +]); + +function isRecord(input: unknown): input is Record { + return typeof input === 'object' && input !== null && !Array.isArray(input); +} + +function isFiniteNonNegativeNumber(value: unknown): value is number { + return typeof value === 'number' && Number.isFinite(value) && value >= 0; +} + +function isNullableFiniteNonNegativeNumber( + value: unknown +): value is number | null { + return value === null || isFiniteNonNegativeNumber(value); +} + +function isNullableReason( + value: unknown, + allowed: ReadonlySet +): value is string | null { + return value === null || (typeof value === 'string' && allowed.has(value)); +} + +function parseEventLoopDelay( + input: unknown +): EventLoopDelayMetrics | null | undefined { + if (input === null) { + return null; + } + if (!isRecord(input)) { + return undefined; + } + const maxMs = input['maxMs']; + const p95Ms = input['p95Ms']; + const p99Ms = input['p99Ms']; + if ( + !isFiniteNonNegativeNumber(maxMs) || + !isFiniteNonNegativeNumber(p95Ms) || + !isFiniteNonNegativeNumber(p99Ms) || + p95Ms > p99Ms || + p99Ms > maxMs + ) { + return undefined; + } + return { maxMs, p95Ms, p99Ms }; +} + +function parseWorkerPerformanceCapture( + input: unknown +): ParsedWorkerPerformanceCapture | null { + if (!isRecord(input)) { + return null; + } + + const eventLoopDelay = parseEventLoopDelay(input['eventLoopDelay']); + const eventLoopDelayUnavailableReason = + input['eventLoopDelayUnavailableReason']; + const eventLoopUtilization = input['eventLoopUtilization']; + const eventLoopUtilizationUnavailableReason = + input['eventLoopUtilizationUnavailableReason']; + const histogramFlushedEpochMs = input['histogramFlushedEpochMs']; + const invalidReason = input['invalidReason']; + const requestReceivedEpochMs = input['requestReceivedEpochMs']; + const threadCpuSystemMicros = input['threadCpuSystemMicros']; + const threadCpuUnavailableReason = input['threadCpuUnavailableReason']; + const threadCpuUserMicros = input['threadCpuUserMicros']; + const workEndedEpochMs = input['workEndedEpochMs']; + const workStartedEpochMs = input['workStartedEpochMs']; + + if ( + eventLoopDelay === undefined || + !isNullableReason( + eventLoopDelayUnavailableReason, + EVENT_LOOP_DELAY_REASONS + ) || + !isNullableFiniteNonNegativeNumber(eventLoopUtilization) || + (eventLoopUtilization !== null && eventLoopUtilization > 1) || + !isNullableReason( + eventLoopUtilizationUnavailableReason, + EVENT_LOOP_UTILIZATION_REASONS + ) || + !isNullableFiniteNonNegativeNumber(histogramFlushedEpochMs) || + !isNullableReason(invalidReason, INVALID_REASONS) || + !isFiniteNonNegativeNumber(requestReceivedEpochMs) || + !isNullableFiniteNonNegativeNumber(threadCpuSystemMicros) || + !isNullableReason(threadCpuUnavailableReason, THREAD_CPU_REASONS) || + !isNullableFiniteNonNegativeNumber(threadCpuUserMicros) || + !isFiniteNonNegativeNumber(workEndedEpochMs) || + !isFiniteNonNegativeNumber(workStartedEpochMs) + ) { + return null; + } + + const timestampsAreOrdered = + requestReceivedEpochMs <= workStartedEpochMs && + workStartedEpochMs <= workEndedEpochMs && + (histogramFlushedEpochMs === null || + workEndedEpochMs <= histogramFlushedEpochMs); + const eventLoopDelayIsCoherent = + eventLoopDelay === null + ? eventLoopDelayUnavailableReason !== null && + (eventLoopDelayUnavailableReason === 'event-loop-delay-invalid' + ? histogramFlushedEpochMs !== null + : histogramFlushedEpochMs === null) + : eventLoopDelayUnavailableReason === null && + histogramFlushedEpochMs !== null; + const eventLoopUtilizationIsCoherent = + eventLoopUtilization === null + ? eventLoopUtilizationUnavailableReason !== null + : eventLoopUtilizationUnavailableReason === null; + const threadCpuIsCoherent = + threadCpuSystemMicros === null && threadCpuUserMicros === null + ? threadCpuUnavailableReason !== null + : threadCpuSystemMicros !== null && + threadCpuUserMicros !== null && + threadCpuUnavailableReason === null; + const invalidCaptureIsCoherent = + invalidReason === null || + (eventLoopDelay === null && + eventLoopDelayUnavailableReason === invalidReason && + eventLoopUtilization === null && + eventLoopUtilizationUnavailableReason === invalidReason && + histogramFlushedEpochMs === null && + threadCpuSystemMicros === null && + threadCpuUnavailableReason === invalidReason && + threadCpuUserMicros === null); + if ( + !timestampsAreOrdered || + !eventLoopDelayIsCoherent || + !eventLoopUtilizationIsCoherent || + !threadCpuIsCoherent || + !invalidCaptureIsCoherent + ) { + return null; + } + + return { + eventLoopDelay, + eventLoopDelayUnavailableReason, + eventLoopUtilization, + eventLoopUtilizationUnavailableReason, + histogramFlushedEpochMs, + invalidReason, + requestReceivedEpochMs, + threadCpuSystemMicros, + threadCpuUnavailableReason, + threadCpuUserMicros, + workEndedEpochMs, + workStartedEpochMs, + }; +} + +export function normalizeWorkerRequestPerformanceOutcome( + input: WorkerRequestPerformanceOutcomeTransport +): WorkerRequestPerformanceMetrics { + const identity = { + operation: input.operation, + operationId: input.operationId, + operationIdUnavailableReason: input.operationIdUnavailableReason, + playlistId: input.playlistId, + requestId: input.requestId, + responseEpochMs: input.responseEpochMs, + success: input.success, + }; + const unavailable = ( + reason: + | 'worker-performance-capture-invalid' + | 'worker-performance-capture-missing' + ): WorkerRequestPerformanceMetrics => ({ + eventLoopDelay: null, + eventLoopDelayUnavailableReason: null, + eventLoopUtilization: null, + eventLoopUtilizationUnavailableReason: null, + histogramFlushedEpochMs: null, + invalidReason: null, + ...identity, + performanceCaptureUnavailableReason: reason, + requestReceivedEpochMs: null, + threadCpuSystemMicros: null, + threadCpuUnavailableReason: null, + threadCpuUserMicros: null, + workEndedEpochMs: null, + workStartedEpochMs: null, + }); + + if ( + input.performanceCapture === undefined || + input.performanceCapture === null + ) { + return unavailable('worker-performance-capture-missing'); + } + + const performanceCapture = parseWorkerPerformanceCapture( + input.performanceCapture + ); + if (!performanceCapture) { + return unavailable('worker-performance-capture-invalid'); + } + + return { + ...performanceCapture, + ...identity, + performanceCaptureUnavailableReason: null, + }; +} + +export function selectMainCaptureGeneration( + input: MainCaptureGenerationTransport +): MainCaptureMetrics { + const workers = input.workers + .filter( + (worker) => worker.captureGeneration === input.captureGeneration + ) + .map((worker): WorkerCaptureMetrics => { + const requests = worker.requests.map( + normalizeWorkerRequestPerformanceOutcome + ); + const onlyRequest = + requests.length === 1 ? (requests[0] ?? null) : null; + return { + ...worker.metrics, + eventLoopDelay: onlyRequest?.eventLoopDelay ?? null, + eventLoopDelayUnavailableReason: onlyRequest + ? (onlyRequest.eventLoopDelayUnavailableReason ?? + onlyRequest.performanceCaptureUnavailableReason ?? + 'worker-self-profile-result-unavailable') + : requests.length > 1 + ? 'multiple-request-scoped-captures' + : worker.metrics.terminatedEpochMs !== null + ? 'worker-terminated-before-profile-flush' + : 'worker-self-profile-result-unavailable', + requests, + }; + }); + return { + ...input.metrics, + workers, + }; +} diff --git a/apps/electron-backend/src/app/workers/database.worker.ts b/apps/electron-backend/src/app/workers/database.worker.ts index b3b4d968a..2310d3398 100644 --- a/apps/electron-backend/src/app/workers/database.worker.ts +++ b/apps/electron-backend/src/app/workers/database.worker.ts @@ -90,8 +90,11 @@ import { } from '../database/operations/xtream.operations'; import { armWorkerPerformanceCapture, - finishWorkerPerformanceCapture, + executeWithWorkerPerformanceCapture, + registerDatabaseWorkerPerformanceCapture, + releaseDatabaseWorkerPerformanceCapture, startWorkerPerformanceCapture, + type WorkerPerformanceCapture, } from './worker-performance-capture'; const loggerLabel = '[DB Worker]'; @@ -104,6 +107,11 @@ type ActiveOperationState = { cancelled: boolean; }; +type PreRegisteredActiveOperation = { + operationId: string; + state: ActiveOperationState; +}; + type OperationController = { control: { checkpoint: () => Promise; @@ -129,6 +137,7 @@ type OperationController = { }; const activeOperations = new Map(); +const activePerformanceCaptures = new Set(); if (!parentPort) { throw new Error('Database worker must be started with a parent port'); @@ -178,17 +187,22 @@ async function pauseBetweenBatches(): Promise { await new Promise((resolve) => setTimeout(resolve, batchDelayMs)); } -function createOperationController(config: { - operationId?: string; - operation: string; - playlistId?: string; - requestId: string; - cancellable?: boolean; -}): OperationController { +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 ? { cancelled: false } : null; + operationId && cancellable + ? (preRegisteredState ?? { cancelled: false }) + : null; if (operationId && activeState) { activeOperations.set(operationId, activeState); @@ -255,7 +269,10 @@ function createOperationController(config: { }); }, cleanup: () => { - if (operationId) { + if ( + operationId && + activeOperations.get(operationId) === activeState + ) { activeOperations.delete(operationId); } }, @@ -264,9 +281,10 @@ function createOperationController(config: { async function executeTrackedOperation( config: Parameters[0], - handler: (controller: OperationController) => Promise + handler: (controller: OperationController) => Promise, + preRegisteredState?: ActiveOperationState ): Promise { - const controller = createOperationController(config); + const controller = createOperationController(config, preRegisteredState); try { return await handler(controller); @@ -282,7 +300,53 @@ async function executeTrackedOperation( } } -async function executeRequest(message: DbWorkerRequestMessage) { +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)[ + '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 +) { const db = await getWorkerDatabase(); switch (message.operation) { @@ -413,7 +477,8 @@ async function executeRequest(message: DbWorkerRequestMessage) { }); return result; - } + }, + preRegisteredState ); } @@ -586,7 +651,8 @@ async function executeRequest(message: DbWorkerRequestMessage) { }); return result; - } + }, + preRegisteredState ); } @@ -695,7 +761,8 @@ async function executeRequest(message: DbWorkerRequestMessage) { }); return result; - } + }, + preRegisteredState ); } @@ -739,7 +806,8 @@ async function executeRequest(message: DbWorkerRequestMessage) { }); return result; - } + }, + preRegisteredState ); } @@ -929,33 +997,49 @@ parentPort.on('message', async (message: DbWorkerIncomingMessage) => { return; } - const performanceCapture = startWorkerPerformanceCapture(); - await armWorkerPerformanceCapture(performanceCapture); - try { - const result = await executeRequest(message); - const performance = - await finishWorkerPerformanceCapture(performanceCapture); + const preRegisteredOperation = preRegisterCancellableOperation(message); + let performanceCapture: WorkerPerformanceCapture | null = null; + const execution = await (async () => { + try { + performanceCapture = startWorkerPerformanceCapture(); + registerDatabaseWorkerPerformanceCapture( + activePerformanceCaptures, + performanceCapture + ); + await armWorkerPerformanceCapture(performanceCapture); + return await executeWithWorkerPerformanceCapture( + performanceCapture, + () => executeRequest(message, preRegisteredOperation?.state) + ); + } finally { + releaseDatabaseWorkerPerformanceCapture( + activePerformanceCaptures, + performanceCapture + ); + releasePreRegisteredOperation(preRegisteredOperation); + } + })(); + + if (execution.success) { postMessage({ type: 'response', requestId: message.requestId, success: true, - result, - performance, + result: execution.result, + performance: execution.performance, }); - } catch (error) { + } else { console.error( loggerLabel, `Error handling ${message.operation}:`, - error + execution.error ); - const performance = - await finishWorkerPerformanceCapture(performanceCapture); postMessage({ type: 'response', requestId: message.requestId, success: false, - error: serializeError(error), - performance, + error: serializeError(execution.error), + performance: execution.performance, }); } }); diff --git a/apps/electron-backend/src/app/workers/playlist-refresh.worker.ts b/apps/electron-backend/src/app/workers/playlist-refresh.worker.ts index 790d6ac55..3bd0059ba 100644 --- a/apps/electron-backend/src/app/workers/playlist-refresh.worker.ts +++ b/apps/electron-backend/src/app/workers/playlist-refresh.worker.ts @@ -24,7 +24,7 @@ import { requestWithValidatedRedirects } from '../util/validated-axios'; import { PLAYLIST_FETCH_TIMEOUT_MS } from '../events/playlist-source'; import { armWorkerPerformanceCapture, - finishWorkerPerformanceCapture, + executeWithWorkerPerformanceCapture, startWorkerPerformanceCapture, } from './worker-performance-capture'; @@ -86,6 +86,15 @@ function checkpoint(payload: PlaylistRefreshPayload): void { } } +function releaseActiveRefresh( + operationId: string, + activeRefresh: ActiveRefreshState +): void { + if (activeRefreshes.get(operationId) === activeRefresh) { + activeRefreshes.delete(operationId); + } +} + async function fetchPlaylistFromUrl( payload: PlaylistRefreshPayload, controller: AbortController @@ -159,23 +168,18 @@ async function fetchPlaylistFromFile( } async function executeRefresh( - payload: PlaylistRefreshPayload + payload: PlaylistRefreshPayload, + activeRefresh: ActiveRefreshState ): Promise { - const controller = new AbortController(); - activeRefreshes.set(payload.operationId, { - cancelled: false, - controller, - }); - try { const playlist = payload.url - ? await fetchPlaylistFromUrl(payload, controller) + ? await fetchPlaylistFromUrl(payload, activeRefresh.controller) : await fetchPlaylistFromFile(payload); checkpoint(payload); return playlist; } finally { - activeRefreshes.delete(payload.operationId); + releaseActiveRefresh(payload.operationId, activeRefresh); } } @@ -191,39 +195,55 @@ parentPort.on( return; } - const performanceCapture = startWorkerPerformanceCapture(); - await armWorkerPerformanceCapture(performanceCapture); + const payload = message.payload; + const activeRefresh: ActiveRefreshState = { + cancelled: false, + controller: new AbortController(), + }; + activeRefreshes.set(payload.operationId, activeRefresh); + try { - const result = await executeRefresh(message.payload); - const performance = - await finishWorkerPerformanceCapture(performanceCapture); - postMessage({ - type: 'response', - success: true, - result, - performance, - }); - } catch (error) { - const payload = message.payload; - if (error instanceof Error && error.name === 'AbortError') { - emitEvent(payload, { status: 'cancelled', phase: 'parsing' }); + const performanceCapture = startWorkerPerformanceCapture(); + await armWorkerPerformanceCapture(performanceCapture); + const execution = await executeWithWorkerPerformanceCapture( + performanceCapture, + () => executeRefresh(payload, activeRefresh) + ); + + if (execution.success) { + postMessage({ + type: 'response', + success: true, + result: execution.result, + performance: execution.performance, + }); } else { - emitEvent(payload, { - status: 'error', - phase: payload.url ? 'fetching' : 'reading-file', - error: - error instanceof Error ? error.message : String(error), + const error = execution.error; + 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), + performance: execution.performance, }); } - - const performance = - await finishWorkerPerformanceCapture(performanceCapture); - postMessage({ - type: 'response', - success: false, - error: serializeError(error), - performance, - }); + } finally { + releaseActiveRefresh(payload.operationId, activeRefresh); } } ); diff --git a/apps/electron-backend/src/app/workers/worker-performance-cancellation.spec.ts b/apps/electron-backend/src/app/workers/worker-performance-cancellation.spec.ts new file mode 100644 index 000000000..eef0171f3 --- /dev/null +++ b/apps/electron-backend/src/app/workers/worker-performance-cancellation.spec.ts @@ -0,0 +1,284 @@ +import type { + DbWorkerIncomingMessage, + DbWorkerMessage, +} from './database-worker.types'; +import type { + PlaylistRefreshWorkerIncomingMessage, + PlaylistRefreshWorkerMessage, +} from './playlist-refresh.worker.types'; + +type WorkerMessageHandler = ( + message: TMessage +) => Promise | void; + +interface ArmGate { + readonly promise: Promise; + readonly release: () => void; +} + +function createArmGate(): ArmGate { + let release = (): void => undefined; + const promise = new Promise((resolvePromise) => { + release = resolvePromise; + }); + return { promise, release }; +} + +function createParentPortHarness() { + let messageHandler: WorkerMessageHandler | null = null; + const postMessage = jest.fn(); + const parentPort = { + on: jest.fn( + (event: string, handler: WorkerMessageHandler): void => { + if (event === 'message') { + messageHandler = handler; + } + } + ), + postMessage, + }; + + return { + getMessageHandler: (): WorkerMessageHandler => { + if (!messageHandler) { + throw new Error('Worker message handler was not registered'); + } + return messageHandler; + }, + parentPort, + postMessage, + }; +} + +function mockPerformanceCapture(armGate: ArmGate): void { + jest.doMock('./worker-performance-capture', () => ({ + armWorkerPerformanceCapture: jest.fn(() => armGate.promise), + executeWithWorkerPerformanceCapture: jest.fn( + async (_capture: unknown, execute: () => Promise) => { + try { + return { + error: null, + performance: undefined, + result: await execute(), + success: true, + }; + } catch (error) { + return { + error, + performance: undefined, + result: undefined, + success: false, + }; + } + } + ), + registerDatabaseWorkerPerformanceCapture: jest.fn(), + releaseDatabaseWorkerPerformanceCapture: jest.fn(), + startWorkerPerformanceCapture: jest.fn(() => ({})), + })); +} + +describe('worker cancellation while performance capture arms', () => { + afterEach(() => { + jest.restoreAllMocks(); + jest.resetModules(); + }); + + it('cancels playlist work before file access when cancel arrives during arming', async () => { + const armGate = createArmGate(); + const port = createParentPortHarness< + PlaylistRefreshWorkerIncomingMessage, + PlaylistRefreshWorkerMessage + >(); + const readFile = jest.fn().mockResolvedValue('#EXTM3U'); + mockPerformanceCapture(armGate); + jest.doMock('worker_threads', () => ({ + parentPort: port.parentPort, + })); + jest.doMock('node:fs/promises', () => ({ readFile })); + jest.doMock('iptv-playlist-parser', () => ({ + parse: jest.fn(() => ({ items: [] })), + })); + jest.doMock('@iptvnator/shared/m3u-utils', () => ({ + createPlaylistObject: jest.fn(() => ({ id: 'playlist-1' })), + getFilenameFromUrl: jest.fn(() => 'playlist.m3u'), + })); + + await import('./playlist-refresh.worker'); + port.postMessage.mockClear(); + const handleMessage = port.getMessageHandler(); + const requestPromise = handleMessage({ + type: 'request', + payload: { + filePath: '/fixtures/playlist.m3u', + operationId: 'refresh-1', + playlistId: 'playlist-1', + title: 'Fixture', + }, + }); + + await handleMessage({ + type: 'cancel', + operationId: 'refresh-1', + }); + armGate.release(); + await requestPromise; + + expect(readFile).not.toHaveBeenCalled(); + expect(port.postMessage).toHaveBeenCalledWith( + expect.objectContaining({ + event: expect.objectContaining({ + operationId: 'refresh-1', + status: 'cancelled', + }), + type: 'event', + }) + ); + expect(port.postMessage).toHaveBeenCalledWith( + expect.objectContaining({ + error: expect.objectContaining({ name: 'AbortError' }), + success: false, + type: 'response', + }) + ); + }); + + it('cancels cancellable database work when cancel arrives during arming', async () => { + const armGate = createArmGate(); + const port = createParentPortHarness< + DbWorkerIncomingMessage, + DbWorkerMessage + >(); + const saveContent = jest.fn( + async ( + _db: unknown, + _playlistId: string, + _streams: unknown[], + _type: string, + control: { checkpoint: () => Promise } + ) => { + await control.checkpoint(); + return { count: 0 }; + } + ); + mockPerformanceCapture(armGate); + jest.doMock('worker_threads', () => ({ + parentPort: port.parentPort, + })); + jest.doMock('./database.worker-connection', () => ({ + closeWorkerDatabase: jest.fn(), + getWorkerDatabase: jest.fn().mockResolvedValue({}), + })); + jest.doMock('../database/operations/content.operations', () => ({ + saveContent, + })); + jest.spyOn(console, 'error').mockImplementation(() => undefined); + + await import('./database.worker'); + port.postMessage.mockClear(); + const handleMessage = port.getMessageHandler(); + const requestPromise = handleMessage({ + type: 'request', + operation: 'DB_SAVE_CONTENT', + payload: { + operationId: 'database-operation-1', + playlistId: 'playlist-1', + streams: [], + type: 'movie', + }, + requestId: 'request-1', + }); + + await handleMessage({ + type: 'cancel', + operationId: 'database-operation-1', + }); + armGate.release(); + await requestPromise; + + expect(saveContent).toHaveBeenCalledTimes(1); + expect(port.postMessage).toHaveBeenCalledWith( + expect.objectContaining({ + event: expect.objectContaining({ + operationId: 'database-operation-1', + status: 'cancelled', + }), + requestId: 'request-1', + type: 'event', + }) + ); + expect(port.postMessage).toHaveBeenCalledWith( + expect.objectContaining({ + error: expect.objectContaining({ name: 'AbortError' }), + requestId: 'request-1', + success: false, + type: 'response', + }) + ); + }); + + it('does not make an explicitly non-cancellable database operation cancellable during arming', async () => { + const armGate = createArmGate(); + const port = createParentPortHarness< + DbWorkerIncomingMessage, + DbWorkerMessage + >(); + const deleteAllPlaylists = jest.fn( + async ( + _db: unknown, + control: { checkpoint: () => Promise } + ) => { + await control.checkpoint(); + return { success: true }; + } + ); + mockPerformanceCapture(armGate); + jest.doMock('worker_threads', () => ({ + parentPort: port.parentPort, + })); + jest.doMock('./database.worker-connection', () => ({ + closeWorkerDatabase: jest.fn(), + getWorkerDatabase: jest.fn().mockResolvedValue({}), + })); + jest.doMock('../database/operations/playlist.operations', () => ({ + deleteAllPlaylists, + })); + jest.spyOn(console, 'error').mockImplementation(() => undefined); + + await import('./database.worker'); + port.postMessage.mockClear(); + const handleMessage = port.getMessageHandler(); + const requestPromise = handleMessage({ + type: 'request', + operation: 'DB_DELETE_ALL_PLAYLISTS', + payload: { + operationId: 'database-operation-2', + }, + requestId: 'request-2', + }); + + await handleMessage({ + type: 'cancel', + operationId: 'database-operation-2', + }); + armGate.release(); + await requestPromise; + + expect(deleteAllPlaylists).toHaveBeenCalledTimes(1); + expect(port.postMessage).toHaveBeenCalledWith( + expect.objectContaining({ + requestId: 'request-2', + result: { success: true }, + success: true, + type: 'response', + }) + ); + expect(port.postMessage).not.toHaveBeenCalledWith( + expect.objectContaining({ + event: expect.objectContaining({ status: 'cancelled' }), + requestId: 'request-2', + type: 'event', + }) + ); + }); +}); diff --git a/apps/electron-backend/src/app/workers/worker-performance-capture.concurrency.spec.ts b/apps/electron-backend/src/app/workers/worker-performance-capture.concurrency.spec.ts new file mode 100644 index 000000000..900f4ccde --- /dev/null +++ b/apps/electron-backend/src/app/workers/worker-performance-capture.concurrency.spec.ts @@ -0,0 +1,95 @@ +import { + WORKER_PERFORMANCE_INVALID_REASON, + armWorkerPerformanceCapture, + executeWithWorkerPerformanceCapture, + registerDatabaseWorkerPerformanceCapture, + releaseDatabaseWorkerPerformanceCapture, + startWorkerPerformanceCapture, + type WorkerPerformanceCapture, +} from './worker-performance-capture'; +import { createFakeRuntime } from './worker-performance-capture.test-harness'; + +const PROFILING_ENV = 'IPTVNATOR_PERF_WORKER_PROFILING'; + +describe('worker performance capture concurrency and real timers', () => { + const originalProfilingValue = process.env[PROFILING_ENV]; + + afterEach(() => { + if (originalProfilingValue === undefined) { + delete process.env[PROFILING_ENV]; + } else { + process.env[PROFILING_ENV] = originalProfilingValue; + } + }); + + it('fails every overlapping database request capture closed without serializing execution', async () => { + const activeCaptures = new Set(); + const firstHarness = createFakeRuntime(); + const secondHarness = createFakeRuntime(); + const first = startWorkerPerformanceCapture({ + enabled: true, + runtime: firstHarness.runtime, + }); + const second = startWorkerPerformanceCapture({ + enabled: true, + runtime: secondHarness.runtime, + }); + registerDatabaseWorkerPerformanceCapture(activeCaptures, first); + registerDatabaseWorkerPerformanceCapture(activeCaptures, second); + await Promise.all([ + armWorkerPerformanceCapture(first), + armWorkerPerformanceCapture(second), + ]); + + const [firstExecution, secondExecution] = await Promise.all([ + executeWithWorkerPerformanceCapture(first, async () => 'first'), + executeWithWorkerPerformanceCapture(second, async () => 'second'), + ]); + releaseDatabaseWorkerPerformanceCapture(activeCaptures, first); + releaseDatabaseWorkerPerformanceCapture(activeCaptures, second); + + for (const execution of [firstExecution, secondExecution]) { + expect(execution.performance).toMatchObject({ + eventLoopDelay: null, + eventLoopDelayUnavailableReason: + WORKER_PERFORMANCE_INVALID_REASON.OVERLAPPING_DATABASE_WORKER_REQUESTS, + eventLoopUtilization: null, + eventLoopUtilizationUnavailableReason: + WORKER_PERFORMANCE_INVALID_REASON.OVERLAPPING_DATABASE_WORKER_REQUESTS, + histogramFlushedEpochMs: null, + invalidReason: + WORKER_PERFORMANCE_INVALID_REASON.OVERLAPPING_DATABASE_WORKER_REQUESTS, + threadCpuSystemMicros: null, + threadCpuUnavailableReason: + WORKER_PERFORMANCE_INVALID_REASON.OVERLAPPING_DATABASE_WORKER_REQUESTS, + threadCpuUserMicros: null, + }); + } + expect(firstExecution.result).toBe('first'); + expect(secondExecution.result).toBe('second'); + expect(activeCaptures.size).toBe(0); + }); + + it('reliably observes a real 20ms event-loop block after condition-based arming', async () => { + process.env[PROFILING_ENV] = '1'; + + const capture = startWorkerPerformanceCapture(); + await armWorkerPerformanceCapture(capture); + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => { + const blockStartedAt = performance.now(); + while (performance.now() - blockStartedAt < 20) { + // This deliberate block is the behavior under measurement. + } + } + ); + + expect(execution.success).toBe(true); + expect(execution.performance?.eventLoopDelay).not.toBeNull(); + expect(execution.performance?.eventLoopDelay?.maxMs).toBeGreaterThan( + 10 + ); + expect(execution.performance?.histogramFlushedEpochMs).not.toBeNull(); + }); +}); diff --git a/apps/electron-backend/src/app/workers/worker-performance-capture.histogram.spec.ts b/apps/electron-backend/src/app/workers/worker-performance-capture.histogram.spec.ts new file mode 100644 index 000000000..19462fe5a --- /dev/null +++ b/apps/electron-backend/src/app/workers/worker-performance-capture.histogram.spec.ts @@ -0,0 +1,39 @@ +import { + armWorkerPerformanceCapture, + executeWithWorkerPerformanceCapture, + startWorkerPerformanceCapture, +} from './worker-performance-capture'; +import { createFakeRuntime } from './worker-performance-capture.test-harness'; + +describe('worker performance histogram lifecycle', () => { + it('accepts a nonthrowing false disable result because the histogram is already disabled', async () => { + const harness = createFakeRuntime({ + histogramDisableResult: false, + }); + const capture = startWorkerPerformanceCapture({ + enabled: true, + runtime: harness.runtime, + }); + await armWorkerPerformanceCapture(capture); + + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => 'business-result' + ); + + expect(execution).toMatchObject({ + error: null, + result: 'business-result', + success: true, + performance: { + eventLoopDelay: { + maxMs: 24, + p95Ms: 18, + p99Ms: 22, + }, + eventLoopDelayUnavailableReason: null, + histogramFlushedEpochMs: 145, + }, + }); + }); +}); diff --git a/apps/electron-backend/src/app/workers/worker-performance-capture.histogram.ts b/apps/electron-backend/src/app/workers/worker-performance-capture.histogram.ts new file mode 100644 index 000000000..a4951d6cc --- /dev/null +++ b/apps/electron-backend/src/app/workers/worker-performance-capture.histogram.ts @@ -0,0 +1,136 @@ +import { + WORKER_PERFORMANCE_UNAVAILABLE_REASON, + type WorkerEventLoopDelayHistogram, + type WorkerEventLoopDelayMetrics, + type WorkerPerformanceCapture, + type WorkerPerformanceMetricUnavailableReason, +} from './worker-performance-capture.model'; +import { + HISTOGRAM_POLL_INTERVAL_MS, + HISTOGRAM_WAIT_CAP_MS, +} from './worker-performance-capture.runtime'; + +export function disableWorkerPerformanceHistogram( + capture: WorkerPerformanceCapture +): boolean { + if (!capture.eventLoopDelay || capture.histogramDisabled) { + return true; + } + + try { + capture.eventLoopDelay.disable(); + capture.histogramDisabled = true; + return true; + } catch { + // Performance instrumentation must never affect worker execution. + return false; + } +} + +export function markEventLoopDelayUnavailable( + capture: WorkerPerformanceCapture, + reason: WorkerPerformanceMetricUnavailableReason +): void { + capture.eventLoopDelayUnavailableReason ??= reason; + disableWorkerPerformanceHistogram(capture); +} + +export async function waitForHistogramCondition( + capture: WorkerPerformanceCapture, + condition: (histogram: WorkerEventLoopDelayHistogram) => boolean +): Promise { + const histogram = capture.eventLoopDelay; + if ( + !histogram || + capture.histogramDisabled || + capture.invalidReason !== null + ) { + return false; + } + + let startedAt: number; + try { + startedAt = capture.runtime.readMonotonicMs(); + } catch { + return false; + } + + const maximumPolls = Math.ceil( + HISTOGRAM_WAIT_CAP_MS / HISTOGRAM_POLL_INTERVAL_MS + ); + let pollCount = 0; + + while (true) { + if (capture.invalidReason !== null || capture.histogramDisabled) { + return false; + } + try { + if (condition(histogram)) { + return true; + } + } catch { + return false; + } + + let elapsedMs: number; + try { + elapsedMs = capture.runtime.readMonotonicMs() - startedAt; + } catch { + return false; + } + if ( + !Number.isFinite(elapsedMs) || + elapsedMs >= HISTOGRAM_WAIT_CAP_MS || + pollCount >= maximumPolls + ) { + return false; + } + + const delayMs = Math.min( + HISTOGRAM_POLL_INTERVAL_MS, + HISTOGRAM_WAIT_CAP_MS - Math.max(0, elapsedMs) + ); + try { + pollCount += 1; + await new Promise((resolvePromise) => { + capture.runtime.scheduleTimeout(resolvePromise, delayMs); + }); + } catch { + return false; + } + } +} + +export function readWorkerEventLoopDelay( + capture: WorkerPerformanceCapture +): WorkerEventLoopDelayMetrics | null { + if ( + capture.invalidReason !== null || + capture.eventLoopDelayUnavailableReason !== null || + capture.histogramFlushedEpochMs === null || + !capture.eventLoopDelay + ) { + return null; + } + + try { + const metrics: WorkerEventLoopDelayMetrics = { + maxMs: capture.eventLoopDelay.max / 1e6, + p95Ms: capture.eventLoopDelay.percentile(95) / 1e6, + p99Ms: capture.eventLoopDelay.percentile(99) / 1e6, + }; + if ( + Object.values(metrics).every( + (value) => Number.isFinite(value) && value >= 0 + ) + ) { + return metrics; + } + } catch { + // The fixed reason below is intentionally payload-free. + } + + capture.eventLoopDelayUnavailableReason = + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_INVALID; + return null; +} diff --git a/apps/electron-backend/src/app/workers/worker-performance-capture.metrics.ts b/apps/electron-backend/src/app/workers/worker-performance-capture.metrics.ts new file mode 100644 index 000000000..4d1fb3a1f --- /dev/null +++ b/apps/electron-backend/src/app/workers/worker-performance-capture.metrics.ts @@ -0,0 +1,235 @@ +import type { EventLoopUtilization } from 'node:perf_hooks'; +import { + WORKER_PERFORMANCE_UNAVAILABLE_REASON, + type WorkerPerformanceCapture, + type WorkerPerformanceCaptureResult, +} from './worker-performance-capture.model'; +import { readWorkerEventLoopDelay } from './worker-performance-capture.histogram'; +import { safeReadEpochMs } from './worker-performance-capture.runtime'; + +function readThreadCpuBoundary( + capture: WorkerPerformanceCapture +): NodeJS.CpuUsage | null { + if (!capture.runtime.readThreadCpuUsage) { + capture.threadCpuUnavailableReason ??= + WORKER_PERFORMANCE_UNAVAILABLE_REASON.THREAD_CPU_USAGE_UNAVAILABLE; + return null; + } + + try { + const usage = capture.runtime.readThreadCpuUsage(); + if ( + !Number.isFinite(usage.system) || + !Number.isFinite(usage.user) || + usage.system < 0 || + usage.user < 0 + ) { + capture.threadCpuUnavailableReason = + WORKER_PERFORMANCE_UNAVAILABLE_REASON.THREAD_CPU_USAGE_INVALID; + return null; + } + return usage; + } catch { + capture.threadCpuUnavailableReason = + WORKER_PERFORMANCE_UNAVAILABLE_REASON.THREAD_CPU_USAGE_UNAVAILABLE; + return null; + } +} + +function readEventLoopUtilizationBoundary( + capture: WorkerPerformanceCapture +): EventLoopUtilization | null { + try { + const utilization = capture.runtime.readEventLoopUtilization(); + if ( + !Number.isFinite(utilization.active) || + !Number.isFinite(utilization.idle) + ) { + capture.eventLoopUtilizationUnavailableReason = + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_UTILIZATION_UNAVAILABLE; + return null; + } + return utilization; + } catch { + capture.eventLoopUtilizationUnavailableReason = + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_UTILIZATION_UNAVAILABLE; + return null; + } +} + +export function beginWorkerPerformanceCapture( + capture: WorkerPerformanceCapture | null +): void { + if (!capture) { + return; + } + + capture.workStartedEpochMs = safeReadEpochMs(capture.runtime); + capture.threadCpuStart = readThreadCpuBoundary(capture); + capture.eventLoopUtilizationStart = + readEventLoopUtilizationBoundary(capture); +} + +function calculateThreadCpu( + capture: WorkerPerformanceCapture +): Pick< + WorkerPerformanceCaptureResult, + 'threadCpuSystemMicros' | 'threadCpuUserMicros' +> { + if ( + capture.invalidReason !== null || + capture.threadCpuUnavailableReason !== null || + !capture.threadCpuStart || + !capture.threadCpuEnd + ) { + return { + threadCpuSystemMicros: null, + threadCpuUserMicros: null, + }; + } + + try { + const system = + capture.threadCpuEnd.system - capture.threadCpuStart.system; + const user = capture.threadCpuEnd.user - capture.threadCpuStart.user; + if ( + !Number.isFinite(system) || + !Number.isFinite(user) || + system < 0 || + user < 0 + ) { + capture.threadCpuUnavailableReason = + WORKER_PERFORMANCE_UNAVAILABLE_REASON.THREAD_CPU_USAGE_INVALID; + return { + threadCpuSystemMicros: null, + threadCpuUserMicros: null, + }; + } + + return { + threadCpuSystemMicros: system, + threadCpuUserMicros: user, + }; + } catch { + capture.threadCpuUnavailableReason = + WORKER_PERFORMANCE_UNAVAILABLE_REASON.THREAD_CPU_USAGE_INVALID; + return { + threadCpuSystemMicros: null, + threadCpuUserMicros: null, + }; + } +} + +function calculateEventLoopUtilization( + capture: WorkerPerformanceCapture +): number | null { + if ( + capture.invalidReason !== null || + capture.eventLoopUtilizationUnavailableReason !== null || + !capture.eventLoopUtilizationStart || + !capture.eventLoopUtilizationEnd + ) { + return null; + } + + try { + const active = + capture.eventLoopUtilizationEnd.active - + capture.eventLoopUtilizationStart.active; + const idle = + capture.eventLoopUtilizationEnd.idle - + capture.eventLoopUtilizationStart.idle; + const total = active + idle; + if ( + !Number.isFinite(active) || + !Number.isFinite(idle) || + active < 0 || + idle < 0 || + total < 0 + ) { + capture.eventLoopUtilizationUnavailableReason = + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_UTILIZATION_UNAVAILABLE; + return null; + } + + return total === 0 ? 0 : active / total; + } catch { + capture.eventLoopUtilizationUnavailableReason = + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_UTILIZATION_UNAVAILABLE; + return null; + } +} + +function createFallbackResult( + capture: WorkerPerformanceCapture +): WorkerPerformanceCaptureResult { + const invalidReason = capture.invalidReason; + return { + eventLoopDelay: null, + eventLoopDelayUnavailableReason: + invalidReason ?? + capture.eventLoopDelayUnavailableReason ?? + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_CAPTURE_UNAVAILABLE, + eventLoopUtilization: null, + eventLoopUtilizationUnavailableReason: + invalidReason ?? + capture.eventLoopUtilizationUnavailableReason ?? + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_UTILIZATION_UNAVAILABLE, + histogramFlushedEpochMs: null, + invalidReason, + requestReceivedEpochMs: capture.requestReceivedEpochMs, + threadCpuSystemMicros: null, + threadCpuUnavailableReason: + invalidReason ?? + capture.threadCpuUnavailableReason ?? + WORKER_PERFORMANCE_UNAVAILABLE_REASON.THREAD_CPU_USAGE_UNAVAILABLE, + threadCpuUserMicros: null, + workEndedEpochMs: + capture.workEndedEpochMs ?? capture.requestReceivedEpochMs, + workStartedEpochMs: + capture.workStartedEpochMs ?? capture.requestReceivedEpochMs, + }; +} + +export function createWorkerPerformanceCaptureResult( + capture: WorkerPerformanceCapture +): WorkerPerformanceCaptureResult { + try { + const overlapReason = capture.invalidReason; + if (overlapReason) { + capture.eventLoopDelayUnavailableReason = overlapReason; + capture.eventLoopUtilizationUnavailableReason = overlapReason; + capture.threadCpuUnavailableReason = overlapReason; + capture.histogramFlushedEpochMs = null; + } + const threadCpu = calculateThreadCpu(capture); + + return { + eventLoopDelay: readWorkerEventLoopDelay(capture), + eventLoopDelayUnavailableReason: + capture.eventLoopDelayUnavailableReason, + eventLoopUtilization: calculateEventLoopUtilization(capture), + eventLoopUtilizationUnavailableReason: + capture.eventLoopUtilizationUnavailableReason, + histogramFlushedEpochMs: capture.histogramFlushedEpochMs, + invalidReason: capture.invalidReason, + requestReceivedEpochMs: capture.requestReceivedEpochMs, + ...threadCpu, + threadCpuUnavailableReason: capture.threadCpuUnavailableReason, + workEndedEpochMs: + capture.workEndedEpochMs ?? capture.requestReceivedEpochMs, + workStartedEpochMs: + capture.workStartedEpochMs ?? capture.requestReceivedEpochMs, + }; + } catch { + return createFallbackResult(capture); + } +} + +export function finishWorkerPerformanceMetricBoundaries( + capture: WorkerPerformanceCapture +): void { + capture.workEndedEpochMs = safeReadEpochMs(capture.runtime); + capture.threadCpuEnd = readThreadCpuBoundary(capture); + capture.eventLoopUtilizationEnd = readEventLoopUtilizationBoundary(capture); +} diff --git a/apps/electron-backend/src/app/workers/worker-performance-capture.model.ts b/apps/electron-backend/src/app/workers/worker-performance-capture.model.ts new file mode 100644 index 000000000..f86576eda --- /dev/null +++ b/apps/electron-backend/src/app/workers/worker-performance-capture.model.ts @@ -0,0 +1,94 @@ +import type { EventLoopUtilization } from 'node:perf_hooks'; + +export const WORKER_PERFORMANCE_INVALID_REASON = { + OVERLAPPING_DATABASE_WORKER_REQUESTS: + 'overlapping-database-worker-requests', +} as const; + +export type WorkerPerformanceInvalidReason = + (typeof WORKER_PERFORMANCE_INVALID_REASON)[keyof typeof WORKER_PERFORMANCE_INVALID_REASON]; + +export const WORKER_PERFORMANCE_UNAVAILABLE_REASON = { + EVENT_LOOP_DELAY_ARM_TIMEOUT: 'event-loop-delay-arm-timeout', + EVENT_LOOP_DELAY_CAPTURE_UNAVAILABLE: + 'event-loop-delay-capture-unavailable', + EVENT_LOOP_DELAY_FLUSH_TIMEOUT: 'event-loop-delay-flush-timeout', + EVENT_LOOP_DELAY_INVALID: 'event-loop-delay-invalid', + EVENT_LOOP_UTILIZATION_UNAVAILABLE: 'event-loop-utilization-unavailable', + THREAD_CPU_USAGE_INVALID: 'thread-cpu-usage-invalid', + THREAD_CPU_USAGE_UNAVAILABLE: 'thread-cpu-usage-unavailable', +} as const; + +export type WorkerPerformanceUnavailableReason = + (typeof WORKER_PERFORMANCE_UNAVAILABLE_REASON)[keyof typeof WORKER_PERFORMANCE_UNAVAILABLE_REASON]; + +export type WorkerPerformanceMetricUnavailableReason = + WorkerPerformanceInvalidReason | WorkerPerformanceUnavailableReason; + +export interface WorkerEventLoopDelayMetrics { + maxMs: number; + p95Ms: number; + p99Ms: number; +} + +export interface WorkerEventLoopDelayHistogram { + readonly count: number; + readonly max: number; + disable: () => boolean; + enable: () => boolean; + percentile: (percentile: number) => number; +} + +export interface WorkerPerformanceCaptureResult { + eventLoopDelay: WorkerEventLoopDelayMetrics | null; + eventLoopDelayUnavailableReason: WorkerPerformanceMetricUnavailableReason | null; + eventLoopUtilization: number | null; + eventLoopUtilizationUnavailableReason: WorkerPerformanceMetricUnavailableReason | null; + histogramFlushedEpochMs: number | null; + invalidReason: WorkerPerformanceInvalidReason | null; + requestReceivedEpochMs: number; + threadCpuSystemMicros: number | null; + threadCpuUnavailableReason: WorkerPerformanceMetricUnavailableReason | null; + threadCpuUserMicros: number | null; + workEndedEpochMs: number; + workStartedEpochMs: number; +} + +export interface WorkerPerformanceCaptureRuntime { + createEventLoopDelayHistogram: () => WorkerEventLoopDelayHistogram; + readEpochMs: () => number; + readEventLoopUtilization: () => EventLoopUtilization; + readMonotonicMs: () => number; + readThreadCpuUsage: (() => NodeJS.CpuUsage) | null; + scheduleTimeout: (callback: () => void, delayMs: number) => void; +} + +export interface WorkerPerformanceCaptureOptions { + enabled?: boolean; + runtime?: WorkerPerformanceCaptureRuntime; +} + +export interface WorkerPerformanceExecutionResult { + error: unknown | null; + performance: WorkerPerformanceCaptureResult | undefined; + result: TResult | undefined; + success: boolean; +} + +export interface WorkerPerformanceCapture { + eventLoopDelay: WorkerEventLoopDelayHistogram | null; + eventLoopDelayUnavailableReason: WorkerPerformanceMetricUnavailableReason | null; + eventLoopUtilizationEnd: EventLoopUtilization | null; + eventLoopUtilizationStart: EventLoopUtilization | null; + eventLoopUtilizationUnavailableReason: WorkerPerformanceMetricUnavailableReason | null; + histogramDisabled: boolean; + histogramFlushedEpochMs: number | null; + invalidReason: WorkerPerformanceInvalidReason | null; + requestReceivedEpochMs: number; + runtime: WorkerPerformanceCaptureRuntime; + threadCpuEnd: NodeJS.CpuUsage | null; + threadCpuStart: NodeJS.CpuUsage | null; + threadCpuUnavailableReason: WorkerPerformanceMetricUnavailableReason | null; + workEndedEpochMs: number | null; + workStartedEpochMs: number | null; +} diff --git a/apps/electron-backend/src/app/workers/worker-performance-capture.runtime.ts b/apps/electron-backend/src/app/workers/worker-performance-capture.runtime.ts new file mode 100644 index 000000000..15c0a4308 --- /dev/null +++ b/apps/electron-backend/src/app/workers/worker-performance-capture.runtime.ts @@ -0,0 +1,45 @@ +import { monitorEventLoopDelay, performance } from 'node:perf_hooks'; +import type { + WorkerPerformanceCaptureRuntime, + WorkerEventLoopDelayHistogram, +} from './worker-performance-capture.model'; + +export const WORKER_PROFILING_ENV = 'IPTVNATOR_PERF_WORKER_PROFILING'; +export const HISTOGRAM_POLL_INTERVAL_MS = 1; +export const HISTOGRAM_WAIT_CAP_MS = 50; + +function readRuntimeThreadCpuUsage(): (() => NodeJS.CpuUsage) | null { + const candidate: unknown = Reflect.get(process, 'threadCpuUsage'); + if (typeof candidate !== 'function') { + return null; + } + + return () => + (candidate as (this: NodeJS.Process) => NodeJS.CpuUsage).call(process); +} + +export const DEFAULT_WORKER_PERFORMANCE_RUNTIME: WorkerPerformanceCaptureRuntime = + { + createEventLoopDelayHistogram: (): WorkerEventLoopDelayHistogram => + monitorEventLoopDelay({ + resolution: HISTOGRAM_POLL_INTERVAL_MS, + }), + readEpochMs: () => performance.timeOrigin + performance.now(), + readEventLoopUtilization: () => performance.eventLoopUtilization(), + readMonotonicMs: () => performance.now(), + readThreadCpuUsage: readRuntimeThreadCpuUsage(), + scheduleTimeout: (callback, delayMs) => { + setTimeout(callback, delayMs); + }, + }; + +export function safeReadEpochMs( + runtime: WorkerPerformanceCaptureRuntime +): number { + try { + const epochMs = runtime.readEpochMs(); + return Number.isFinite(epochMs) ? epochMs : Date.now(); + } catch { + return Date.now(); + } +} diff --git a/apps/electron-backend/src/app/workers/worker-performance-capture.spec.ts b/apps/electron-backend/src/app/workers/worker-performance-capture.spec.ts index 2801935c3..7bebd4fe9 100644 --- a/apps/electron-backend/src/app/workers/worker-performance-capture.spec.ts +++ b/apps/electron-backend/src/app/workers/worker-performance-capture.spec.ts @@ -1,8 +1,10 @@ import { + WORKER_PERFORMANCE_UNAVAILABLE_REASON, armWorkerPerformanceCapture, - finishWorkerPerformanceCapture, + executeWithWorkerPerformanceCapture, startWorkerPerformanceCapture, } from './worker-performance-capture'; +import { createFakeRuntime } from './worker-performance-capture.test-harness'; const PROFILING_ENV = 'IPTVNATOR_PERF_WORKER_PROFILING'; @@ -23,26 +25,355 @@ describe('worker performance capture', () => { expect(startWorkerPerformanceCapture()).toBeNull(); }); - it('reports finite worker event-loop metrics when opted in', async () => { - process.env[PROFILING_ENV] = '1'; - - const capture = startWorkerPerformanceCapture(); - await armWorkerPerformanceCapture(capture); - const blockStartedAt = performance.now(); - while (performance.now() - blockStartedAt < 20) { - // Deliberately block the loop so the histogram contract is tested. - } - const result = await finishWorkerPerformanceCapture(capture); - - expect(result).toEqual({ - eventLoopDelay: { - maxMs: expect.any(Number), - p95Ms: expect.any(Number), - p99Ms: expect.any(Number), - }, - eventLoopUtilization: expect.any(Number), + it('captures exact work-boundary CPU, ELU, timestamps, and a flushed fresh histogram', async () => { + const harness = createFakeRuntime(); + const capture = startWorkerPerformanceCapture({ + enabled: true, + runtime: harness.runtime, + }); + await armWorkerPerformanceCapture(capture); + + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => { + harness.lifecycle.push('execute'); + await Promise.resolve(); + harness.lifecycle.push('execute:settled'); + return 'result'; + } + ); + + expect(execution).toEqual({ + error: null, + performance: { + eventLoopDelay: { + maxMs: 24, + p95Ms: 18, + p99Ms: 22, + }, + eventLoopDelayUnavailableReason: null, + eventLoopUtilization: 0.75, + eventLoopUtilizationUnavailableReason: null, + histogramFlushedEpochMs: 145, + invalidReason: null, + requestReceivedEpochMs: 100, + threadCpuSystemMicros: 30, + threadCpuUnavailableReason: null, + threadCpuUserMicros: 80, + workEndedEpochMs: 140, + workStartedEpochMs: 110, + }, + result: 'result', + success: true, + }); + expect(harness.lifecycle).toEqual([ + 'epoch:100', + 'histogram:create', + 'timeout:1', + 'epoch:110', + 'cpu:1', + 'elu:1', + 'execute', + 'execute:settled', + 'epoch:140', + 'cpu:2', + 'elu:2', + 'timeout:1', + 'epoch:145', + ]); + expect(harness.histogramLifecycle.enableCalls).toBe(1); + expect(harness.histogramLifecycle.disableCalls).toBe(1); + expect('reset' in harness.histogram).toBe(false); + }); + + it('finishes profiling before returning a rejected execution to worker side effects', async () => { + const harness = createFakeRuntime(); + const capture = startWorkerPerformanceCapture({ + enabled: true, + runtime: harness.runtime, + }); + const failure = new Error('worker failed'); + await armWorkerPerformanceCapture(capture); + + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => { + harness.lifecycle.push('execute'); + await Promise.resolve(); + harness.lifecycle.push('execute:rejected'); + throw failure; + } + ); + harness.lifecycle.push('worker:emit-error'); + + expect(execution).toMatchObject({ + error: failure, + success: false, + }); + expect(harness.lifecycle.indexOf('epoch:140')).toBe( + harness.lifecycle.indexOf('execute:rejected') + 1 + ); + expect(harness.lifecycle.indexOf('worker:emit-error')).toBeGreaterThan( + harness.lifecycle.indexOf('epoch:145') + ); + }); + + it('times out arming after a monotonic 50ms without blocking work metrics', async () => { + const harness = createFakeRuntime({ + armHistogram: false, + }); + const capture = startWorkerPerformanceCapture({ + enabled: true, + runtime: harness.runtime, + }); + + await armWorkerPerformanceCapture(capture); + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => 'result' + ); + + expect( + harness.scheduledTimeouts.reduce( + (total, timeout) => total + timeout, + 0 + ) + ).toBe(50); + expect(execution.performance).toMatchObject({ + eventLoopDelay: null, + eventLoopDelayUnavailableReason: + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_ARM_TIMEOUT, + eventLoopUtilization: 0.75, + histogramFlushedEpochMs: null, + threadCpuSystemMicros: 30, + threadCpuUserMicros: 80, + }); + }); + + it('cannot poll forever when the monotonic runtime clock stalls', async () => { + const harness = createFakeRuntime({ + armHistogram: false, + stallMonotonicClock: true, + }); + const capture = startWorkerPerformanceCapture({ + enabled: true, + runtime: harness.runtime, + }); + + await armWorkerPerformanceCapture(capture); + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => 'result' + ); + + expect(harness.scheduledTimeouts).toHaveLength(50); + expect(execution).toMatchObject({ + result: 'result', + success: true, + performance: { + eventLoopDelay: null, + eventLoopDelayUnavailableReason: + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_ARM_TIMEOUT, + }, + }); + }); + + it('times out flushing after a monotonic 50ms while preserving work CPU and ELU', async () => { + const harness = createFakeRuntime({ + flushHistogram: false, + }); + const capture = startWorkerPerformanceCapture({ + enabled: true, + runtime: harness.runtime, + }); + await armWorkerPerformanceCapture(capture); + + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => 'result' + ); + + expect( + harness.scheduledTimeouts + .slice(1) + .reduce((total, timeout) => total + timeout, 0) + ).toBe(50); + expect(execution.performance).toMatchObject({ + eventLoopDelay: null, + eventLoopDelayUnavailableReason: + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_FLUSH_TIMEOUT, + eventLoopUtilization: 0.75, + histogramFlushedEpochMs: null, + threadCpuSystemMicros: 30, + threadCpuUserMicros: 80, + }); + }); + + it('reports unavailable thread CPU without falling back to process CPU usage', async () => { + const harness = createFakeRuntime({ + threadCpuAvailable: false, + }); + const capture = startWorkerPerformanceCapture({ + enabled: true, + runtime: harness.runtime, + }); + await armWorkerPerformanceCapture(capture); + + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => 'result' + ); + + expect(execution.performance).toMatchObject({ + eventLoopUtilization: 0.75, + threadCpuSystemMicros: null, + threadCpuUnavailableReason: + WORKER_PERFORMANCE_UNAVAILABLE_REASON.THREAD_CPU_USAGE_UNAVAILABLE, + threadCpuUserMicros: null, + }); + expect(harness.lifecycle).not.toContain('cpu:1'); + expect(harness.lifecycle).not.toContain('cpu:2'); + }); + + it('preserves the business result when CPU and ELU runtime callbacks throw', async () => { + const harness = createFakeRuntime({ + throwBoundaryCallbacks: true, + }); + const capture = startWorkerPerformanceCapture({ + enabled: true, + runtime: harness.runtime, + }); + await armWorkerPerformanceCapture(capture); + + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => 'business-result' + ); + + expect(execution).toMatchObject({ + error: null, + result: 'business-result', + success: true, + performance: { + eventLoopDelay: { + maxMs: 24, + p95Ms: 18, + p99Ms: 22, + }, + eventLoopUtilization: null, + eventLoopUtilizationUnavailableReason: + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_UTILIZATION_UNAVAILABLE, + threadCpuSystemMicros: null, + threadCpuUnavailableReason: + WORKER_PERFORMANCE_UNAVAILABLE_REASON.THREAD_CPU_USAGE_UNAVAILABLE, + threadCpuUserMicros: null, + }, + }); + }); + + it('preserves work metrics and the business result when timer scheduling throws', async () => { + const harness = createFakeRuntime({ + throwScheduleTimeout: true, + }); + const capture = startWorkerPerformanceCapture({ + enabled: true, + runtime: harness.runtime, + }); + await armWorkerPerformanceCapture(capture); + + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => 'business-result' + ); + + expect(execution).toMatchObject({ + error: null, + result: 'business-result', + success: true, + performance: { + eventLoopDelay: null, + eventLoopDelayUnavailableReason: + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_ARM_TIMEOUT, + eventLoopUtilization: 0.75, + threadCpuSystemMicros: 30, + threadCpuUserMicros: 80, + }, + }); + }); + + it.each([ + { + label: 'count getter during finish', + options: { + throwHistogramCountAtRead: 3, + }, + reason: WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_CAPTURE_UNAVAILABLE, + }, + { + label: 'metric getter after flush', + options: { + throwHistogramMetric: true, + }, + reason: WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_INVALID, + }, + ])( + 'preserves the business result when the histogram $label throws', + async ({ options, reason }) => { + const harness = createFakeRuntime(options); + const capture = startWorkerPerformanceCapture({ + enabled: true, + runtime: harness.runtime, + }); + await armWorkerPerformanceCapture(capture); + + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => 'business-result' + ); + + expect(execution).toMatchObject({ + error: null, + result: 'business-result', + success: true, + performance: { + eventLoopDelay: null, + eventLoopDelayUnavailableReason: reason, + eventLoopUtilization: 0.75, + threadCpuSystemMicros: 30, + threadCpuUserMicros: 80, + }, + }); + } + ); + + it('keeps the business result and rejects delay metrics when histogram disable throws', async () => { + const harness = createFakeRuntime({ + throwHistogramDisable: true, + }); + const capture = startWorkerPerformanceCapture({ + enabled: true, + runtime: harness.runtime, + }); + await armWorkerPerformanceCapture(capture); + + const execution = await executeWithWorkerPerformanceCapture( + capture, + async () => 'business-result' + ); + + expect(execution).toMatchObject({ + error: null, + result: 'business-result', + success: true, + performance: { + eventLoopDelay: null, + eventLoopDelayUnavailableReason: + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_CAPTURE_UNAVAILABLE, + eventLoopUtilization: 0.75, + histogramFlushedEpochMs: null, + threadCpuSystemMicros: 30, + threadCpuUserMicros: 80, + }, }); - expect(result?.eventLoopUtilization).toBeGreaterThanOrEqual(0); - expect(result?.eventLoopDelay.maxMs).toBeGreaterThan(10); }); }); diff --git a/apps/electron-backend/src/app/workers/worker-performance-capture.test-harness.ts b/apps/electron-backend/src/app/workers/worker-performance-capture.test-harness.ts new file mode 100644 index 000000000..54517e637 --- /dev/null +++ b/apps/electron-backend/src/app/workers/worker-performance-capture.test-harness.ts @@ -0,0 +1,158 @@ +import type { + WorkerEventLoopDelayHistogram, + WorkerPerformanceCaptureRuntime, +} from './worker-performance-capture'; + +interface FakeRuntimeOptions { + armHistogram?: boolean; + flushHistogram?: boolean; + histogramDisableResult?: boolean; + stallMonotonicClock?: boolean; + threadCpuAvailable?: boolean; + throwBoundaryCallbacks?: boolean; + throwHistogramDisable?: boolean; + throwHistogramCountAtRead?: number; + throwHistogramMetric?: boolean; + throwScheduleTimeout?: boolean; +} + +export interface FakeRuntimeHarness { + histogram: WorkerEventLoopDelayHistogram; + histogramLifecycle: { + disableCalls: number; + enableCalls: number; + }; + lifecycle: string[]; + runtime: WorkerPerformanceCaptureRuntime; + scheduledTimeouts: number[]; +} + +export function createFakeRuntime( + options: FakeRuntimeOptions = {} +): FakeRuntimeHarness { + const lifecycle: string[] = []; + const scheduledTimeouts: number[] = []; + let epochIndex = 0; + let monotonicMs = 0; + let timeoutCount = 0; + let cpuIndex = 0; + let eluIndex = 0; + let histogramCount = 0; + let histogramCountReads = 0; + const histogramLifecycle = { + disableCalls: 0, + enableCalls: 0, + }; + const epochValues = [100, 110, 140, 145]; + const cpuValues = [ + { system: 50, user: 100 }, + { system: 80, user: 180 }, + ]; + const eluValues = [ + { active: 10, idle: 20, utilization: 1 / 3 }, + { active: 40, idle: 30, utilization: 4 / 7 }, + ]; + const histogram = { + get count() { + histogramCountReads += 1; + if (histogramCountReads === options.throwHistogramCountAtRead) { + throw new Error('histogram count unavailable'); + } + return histogramCount; + }, + disable: () => { + histogramLifecycle.disableCalls += 1; + if (options.throwHistogramDisable === true) { + throw new Error('histogram disable unavailable'); + } + return options.histogramDisableResult ?? true; + }, + enable: () => { + histogramLifecycle.enableCalls += 1; + return true; + }, + get max() { + if (options.throwHistogramMetric === true) { + throw new Error('histogram max unavailable'); + } + return 24_000_000; + }, + percentile: (percentile: number) => { + if (options.throwHistogramMetric === true) { + throw new Error('histogram percentile unavailable'); + } + if (percentile === 95) { + return 18_000_000; + } + if (percentile === 99) { + return 22_000_000; + } + return 0; + }, + } satisfies WorkerEventLoopDelayHistogram; + const runtime: WorkerPerformanceCaptureRuntime = { + createEventLoopDelayHistogram: () => { + lifecycle.push('histogram:create'); + return histogram; + }, + readEpochMs: () => { + const value = epochValues[epochIndex] ?? 145 + epochIndex; + epochIndex += 1; + lifecycle.push(`epoch:${value}`); + return value; + }, + readEventLoopUtilization: () => { + if (options.throwBoundaryCallbacks === true) { + throw new Error('ELU callback unavailable'); + } + const value = eluValues[eluIndex] ?? eluValues.at(-1)!; + eluIndex += 1; + lifecycle.push(`elu:${eluIndex}`); + return value; + }, + readMonotonicMs: () => monotonicMs, + readThreadCpuUsage: + options.threadCpuAvailable === false + ? null + : () => { + if (options.throwBoundaryCallbacks === true) { + throw new Error('thread CPU callback unavailable'); + } + const value = cpuValues[cpuIndex] ?? cpuValues.at(-1)!; + cpuIndex += 1; + lifecycle.push(`cpu:${cpuIndex}`); + return value; + }, + scheduleTimeout: (callback, delayMs) => { + lifecycle.push(`timeout:${delayMs}`); + scheduledTimeouts.push(delayMs); + if (options.throwScheduleTimeout === true) { + throw new Error('timer callback unavailable'); + } + if (options.stallMonotonicClock !== true) { + monotonicMs += delayMs; + } else if (timeoutCount >= 55) { + throw new Error('test scheduler fail-safe'); + } + timeoutCount += 1; + if (timeoutCount === 1 && options.armHistogram !== false) { + histogramCount += 1; + } else if ( + timeoutCount > 1 && + options.armHistogram !== false && + options.flushHistogram !== false + ) { + histogramCount += 1; + } + callback(); + }, + }; + + return { + histogram, + histogramLifecycle, + lifecycle, + runtime, + scheduledTimeouts, + }; +} diff --git a/apps/electron-backend/src/app/workers/worker-performance-capture.ts b/apps/electron-backend/src/app/workers/worker-performance-capture.ts index db86f234d..a29d97f0a 100644 --- a/apps/electron-backend/src/app/workers/worker-performance-capture.ts +++ b/apps/electron-backend/src/app/workers/worker-performance-capture.ts @@ -1,47 +1,119 @@ import { - monitorEventLoopDelay, - performance, - type IntervalHistogram, -} from 'node:perf_hooks'; + WORKER_PERFORMANCE_INVALID_REASON, + WORKER_PERFORMANCE_UNAVAILABLE_REASON, + type WorkerPerformanceCapture, + type WorkerPerformanceCaptureOptions, + type WorkerPerformanceCaptureResult, + type WorkerPerformanceExecutionResult, +} from './worker-performance-capture.model'; +import { + disableWorkerPerformanceHistogram, + markEventLoopDelayUnavailable, + waitForHistogramCondition, +} from './worker-performance-capture.histogram'; +import { + beginWorkerPerformanceCapture, + createWorkerPerformanceCaptureResult, + finishWorkerPerformanceMetricBoundaries, +} from './worker-performance-capture.metrics'; +import { + DEFAULT_WORKER_PERFORMANCE_RUNTIME, + safeReadEpochMs, + WORKER_PROFILING_ENV, +} from './worker-performance-capture.runtime'; -const WORKER_PROFILING_ENV = 'IPTVNATOR_PERF_WORKER_PROFILING'; +export { + WORKER_PERFORMANCE_INVALID_REASON, + WORKER_PERFORMANCE_UNAVAILABLE_REASON, +} from './worker-performance-capture.model'; +export type { + WorkerEventLoopDelayHistogram, + WorkerEventLoopDelayMetrics, + WorkerPerformanceCapture, + WorkerPerformanceCaptureOptions, + WorkerPerformanceCaptureResult, + WorkerPerformanceCaptureRuntime, + WorkerPerformanceExecutionResult, + WorkerPerformanceInvalidReason, + WorkerPerformanceUnavailableReason, +} from './worker-performance-capture.model'; -export interface WorkerPerformanceCaptureResult { - eventLoopDelay: { - maxMs: number; - p95Ms: number; - p99Ms: number; - }; - eventLoopUtilization: number; -} - -interface WorkerPerformanceCapture { - eventLoopDelay: IntervalHistogram; - eventLoopUtilizationStart: ReturnType< - typeof performance.eventLoopUtilization - >; -} - -export function startWorkerPerformanceCapture(): WorkerPerformanceCapture | null { - if (process.env[WORKER_PROFILING_ENV] !== '1') { +export function startWorkerPerformanceCapture( + options: WorkerPerformanceCaptureOptions = {} +): WorkerPerformanceCapture | null { + const enabled = + options.enabled ?? process.env[WORKER_PROFILING_ENV] === '1'; + if (!enabled) { return null; } - const eventLoopDelay = monitorEventLoopDelay({ resolution: 1 }); - eventLoopDelay.enable(); - return { - eventLoopDelay, - eventLoopUtilizationStart: performance.eventLoopUtilization(), + const runtime = options.runtime ?? DEFAULT_WORKER_PERFORMANCE_RUNTIME; + const capture: WorkerPerformanceCapture = { + eventLoopDelay: null, + eventLoopDelayUnavailableReason: null, + eventLoopUtilizationEnd: null, + eventLoopUtilizationStart: null, + eventLoopUtilizationUnavailableReason: null, + histogramDisabled: false, + histogramFlushedEpochMs: null, + invalidReason: null, + requestReceivedEpochMs: safeReadEpochMs(runtime), + runtime, + threadCpuEnd: null, + threadCpuStart: null, + threadCpuUnavailableReason: + runtime.readThreadCpuUsage === null + ? WORKER_PERFORMANCE_UNAVAILABLE_REASON.THREAD_CPU_USAGE_UNAVAILABLE + : null, + workEndedEpochMs: null, + workStartedEpochMs: null, }; + + try { + capture.eventLoopDelay = runtime.createEventLoopDelayHistogram(); + capture.eventLoopDelay.enable(); + } catch { + capture.eventLoopDelay = null; + capture.eventLoopDelayUnavailableReason = + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_CAPTURE_UNAVAILABLE; + } + + return capture; } export async function armWorkerPerformanceCapture( capture: WorkerPerformanceCapture | null ): Promise { - if (capture) { - await new Promise((resolvePromise) => - setTimeout(resolvePromise, 2) + if (!capture) { + return; + } + + try { + if ( + capture.invalidReason !== null || + capture.eventLoopDelayUnavailableReason !== null + ) { + return; + } + const armed = await waitForHistogramCondition( + capture, + (histogram) => histogram.count > 0 ); + if (!armed && capture.invalidReason === null) { + markEventLoopDelayUnavailable( + capture, + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_ARM_TIMEOUT + ); + } + } catch { + try { + markEventLoopDelayUnavailable( + capture, + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_CAPTURE_UNAVAILABLE + ); + } catch { + // Instrumentation must never prevent business execution. + } } } @@ -52,17 +124,148 @@ export async function finishWorkerPerformanceCapture( return undefined; } - await new Promise((resolvePromise) => setTimeout(resolvePromise, 2)); - capture.eventLoopDelay.disable(); - const eventLoopUtilization = performance.eventLoopUtilization( - capture.eventLoopUtilizationStart + finishWorkerPerformanceMetricBoundaries(capture); + + if ( + capture.invalidReason === null && + capture.eventLoopDelayUnavailableReason === null && + capture.eventLoopDelay + ) { + try { + const workEndedHistogramCount = capture.eventLoopDelay.count; + const flushed = await waitForHistogramCondition( + capture, + (histogram) => histogram.count > workEndedHistogramCount + ); + if (flushed) { + const disabled = disableWorkerPerformanceHistogram(capture); + if (disabled) { + capture.histogramFlushedEpochMs = safeReadEpochMs( + capture.runtime + ); + } else { + markEventLoopDelayUnavailable( + capture, + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_CAPTURE_UNAVAILABLE + ); + } + } else if ( + capture.invalidReason === null && + capture.eventLoopDelayUnavailableReason === null + ) { + markEventLoopDelayUnavailable( + capture, + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_FLUSH_TIMEOUT + ); + } + } catch { + markEventLoopDelayUnavailable( + capture, + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_CAPTURE_UNAVAILABLE + ); + } + } else { + disableWorkerPerformanceHistogram(capture); + } + + return createWorkerPerformanceCaptureResult(capture); +} + +export async function executeWithWorkerPerformanceCapture( + capture: WorkerPerformanceCapture | null, + execute: () => Promise, + finish: ( + capture: WorkerPerformanceCapture | null + ) => Promise< + WorkerPerformanceCaptureResult | undefined + > = finishWorkerPerformanceCapture +): Promise> { + try { + beginWorkerPerformanceCapture(capture); + } catch { + // Instrumentation setup must never prevent business execution. + } + let result: TResult; + try { + result = await execute(); + } catch (error) { + const performance = await safelyFinishWorkerPerformanceCapture( + capture, + finish + ); + return { + error, + performance, + result: undefined, + success: false, + }; + } + + const performance = await safelyFinishWorkerPerformanceCapture( + capture, + finish ); return { - eventLoopDelay: { - maxMs: capture.eventLoopDelay.max / 1e6, - p95Ms: capture.eventLoopDelay.percentile(95) / 1e6, - p99Ms: capture.eventLoopDelay.percentile(99) / 1e6, - }, - eventLoopUtilization: eventLoopUtilization.utilization, + error: null, + performance, + result, + success: true, }; } + +async function safelyFinishWorkerPerformanceCapture( + capture: WorkerPerformanceCapture | null, + finish: ( + capture: WorkerPerformanceCapture | null + ) => Promise +): Promise { + try { + return await finish(capture); + } catch { + if (!capture) { + return undefined; + } + try { + if (capture.workEndedEpochMs === null) { + finishWorkerPerformanceMetricBoundaries(capture); + } + markEventLoopDelayUnavailable( + capture, + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_CAPTURE_UNAVAILABLE + ); + return createWorkerPerformanceCaptureResult(capture); + } catch { + return undefined; + } + } +} + +export function registerDatabaseWorkerPerformanceCapture( + activeCaptures: Set, + capture: WorkerPerformanceCapture | null +): void { + if (!capture) { + return; + } + + if (activeCaptures.size > 0) { + for (const activeCapture of activeCaptures) { + activeCapture.invalidReason = + WORKER_PERFORMANCE_INVALID_REASON.OVERLAPPING_DATABASE_WORKER_REQUESTS; + disableWorkerPerformanceHistogram(activeCapture); + } + capture.invalidReason = + WORKER_PERFORMANCE_INVALID_REASON.OVERLAPPING_DATABASE_WORKER_REQUESTS; + disableWorkerPerformanceHistogram(capture); + } + activeCaptures.add(capture); +} + +export function releaseDatabaseWorkerPerformanceCapture( + activeCaptures: Set, + capture: WorkerPerformanceCapture | null +): void { + if (capture) { + activeCaptures.delete(capture); + } +} diff --git a/apps/electron-backend/src/app/workers/worker-performance-integration.spec.ts b/apps/electron-backend/src/app/workers/worker-performance-integration.spec.ts new file mode 100644 index 000000000..e13a11f29 --- /dev/null +++ b/apps/electron-backend/src/app/workers/worker-performance-integration.spec.ts @@ -0,0 +1,118 @@ +import { readFileSync } from 'node:fs'; +import { join } from 'node:path'; +import { + WORKER_PERFORMANCE_UNAVAILABLE_REASON, + armWorkerPerformanceCapture, + executeWithWorkerPerformanceCapture, + startWorkerPerformanceCapture, +} from './worker-performance-capture'; +import { createFakeRuntime } from './worker-performance-capture.test-harness'; + +function readWorkerSource(filename: string): string { + return readFileSync( + join(process.cwd(), 'apps/electron-backend/src/app/workers', filename), + 'utf8' + ); +} + +describe('worker performance integration order', () => { + it('profiles overlapping database requests without serializing their execution', () => { + const source = readWorkerSource('database.worker.ts'); + const handler = source.slice( + source.lastIndexOf("parentPort.on('message'") + ); + + expect(handler).toContain('registerDatabaseWorkerPerformanceCapture'); + expect(handler).toContain('executeWithWorkerPerformanceCapture'); + expect(handler).toContain('releaseDatabaseWorkerPerformanceCapture'); + expect(handler).toContain('finally'); + expect(handler).not.toContain('finishWorkerPerformanceCapture'); + expect( + handler.indexOf('executeWithWorkerPerformanceCapture') + ).toBeLessThan(handler.indexOf('console.error')); + expect( + handler.indexOf('executeWithWorkerPerformanceCapture') + ).toBeLessThan(handler.indexOf('postMessage')); + expect(handler).not.toContain('await previous'); + expect(handler).not.toContain('requestQueue'); + }); + + it('finishes playlist profiling before cancellation/error events and response serialization', () => { + const source = readWorkerSource('playlist-refresh.worker.ts'); + const handler = source.slice(source.lastIndexOf('parentPort.on(')); + const executionIndex = handler.indexOf( + 'executeWithWorkerPerformanceCapture' + ); + + expect(executionIndex).toBeGreaterThanOrEqual(0); + expect(handler).not.toContain('finishWorkerPerformanceCapture'); + expect(executionIndex).toBeLessThan( + handler.indexOf('emitEvent', executionIndex) + ); + expect(executionIndex).toBeLessThan( + handler.indexOf('serializeError', executionIndex) + ); + expect(executionIndex).toBeLessThan( + handler.indexOf('postMessage', executionIndex) + ); + }); + + it('preserves success and business errors when the finish seam throws once', async () => { + const successHarness = createFakeRuntime(); + const successCapture = startWorkerPerformanceCapture({ + enabled: true, + runtime: successHarness.runtime, + }); + await armWorkerPerformanceCapture(successCapture); + const failingSuccessFinish = jest.fn(async () => { + throw new Error('finish failed'); + }); + + const success = await executeWithWorkerPerformanceCapture( + successCapture, + async () => 'business-result', + failingSuccessFinish + ); + + expect(success).toMatchObject({ + error: null, + result: 'business-result', + success: true, + performance: { + eventLoopDelay: null, + eventLoopDelayUnavailableReason: + WORKER_PERFORMANCE_UNAVAILABLE_REASON.EVENT_LOOP_DELAY_CAPTURE_UNAVAILABLE, + eventLoopUtilization: 0.75, + threadCpuSystemMicros: 30, + threadCpuUserMicros: 80, + }, + }); + expect(failingSuccessFinish).toHaveBeenCalledTimes(1); + + const errorHarness = createFakeRuntime(); + const errorCapture = startWorkerPerformanceCapture({ + enabled: true, + runtime: errorHarness.runtime, + }); + await armWorkerPerformanceCapture(errorCapture); + const businessError = new Error('business failed'); + const failingErrorFinish = jest.fn(async () => { + throw new Error('finish failed'); + }); + + const failure = await executeWithWorkerPerformanceCapture( + errorCapture, + async () => { + throw businessError; + }, + failingErrorFinish + ); + + expect(failure).toMatchObject({ + error: businessError, + result: undefined, + success: false, + }); + expect(failingErrorFinish).toHaveBeenCalledTimes(1); + }); +}); diff --git a/docs/architecture/m3u-playlist-module.md b/docs/architecture/m3u-playlist-module.md index 48a4a468b..855da02a3 100644 --- a/docs/architecture/m3u-playlist-module.md +++ b/docs/architecture/m3u-playlist-module.md @@ -121,11 +121,27 @@ before/after claims. The target reserves and verifies CDP port 9222, freezes renderer long-task, frame-gap, and heartbeat probes before forced post-GC heap collection, and -enables opt-in worker profiling. Worker event-loop delay is read from a -request-scoped `node:perf_hooks` capture; a worker terminated before it can flush -the capture reports the metric as unavailable rather than zero. Diagnostic CPU -profiles, heap snapshots, and Chromium traces are excluded from the five-run -headline distributions. +enables opt-in worker profiling. Each worker response retains a raw +request-scoped record containing request/operation identity, received/work/flush +timestamps, thread CPU, event-loop utilization, event-loop delay, and fixed +unavailability or invalid reasons. Missing or malformed profiling metadata +still produces an outcome with its request identity, nullable metrics, and a +fixed capture-unavailable reason. A worker terminated before it can flush the +capture reports the metric as unavailable rather than zero. + +The harness excludes workers created by pre-capture seed setup through an exact +capture-generation marker. Since the production M3U database payload +deliberately omits the refresh operation ID, the benchmark correlates only a +valid preload `start` marker to the exact database operation and playlist. A +missing or ambiguous marker leaves the raw operation ID `null` with a fixed +reason; it does not infer identity from timing alone. + +Headline worker fields named `*WorkerRequest*` are distributions of individual +request metrics. In particular, request p95/p99 distributions are not presented +as an operation-wide or process-wide percentile, and database request +percentiles are never collapsed with `Math.max`. Diagnostic CPU profiles, heap +snapshots, and Chromium traces are excluded from the five-run headline +distributions. ### Reporting The Auto-Update Result diff --git a/docs/architecture/sqlite-db-worker.md b/docs/architecture/sqlite-db-worker.md index 0e2a3705b..52bdbb860 100644 --- a/docs/architecture/sqlite-db-worker.md +++ b/docs/architecture/sqlite-db-worker.md @@ -118,6 +118,35 @@ The worker contract lives in 3. `DbWorkerEventMessage` 4. `DbOperationEvent` +### Opt-in request performance capture + +`IPTVNATOR_PERF_WORKER_PROFILING=1` adds development/test-only performance +metadata to each database-worker response. It is disabled by default and must +stay disabled for production launches. + +Each enabled request gets a fresh event-loop-delay histogram and records: + +- `requestReceivedEpochMs`, `workStartedEpochMs`, `workEndedEpochMs`, and + `histogramFlushedEpochMs` +- worker-thread CPU user/system microseconds from `process.threadCpuUsage()` +- event-loop utilization across the exact work interval +- event-loop-delay max/p95/p99 from the request's own histogram +- fixed invalid or unavailable reasons whenever a metric cannot be attributed + +Histogram arming waits until the histogram has a sample; flushing waits for its +sample count to advance after work ends. Both waits use condition-based timer +polling. Each wait stops after 50 ms of observed monotonic time or its bounded +poll count; arming and flushing have separate caps, and timer scheduling may +overshoot wall-clock time. A timeout or profiling API failure never replaces +the business response: timestamps and any independently available CPU/ELU +metrics remain valid, while event-loop delay is `null` with a fixed reason. + +The long-lived database worker still executes concurrent requests without a +profiling queue. If captures overlap, every overlapping response carries +`invalidReason: "overlapping-database-worker-requests"` and all attributable +CPU, ELU, and event-loop-delay values are `null`. This avoids assigning shared +worker activity to one request while preserving normal worker concurrency. + ### Progress event contract The worker now emits request-scoped events with: