diff --git a/AGENTS.md b/AGENTS.md index c1b428902..2e731d993 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, 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 + - `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, plus the database worker's idle-only one-shot post-GC heap probe; 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 bd872ab79..bda15edcc 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, 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 +- `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, plus the database worker's idle-only one-shot post-GC heap probe; 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 @@ -353,10 +353,11 @@ Key patterns: - **Factory injection**: `provideXtreamDataSource()` selects Electron or PWA implementation at runtime Data strategies by environment: -| Environment | Strategy | -|-------------|----------| + +| Environment | Strategy | +| ------------ | ------------------------------------------------------- | | **Electron** | DB-first: Check DB → fetch API if missing → cache to DB | -| **PWA** | API-only: Always fetch from API, store in memory | +| **PWA** | API-only: Always fetch from API, store in memory | **M3U Playlist Module Architecture**: diff --git a/apps/electron-backend-e2e/src/database-worker-post-gc.e2e.ts b/apps/electron-backend-e2e/src/database-worker-post-gc.e2e.ts new file mode 100644 index 000000000..14c2dc93e --- /dev/null +++ b/apps/electron-backend-e2e/src/database-worker-post-gc.e2e.ts @@ -0,0 +1,122 @@ +import type { ElectronApplication, Page } from '@playwright/test'; + +import { + closeElectronApp, + expect, + launchElectronApp, + test, + type LaunchedElectronApp, +} from './electron-test-fixtures'; +import { + installMainCapture, + rolloverMainCapture, + startMainCapture, + stopMainCapture, +} from './performance/m3u-refresh-main-capture'; +import type { MainCaptureMetrics } from './performance/m3u-refresh-cancellation-contract'; +import type { RendererWindowIdentity } from './performance/renderer-window-rss-session'; + +test('captures the exact database worker post-GC heap in built Electron', async ({ + dataDir, +}) => { + const launchedApps: LaunchedElectronApp[] = []; + + try { + const app = await launchElectronApp(dataDir, { + args: ['--js-flags=--expose-gc'], + env: { + IPTVNATOR_DB_WORKER_BATCH_DELAY_MS: '0', + IPTVNATOR_PERF_CAPTURE: '1', + IPTVNATOR_PERF_WORKER_PROFILING: '1', + }, + }); + launchedApps.push(app); + await installMainCapture(app.electronApp); + const rendererWindowIdentity = await resolveRendererWindowIdentity( + app.electronApp, + app.mainWindow + ); + const captureOptions = { + diagnostic: false, + outputDirectory: dataDir, + rendererWindowIdentity, + } as const; + + await startMainCapture(app.electronApp, captureOptions); + await requestDatabaseWorker(app); + const rollover = await rolloverMainCapture( + app.electronApp, + captureOptions + ); + expect(rollover.nextCaptureStarted).toBe(true); + expect(rollover.nextCaptureUnavailableReason).toBeNull(); + const firstDatabaseWorker = expectExactDatabaseWorker( + rollover.completedCapture + ); + + await requestDatabaseWorker(app); + const secondDatabaseWorker = expectExactDatabaseWorker( + await stopMainCapture(app.electronApp) + ); + expect(secondDatabaseWorker.ordinal).toBe(firstDatabaseWorker.ordinal); + + await startMainCapture(app.electronApp, captureOptions); + const noDatabasePhase = await stopMainCapture(app.electronApp); + expect( + noDatabasePhase.workers.filter( + (worker) => worker.kind === 'database.worker' + ) + ).toHaveLength(0); + } finally { + await Promise.all(launchedApps.map((app) => closeElectronApp(app))); + } +}); + +async function requestDatabaseWorker(app: LaunchedElectronApp): Promise { + await app.mainWindow.evaluate(() => window.electron.dbGetAppPlaylists()); +} + +function expectExactDatabaseWorker(capture: MainCaptureMetrics) { + const databaseWorkers = capture.workers.filter( + (worker) => worker.kind === 'database.worker' + ); + + expect(databaseWorkers).toHaveLength(1); + const databaseWorker = databaseWorkers[0]; + expect(databaseWorker?.ordinal).toBeGreaterThan(0); + expect(databaseWorker?.postGcHeapUnavailableReason).toBeNull(); + expect(databaseWorker?.postGcHeapUsedBytes).toBeGreaterThan(0); + expect(databaseWorker?.requests).toHaveLength(1); + expect(databaseWorker?.requests[0]).toEqual( + expect.objectContaining({ + invalidReason: null, + operation: 'DB_GET_APP_PLAYLISTS', + performanceCaptureUnavailableReason: null, + playlistId: null, + requestId: expect.any(String), + success: true, + }) + ); + return databaseWorker; +} + +async function resolveRendererWindowIdentity( + electronApp: ElectronApplication, + page: Page +): Promise { + const browserWindowHandle = await electronApp.browserWindow(page); + try { + return await browserWindowHandle.evaluate((browserWindow) => { + const exactWindow = browserWindow as unknown as { + readonly id: number; + readonly webContents: { readonly id: number }; + }; + return { + browserWindowId: exactWindow.id, + webContentsId: exactWindow.webContents.id, + }; + }); + } finally { + await browserWindowHandle.dispose(); + } +} diff --git a/apps/electron-backend-e2e/src/performance/database-worker-post-gc-cutoff.spec.ts b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-cutoff.spec.ts new file mode 100644 index 000000000..631081fd4 --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-cutoff.spec.ts @@ -0,0 +1,116 @@ +/* eslint-disable playwright/expect-expect -- This is a Node assertion-based performance contract test. */ +import assert from 'node:assert/strict'; +import { readFileSync } from 'node:fs'; +import test from 'node:test'; + +type RequestDisposition = 'after-cutoff' | 'capture' | 'outside-capture'; + +interface CutoffApi { + beginCapture(): void; + beginStop(): void; + finishStop(): void; + observeDatabaseRequest(): RequestDisposition; + rolloverCapture?(): void; + snapshot(): { + readonly lateRequestCount: number; + readonly phase: 'active' | 'idle' | 'stopping'; + }; +} + +interface CutoffModule { + createDatabaseWorkerPostGcCutoffApi?: () => CutoffApi; +} + +const cutoffModulePromise = import( + new URL('./database-worker-post-gc-cutoff.ts', import.meta.url).href +) + .then((module) => module as CutoffModule) + .catch(() => null); + +async function restoreSerializableApi(): Promise { + const module = await cutoffModulePromise; + assert.ok(module, 'database worker post-GC cutoff module must exist'); + const factory = module.createDatabaseWorkerPostGcCutoffApi; + assert.equal(typeof factory, 'function'); + + const source = factory.toString(); + assert.doesNotMatch(source, /__name/); + const restoredFactory = Function( + `"use strict"; return (${source});` + )() as () => CutoffApi; + return restoredFactory(); +} + +test('classifies requests after the synchronous capture cutoff without carrying state across generations', async () => { + const api = await restoreSerializableApi(); + + assert.equal(api.observeDatabaseRequest(), 'outside-capture'); + api.beginCapture(); + assert.equal(api.observeDatabaseRequest(), 'capture'); + api.beginStop(); + assert.equal(api.observeDatabaseRequest(), 'after-cutoff'); + assert.equal(api.observeDatabaseRequest(), 'after-cutoff'); + assert.deepEqual(api.snapshot(), { + lateRequestCount: 2, + phase: 'stopping', + }); + api.finishStop(); + assert.equal(api.observeDatabaseRequest(), 'outside-capture'); + + api.beginCapture(); + assert.deepEqual(api.snapshot(), { + lateRequestCount: 0, + phase: 'active', + }); +}); + +test('atomically rolls a clean cutoff into the next generation without erasing contamination', async () => { + const clean = await restoreSerializableApi(); + assert.equal(typeof clean.rolloverCapture, 'function'); + clean.beginCapture(); + clean.beginStop(); + clean.rolloverCapture?.(); + assert.deepEqual(clean.snapshot(), { + lateRequestCount: 0, + phase: 'active', + }); + assert.equal(clean.observeDatabaseRequest(), 'capture'); + + const contaminated = await restoreSerializableApi(); + assert.equal(typeof contaminated.rolloverCapture, 'function'); + contaminated.beginCapture(); + contaminated.beginStop(); + assert.equal(contaminated.observeDatabaseRequest(), 'after-cutoff'); + assert.throws( + () => contaminated.rolloverCapture?.(), + /database-worker-post-gc-capture-contaminated/ + ); + assert.deepEqual(contaminated.snapshot(), { + lateRequestCount: 1, + phase: 'stopping', + }); +}); + +test('main capture arms cutoff before awaiting and rejects late DB work before startWorker can restart profiling', () => { + const source = readFileSync( + new URL('./m3u-refresh-main-capture.ts', import.meta.url), + 'utf8' + ); + const stopStart = source.indexOf('stop: async ('); + const cutoffStart = source.indexOf( + 'databaseWorkerPostGcCutoffApi.beginStop()', + stopStart + ); + const activeFalse = source.indexOf('state.active = false', stopStart); + const firstAwait = source.indexOf('await ', stopStart); + const observeRequest = source.indexOf( + 'databaseWorkerPostGcCutoffApi.observeDatabaseRequest()' + ); + const startWorker = source.indexOf('startWorker(record)', observeRequest); + + assert.ok(stopStart >= 0); + assert.ok(cutoffStart > stopStart && cutoffStart < firstAwait); + assert.ok(activeFalse > stopStart && activeFalse < firstAwait); + assert.ok(observeRequest >= 0 && observeRequest < startWorker); + assert.match(source, /database-worker-activity-after-cutoff/); +}); diff --git a/apps/electron-backend-e2e/src/performance/database-worker-post-gc-cutoff.ts b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-cutoff.ts new file mode 100644 index 000000000..997963bdb --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-cutoff.ts @@ -0,0 +1,74 @@ +export type DatabaseWorkerPostGcRequestDisposition = + 'after-cutoff' | 'capture' | 'outside-capture'; + +export interface DatabaseWorkerPostGcCutoffSnapshot { + readonly lateRequestCount: number; + readonly phase: 'active' | 'idle' | 'stopping'; +} + +export interface DatabaseWorkerPostGcCutoffApi { + beginCapture(): void; + beginStop(): void; + finishStop(): void; + observeDatabaseRequest(): DatabaseWorkerPostGcRequestDisposition; + rolloverCapture(): void; + snapshot(): DatabaseWorkerPostGcCutoffSnapshot; +} + +export function createDatabaseWorkerPostGcCutoffApi(): DatabaseWorkerPostGcCutoffApi { + let lateRequestCount = 0; + let phase: DatabaseWorkerPostGcCutoffSnapshot['phase'] = 'idle'; + + const api: DatabaseWorkerPostGcCutoffApi = { + beginCapture(): void { + if (phase !== 'idle') { + throw new Error( + phase === 'stopping' + ? 'database-worker-post-gc-stop-in-progress' + : 'database-worker-post-gc-capture-already-active' + ); + } + lateRequestCount = 0; + phase = 'active'; + }, + + beginStop(): void { + if (phase !== 'active') { + throw new Error('database-worker-post-gc-capture-not-active'); + } + phase = 'stopping'; + }, + + finishStop(): void { + phase = 'idle'; + }, + + observeDatabaseRequest(): DatabaseWorkerPostGcRequestDisposition { + if (phase === 'active') { + return 'capture'; + } + if (phase === 'stopping') { + lateRequestCount += 1; + return 'after-cutoff'; + } + return 'outside-capture'; + }, + + rolloverCapture(): void { + if (phase !== 'stopping') { + throw new Error('database-worker-post-gc-stop-not-in-progress'); + } + if (lateRequestCount > 0) { + throw new Error('database-worker-post-gc-capture-contaminated'); + } + lateRequestCount = 0; + phase = 'active'; + }, + + snapshot(): DatabaseWorkerPostGcCutoffSnapshot { + return Object.freeze({ lateRequestCount, phase }); + }, + }; + + return Object.freeze(api); +} diff --git a/apps/electron-backend-e2e/src/performance/database-worker-post-gc-finalization.spec.ts b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-finalization.spec.ts new file mode 100644 index 000000000..bcd3ac969 --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-finalization.spec.ts @@ -0,0 +1,392 @@ +/* eslint-disable playwright/expect-expect -- This is a Node assertion-based performance contract test. */ +import assert from 'node:assert/strict'; +import { readFileSync } from 'node:fs'; +import test from 'node:test'; + +type AncillaryFailureStage = 'heap-snapshot' | 'post-gc-probe' | 'profile-stop'; + +type PostGcOutcome = + | { + readonly postGcHeapUsedBytes: number; + readonly unavailableReason: null; + } + | { + readonly postGcHeapUsedBytes: null; + readonly unavailableReason: string; + }; + +interface FinalizationInput { + readonly finalizationKey: object; + readonly joinFinalSample: () => Promise; + readonly probePostGc: () => Promise; + readonly reportAncillaryFailure: ( + stage: AncillaryFailureStage, + error: unknown + ) => void; + readonly stopProfile: () => Promise; + readonly stopSampling: () => void; + readonly takeHeapSnapshot?: () => Promise; +} + +interface FinalizationApi { + finalize(input: FinalizationInput): Promise; +} + +interface FinalizationModule { + createDatabaseWorkerPostGcFinalizationApi?: () => FinalizationApi; +} + +const finalizationModulePromise = import( + new URL('./database-worker-post-gc-finalization.ts', import.meta.url).href +) + .then((module) => module as FinalizationModule) + .catch(() => null); + +function deferred(): { + readonly promise: Promise; + readonly resolve: () => void; +} { + let resolve!: () => void; + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise; + }); + return { promise, resolve }; +} + +async function restoreSerializableApi(): Promise { + const module = await finalizationModulePromise; + assert.ok(module, 'database worker post-GC finalization module must exist'); + const factory = module.createDatabaseWorkerPostGcFinalizationApi; + assert.equal(typeof factory, 'function'); + + const source = factory.toString(); + assert.doesNotMatch(source, /__name/); + const restoredFactory = Function( + `"use strict"; return (${source});` + )() as () => FinalizationApi; + return restoredFactory(); +} + +test('coordinates database worker post-GC finalization in order and single-flight', async () => { + const api = await restoreSerializableApi(); + const finalSampleGate = deferred(); + const events: string[] = []; + let probeCalls = 0; + const successfulPostGc = { + postGcHeapUsedBytes: 98_304, + unavailableReason: null, + } as const; + const input: FinalizationInput = { + finalizationKey: {}, + joinFinalSample: async () => { + events.push('final-sample:start'); + await finalSampleGate.promise; + events.push('final-sample:joined'); + }, + probePostGc: async () => { + probeCalls += 1; + events.push('post-gc-probe'); + return successfulPostGc; + }, + reportAncillaryFailure: () => { + assert.fail('the successful path has no ancillary failures'); + }, + stopProfile: async () => { + events.push('profile-stop'); + }, + stopSampling: () => { + events.push('sampling-stop'); + }, + takeHeapSnapshot: async () => { + events.push('heap-snapshot'); + }, + }; + + const firstFinalization = api.finalize(input); + const concurrentFinalization = api.finalize(input); + await Promise.resolve(); + + assert.deepEqual(events, ['sampling-stop', 'final-sample:start']); + assert.equal(probeCalls, 0); + + finalSampleGate.resolve(); + const [firstResult, concurrentResult] = await Promise.all([ + firstFinalization, + concurrentFinalization, + ]); + + assert.deepEqual(events, [ + 'sampling-stop', + 'final-sample:start', + 'final-sample:joined', + 'profile-stop', + 'post-gc-probe', + 'heap-snapshot', + ]); + assert.equal(probeCalls, 1); + assert.deepEqual(firstResult, successfulPostGc); + assert.deepEqual(concurrentResult, successfulPostGc); +}); + +test('reports profile and snapshot failures without erasing successful post-GC', async () => { + const api = await restoreSerializableApi(); + const events: string[] = []; + const profileError = new Error('profile-stop-failed'); + const snapshotError = new Error('heap-snapshot-failed'); + const failures: { + readonly error: unknown; + readonly stage: AncillaryFailureStage; + }[] = []; + + const result = await api.finalize({ + finalizationKey: {}, + joinFinalSample: async () => { + events.push('final-sample:joined'); + }, + probePostGc: async () => { + events.push('post-gc-probe'); + return { + postGcHeapUsedBytes: 131_072, + unavailableReason: null, + }; + }, + reportAncillaryFailure: (stage, error) => { + failures.push({ error, stage }); + }, + stopProfile: async () => { + events.push('profile-stop'); + throw profileError; + }, + stopSampling: () => { + events.push('sampling-stop'); + }, + takeHeapSnapshot: async () => { + events.push('heap-snapshot'); + throw snapshotError; + }, + }); + + assert.deepEqual(events, [ + 'sampling-stop', + 'final-sample:joined', + 'profile-stop', + 'post-gc-probe', + 'heap-snapshot', + ]); + assert.deepEqual(failures, [ + { error: profileError, stage: 'profile-stop' }, + { error: snapshotError, stage: 'heap-snapshot' }, + ]); + assert.deepEqual(result, { + postGcHeapUsedBytes: 131_072, + unavailableReason: null, + }); +}); + +test('fails closed when the post-GC probe throws without relabelling profile artifacts', async () => { + const api = await restoreSerializableApi(); + const probeError = new Error('probe-failed'); + const failures: { + readonly error: unknown; + readonly stage: AncillaryFailureStage; + }[] = []; + + const result = await api.finalize({ + finalizationKey: {}, + joinFinalSample: async () => undefined, + probePostGc: async () => { + throw probeError; + }, + reportAncillaryFailure: (stage, error) => { + failures.push({ error, stage }); + }, + stopProfile: async () => undefined, + stopSampling: () => undefined, + }); + + assert.deepEqual(result, { + postGcHeapUsedBytes: null, + unavailableReason: 'capture-failed', + }); + assert.deepEqual(failures, [{ error: probeError, stage: 'post-gc-probe' }]); +}); + +test('the Electron main capture wires exact DB selection and explicit-GC finalization', () => { + const source = readFileSync( + new URL('./m3u-refresh-main-capture.ts', import.meta.url), + 'utf8' + ); + + assert.match( + source, + /databaseWorkerPostGcSelectionApiFactorySource:\s+createDatabaseWorkerPostGcSelectionApi\.toString\(\)/ + ); + assert.match( + source, + /databaseWorkerPostGcProbeApiFactorySource:\s+createDatabaseWorkerPostGcProbeApi\.toString\(\)/ + ); + assert.match( + source, + /databaseWorkerPostGcFinalizationApiFactorySource:\s+createDatabaseWorkerPostGcFinalizationApi\.toString\(\)/ + ); + assert.match(source, /new workerThreads\.MessageChannel\(\)/); + assert.match( + source, + /databaseWorkerPostGcSelectionApi\.select\(\s*currentWorkerRecords,\s*state\.captureGeneration\s*\)/ + ); + assert.match(source, /databaseWorkerPostGcFinalizationApi\s*\.finalize\(/); + assert.match(source, /pendingCount: 0/); + assert.match(source, /record\.pendingCount \+= 1/); + assert.match( + source, + /request\.record\.pendingCount = Math\.max\(\s*0,\s*request\.record\.pendingCount - 1\s*\)/ + ); + assert.match(source, /samplePromise/); + assert.match(source, /postGcHeapUnavailableReason/); + assert.match(source, /ordinal: nextWorkerOrdinal/); + assert.match( + source, + /`\$\{record\.kind\}-\$\{record\.ordinal\}\.cpuprofile`/ + ); + assert.doesNotMatch(source, /sampleBusy/); + assert.doesNotMatch(source, /const postSnapshot/); +}); + +test('main CPU profiling stops at the operation cutoff before worker artifact finalization', () => { + const source = readFileSync( + new URL('./m3u-refresh-main-capture.ts', import.meta.url), + 'utf8' + ); + const stopStart = source.indexOf('stop: async ('); + const mainProfileStop = source.indexOf("'Profiler.stop'", stopStart); + const databaseSelection = source.indexOf( + 'databaseWorkerPostGcSelectionApi.select', + stopStart + ); + const databaseFinalization = source.indexOf( + 'await finalizeDatabaseWorker', + stopStart + ); + + assert.ok(stopStart >= 0); + assert.ok(mainProfileStop > stopStart); + assert.ok(mainProfileStop < databaseSelection); + assert.ok(mainProfileStop < databaseFinalization); +}); + +test('worker CPU profile serialization runs only after the main CPU profiler stops', () => { + const source = readFileSync( + new URL('./m3u-refresh-main-capture.ts', import.meta.url), + 'utf8' + ); + const profileStopStart = source.indexOf('const stopWorkerProfile'); + const profileStopEnd = source.indexOf( + 'const writeWorkerProfile', + profileStopStart + ); + const profileStop = source.slice(profileStopStart, profileStopEnd); + const stopStart = source.indexOf('stop: async ('); + const mainProfileStop = source.indexOf("'Profiler.stop'", stopStart); + const workerProfileFlush = source.indexOf( + 'flushWorkerProfiles(currentWorkerRecords)', + mainProfileStop + ); + + assert.match(profileStop, /const profileResult = await handle\.stop\(\)/); + assert.match(profileStop, /record\.profileResult = profileResult/); + assert.doesNotMatch(profileStop, /writeFileSync|JSON\.stringify/); + assert.ok(mainProfileStop > stopStart); + assert.ok(workerProfileFlush > mainProfileStop); +}); + +test('measured termination stays unperturbed while diagnostic capture finalizes its profile before termination', () => { + const source = readFileSync( + new URL('./m3u-refresh-main-capture.ts', import.meta.url), + 'utf8' + ); + const terminateStart = source.indexOf( + 'WorkerClass.prototype.terminate = function' + ); + const terminateEnd = source.indexOf('const inspectorPost', terminateStart); + const terminateBlock = source.slice(terminateStart, terminateEnd); + const terminatingFinalizerStart = source.indexOf( + 'const finalizeTerminatingWorker' + ); + const terminatingFinalizerEnd = source.indexOf( + 'WorkerClass.prototype.postMessage', + terminatingFinalizerStart + ); + const terminatingFinalizer = source.slice( + terminatingFinalizerStart, + terminatingFinalizerEnd + ); + + const measuredTerminate = terminateBlock.indexOf( + 'const termination = originalTerminate.call(this)' + ); + const measuredFinalization = terminateBlock.indexOf( + 'void finalizeTerminatingWorker(record, false)' + ); + const diagnosticFinalizationStart = terminateBlock.indexOf( + 'void finalizeTerminatingWorker(record, true)' + ); + const diagnosticFinalizationWait = terminateBlock.indexOf( + 'await waitForTerminatingWorkerFinalization(record)' + ); + const diagnosticTerminate = terminateBlock.indexOf( + 'return originalTerminate.call(this)', + diagnosticFinalizationWait + ); + + assert.ok( + diagnosticFinalizationStart >= 0 && + diagnosticFinalizationStart < diagnosticFinalizationWait && + diagnosticFinalizationWait < diagnosticTerminate, + 'diagnostic capture must start profile finalization, bound its wait, and then terminate the worker' + ); + assert.ok( + measuredTerminate >= 0 && measuredTerminate < measuredFinalization, + 'measured capture must initiate termination before best-effort finalization' + ); + assert.doesNotMatch(terminatingFinalizer, /takeWorkerHeapSnapshot/); + assert.match( + terminatingFinalizer, + /stopWorkerProfile\(record, waitForProfileHandle\)/ + ); + assert.match(source, /resolvedProfileHandle/); + assert.match(source, /worker-profile-finalization-timeout/); + assert.match(source, /finalizationTimedOut/); + assert.match(source, /profileCaptureKey/); + assert.match(source, /record\.profileCaptureKey !== profileCaptureKey/); +}); + +test('worker capture failures use an artifact-neutral timeline label', () => { + const source = readFileSync( + new URL('./m3u-refresh-main-capture.ts', import.meta.url), + 'utf8' + ); + + assert.match(source, /worker-artifact-error:\$\{stage\}/); + assert.doesNotMatch(source, /worker-profile-error:\$\{stage\}/); +}); + +test('capture generation resets artifact paths and always emits current DB records', () => { + const source = readFileSync( + new URL('./m3u-refresh-main-capture.ts', import.meta.url), + 'utf8' + ); + const start = source.indexOf('const startCapture = async ('); + const diagnostic = source.indexOf('if (state.diagnostic)', start); + const profileReset = source.indexOf('state.mainProfilePath = null', start); + const snapshotReset = source.indexOf( + 'state.mainSnapshotPath = null', + start + ); + + assert.ok(profileReset > start && profileReset < diagnostic); + assert.ok(snapshotReset > start && snapshotReset < diagnostic); + assert.match( + source, + /record\.captureGeneration === state\.captureGeneration[\s\S]*record\.kind === 'database\.worker'/ + ); +}); diff --git a/apps/electron-backend-e2e/src/performance/database-worker-post-gc-finalization.ts b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-finalization.ts new file mode 100644 index 000000000..53ff8a226 --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-finalization.ts @@ -0,0 +1,90 @@ +import type { DatabaseWorkerPostGcProbeUnavailableReason } from './database-worker-post-gc-probe'; +import type { DatabaseWorkerPostGcUnavailableReason } from './database-worker-post-gc-selection'; + +export type DatabaseWorkerPostGcAncillaryFailureStage = + 'heap-snapshot' | 'post-gc-probe' | 'profile-stop' | 'profile-write'; + +export type DatabaseWorkerPostGcFinalizationOutcome = + | { + readonly postGcHeapUsedBytes: number; + readonly unavailableReason: null; + } + | { + readonly postGcHeapUsedBytes: null; + readonly unavailableReason: + | DatabaseWorkerPostGcProbeUnavailableReason + | DatabaseWorkerPostGcUnavailableReason; + }; + +export interface DatabaseWorkerPostGcFinalizationInput { + readonly finalizationKey: object; + readonly joinFinalSample: () => Promise; + readonly probePostGc: () => Promise; + readonly reportAncillaryFailure: ( + stage: DatabaseWorkerPostGcAncillaryFailureStage, + error: unknown + ) => void; + readonly stopProfile: () => Promise; + readonly stopSampling: () => void; + readonly takeHeapSnapshot?: () => Promise; +} + +export interface DatabaseWorkerPostGcFinalizationApi { + finalize( + input: DatabaseWorkerPostGcFinalizationInput + ): Promise; +} + +export function createDatabaseWorkerPostGcFinalizationApi(): DatabaseWorkerPostGcFinalizationApi { + const finalizations = new WeakMap< + object, + Promise + >(); + + const helpers = { + async run( + input: DatabaseWorkerPostGcFinalizationInput + ): Promise { + input.stopSampling(); + await input.joinFinalSample(); + try { + await input.stopProfile(); + } catch (error: unknown) { + input.reportAncillaryFailure('profile-stop', error); + } + + let outcome: DatabaseWorkerPostGcFinalizationOutcome; + try { + outcome = await input.probePostGc(); + } catch (error: unknown) { + input.reportAncillaryFailure('post-gc-probe', error); + outcome = { + postGcHeapUsedBytes: null, + unavailableReason: 'capture-failed', + }; + } + if (input.takeHeapSnapshot) { + try { + await input.takeHeapSnapshot(); + } catch (error: unknown) { + input.reportAncillaryFailure('heap-snapshot', error); + } + } + return outcome; + }, + + finalize( + input: DatabaseWorkerPostGcFinalizationInput + ): Promise { + const existing = finalizations.get(input.finalizationKey); + if (existing) { + return existing; + } + const finalization = helpers.run(input); + finalizations.set(input.finalizationKey, finalization); + return finalization; + }, + }; + + return Object.freeze({ finalize: helpers.finalize }); +} diff --git a/apps/electron-backend-e2e/src/performance/database-worker-post-gc-probe.spec.ts b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-probe.spec.ts index 59efb9084..7ddf4e4cf 100644 --- a/apps/electron-backend-e2e/src/performance/database-worker-post-gc-probe.spec.ts +++ b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-probe.spec.ts @@ -63,11 +63,25 @@ interface ProbeModule { createDatabaseWorkerPostGcProbeApi?: () => ProbeApi; } +interface ProductionWorkerPostGcModule { + DATABASE_WORKER_POST_GC_HEAP_UNAVAILABLE_REASON?: Readonly< + Record + >; +} + const probeModulePromise = import( new URL('./database-worker-post-gc-probe.ts', import.meta.url).href ) .then((module) => module as ProbeModule) .catch(() => null); +const productionWorkerPostGcModulePromise = import( + new URL( + '../../../electron-backend/src/app/workers/database-worker-post-gc-heap.ts', + import.meta.url + ).href +) + .then((module) => module as ProductionWorkerPostGcModule) + .catch(() => null); class FakePort extends EventEmitter implements ProbePort { closeCalls = 0; @@ -231,12 +245,21 @@ test('serializes the factory and sends the exact one-shot worker request with it test('preserves every coherent worker-side unavailable result', async () => { const api = await restoreSerializableApi(); - const reasons: readonly WorkerUnavailableReason[] = [ - 'capture-failed', - 'gc-unavailable', - 'profiling-disabled', - 'worker-busy', - ]; + const productionModule = await productionWorkerPostGcModulePromise; + assert.ok(productionModule, 'production worker post-GC module must load'); + const productionReasons = + productionModule.DATABASE_WORKER_POST_GC_HEAP_UNAVAILABLE_REASON; + assert.ok(productionReasons, 'production worker reasons must be exported'); + const reasons = Object.values(productionReasons); + assert.deepEqual( + new Set(reasons), + new Set([ + 'capture-failed', + 'gc-unavailable', + 'profiling-disabled', + 'worker-busy', + ]) + ); for (const unavailableReason of reasons) { const harness = createProbeHarness(); diff --git a/apps/electron-backend-e2e/src/performance/database-worker-post-gc-report.spec.ts b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-report.spec.ts index bdcb900ac..83b221927 100644 --- a/apps/electron-backend-e2e/src/performance/database-worker-post-gc-report.spec.ts +++ b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-report.spec.ts @@ -96,6 +96,30 @@ test('summary includes every database isolate instead of silently taking the fir p95: 290, p99: 298, }); + assert.deepEqual( + ( + summary as unknown as { + readonly validity: { + readonly databaseWorkerPostGc: unknown; + }; + } + ).validity.databaseWorkerPostGc, + { + applicableMeasuredRunCount: 1, + invalidMeasuredRuns: [ + { + databaseWorkerCount: 3, + reason: 'database-worker-unexpected-activity', + runId: 'measured', + }, + ], + measuredRunCount: 1, + notApplicableMeasuredRuns: [], + validForBenchmark: false, + validForComparison: false, + validMeasuredRunCount: 0, + } + ); assert.equal( ( summary.iterations[0]?.main.workers[2] as WorkerCaptureMetrics & { diff --git a/apps/electron-backend-e2e/src/performance/database-worker-post-gc-validity.spec.ts b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-validity.spec.ts new file mode 100644 index 000000000..c26a4163d --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-validity.spec.ts @@ -0,0 +1,392 @@ +/* eslint-disable playwright/expect-expect -- This is a Node assertion-based performance contract test. */ +import assert from 'node:assert/strict'; +import { readFileSync } from 'node:fs'; +import test from 'node:test'; + +type IterationKind = 'diagnostic' | 'measured' | 'warmup'; +type WorkerKind = 'database.worker' | 'playlist-refresh.worker'; + +interface TestWorkerCapture { + readonly kind: WorkerKind; + readonly postGcHeapUnavailableReason: string | null; + readonly postGcHeapUsedBytes: number | null; +} + +interface TestIteration { + readonly cancellationEffectObserved: boolean; + readonly kind: IterationKind; + readonly main: { + readonly timeline: readonly { + readonly type: string; + }[]; + readonly workers: readonly TestWorkerCapture[]; + }; + readonly runId: string; +} + +interface InvalidMeasuredRun { + readonly databaseWorkerCount: number; + readonly reason: string; + readonly runId: string; +} + +interface DatabaseWorkerPostGcValidity { + readonly applicableMeasuredRunCount: number; + readonly invalidMeasuredRuns: readonly InvalidMeasuredRun[]; + readonly measuredRunCount: number; + readonly notApplicableMeasuredRuns: readonly { + readonly reason: string; + readonly runId: string; + }[]; + readonly validForBenchmark: boolean; + readonly validForComparison: boolean; + readonly validMeasuredRunCount: number; +} + +interface ValidityModule { + assessDatabaseWorkerPostGcValidity?: ( + iterations: readonly TestIteration[] + ) => DatabaseWorkerPostGcValidity; +} + +const validityModulePromise = import( + new URL('./database-worker-post-gc-validity.ts', import.meta.url).href +) + .then((module) => module as ValidityModule) + .catch(() => null); + +function worker( + postGcHeapUsedBytes: number | null, + postGcHeapUnavailableReason: string | null, + kind: WorkerKind = 'database.worker' +): TestWorkerCapture { + return { + kind, + postGcHeapUnavailableReason, + postGcHeapUsedBytes, + }; +} + +function iteration( + runId: string, + kind: IterationKind, + workers: readonly TestWorkerCapture[], + timeline: readonly { readonly type: string }[] = [], + cancellationEffectObserved = false +): TestIteration { + return { + cancellationEffectObserved, + kind, + main: { timeline, workers }, + runId, + }; +} + +test('validates database post-GC heap per measured iteration without losing raw reasons', async () => { + const module = await validityModulePromise; + assert.ok(module, 'database worker post-GC validity helper must exist'); + const assess = module.assessDatabaseWorkerPostGcValidity; + assert.equal(typeof assess, 'function'); + + const valid = assess([ + iteration('warmup-broken', 'warmup', []), + iteration('diagnostic-broken', 'diagnostic', [ + worker(null, 'gc-unavailable'), + worker(100, null), + ]), + iteration('measured-valid', 'measured', [ + worker(4_096, null), + worker( + null, + 'worker-force-terminated-before-gc', + 'playlist-refresh.worker' + ), + ]), + ]); + assert.deepEqual(valid, { + applicableMeasuredRunCount: 1, + invalidMeasuredRuns: [], + measuredRunCount: 1, + notApplicableMeasuredRuns: [], + validForBenchmark: true, + validForComparison: true, + validMeasuredRunCount: 1, + }); + + const compensatingCardinality = assess([ + iteration('measured-duplicate', 'measured', [ + worker(1_024, null), + worker(2_048, null), + ]), + iteration( + 'measured-missing', + 'measured', + [worker(3_072, null, 'playlist-refresh.worker')], + [{ type: 'db-request' }] + ), + ]); + assert.deepEqual(compensatingCardinality, { + applicableMeasuredRunCount: 2, + invalidMeasuredRuns: [ + { + databaseWorkerCount: 2, + reason: 'multiple-database-workers', + runId: 'measured-duplicate', + }, + { + databaseWorkerCount: 0, + reason: 'database-worker-missing', + runId: 'measured-missing', + }, + ], + measuredRunCount: 2, + notApplicableMeasuredRuns: [], + validForBenchmark: false, + validForComparison: false, + validMeasuredRunCount: 0, + }); + + const unavailableWorker = worker(null, 'post-gc-probe-timeout'); + const unavailableAndIncoherent = assess([ + iteration('measured-unavailable', 'measured', [unavailableWorker]), + iteration('measured-incoherent', 'measured', [ + worker(8_192, 'gc-unavailable'), + ]), + ]); + assert.deepEqual(unavailableAndIncoherent, { + applicableMeasuredRunCount: 2, + invalidMeasuredRuns: [ + { + databaseWorkerCount: 1, + reason: 'post-gc-probe-timeout', + runId: 'measured-unavailable', + }, + { + databaseWorkerCount: 1, + reason: 'post-gc-capture-invalid', + runId: 'measured-incoherent', + }, + ], + measuredRunCount: 2, + notApplicableMeasuredRuns: [], + validForBenchmark: false, + validForComparison: false, + validMeasuredRunCount: 0, + }); + assert.equal( + unavailableWorker.postGcHeapUnavailableReason, + 'post-gc-probe-timeout' + ); +}); + +test('marks parsing cancellation before persistence as DB N/A without accepting late or missing captured DB work', async () => { + const module = await validityModulePromise; + assert.ok(module); + const assess = module.assessDatabaseWorkerPostGcValidity; + assert.equal(typeof assess, 'function'); + + const noDatabasePhase = assess([ + iteration( + 'run-01', + 'measured', + [ + worker( + null, + 'worker-force-terminated-before-gc', + 'playlist-refresh.worker' + ), + ], + [], + true + ), + iteration( + 'run-02', + 'measured', + [ + worker( + null, + 'worker-force-terminated-before-gc', + 'playlist-refresh.worker' + ), + ], + [], + true + ), + ]); + assert.deepEqual(noDatabasePhase, { + applicableMeasuredRunCount: 0, + invalidMeasuredRuns: [], + measuredRunCount: 2, + notApplicableMeasuredRuns: [ + { + reason: 'operation-cancelled-before-database-phase', + runId: 'run-01', + }, + { + reason: 'operation-cancelled-before-database-phase', + runId: 'run-02', + }, + ], + validForBenchmark: true, + validForComparison: false, + validMeasuredRunCount: 0, + }); + + const mixedApplicability = assess([ + iteration('run-cancelled', 'measured', [], [], true), + iteration('run-persisted', 'measured', [worker(4_096, null)]), + ]); + assert.deepEqual(mixedApplicability, { + applicableMeasuredRunCount: 1, + invalidMeasuredRuns: [], + measuredRunCount: 2, + notApplicableMeasuredRuns: [ + { + reason: 'operation-cancelled-before-database-phase', + runId: 'run-cancelled', + }, + ], + validForBenchmark: true, + validForComparison: false, + validMeasuredRunCount: 1, + }); + + const lateDatabaseWork = assess([ + iteration( + 'run-late', + 'measured', + [], + [{ type: 'db-request-after-capture-cutoff' }], + true + ), + ]); + assert.deepEqual(lateDatabaseWork, { + applicableMeasuredRunCount: 1, + invalidMeasuredRuns: [ + { + databaseWorkerCount: 0, + reason: 'database-worker-activity-after-cutoff', + runId: 'run-late', + }, + ], + measuredRunCount: 1, + notApplicableMeasuredRuns: [], + validForBenchmark: false, + validForComparison: false, + validMeasuredRunCount: 0, + }); + + const unrelatedDatabaseWorker = assess([ + iteration( + 'run-unrelated-db', + 'measured', + [worker(4_096, null)], + [{ type: 'db-request' }], + true + ), + ]); + assert.deepEqual(unrelatedDatabaseWorker, { + applicableMeasuredRunCount: 1, + invalidMeasuredRuns: [ + { + databaseWorkerCount: 1, + reason: 'database-worker-unexpected-activity', + runId: 'run-unrelated-db', + }, + ], + measuredRunCount: 1, + notApplicableMeasuredRuns: [], + validForBenchmark: false, + validForComparison: false, + validMeasuredRunCount: 0, + }); + + const nonCancelLateRequest = assess([ + iteration( + 'run-non-cancel-late', + 'measured', + [worker(8_192, null)], + [{ type: 'db-request-after-capture-cutoff' }] + ), + ]); + assert.deepEqual(nonCancelLateRequest, { + applicableMeasuredRunCount: 1, + invalidMeasuredRuns: [ + { + databaseWorkerCount: 1, + reason: 'database-worker-activity-after-cutoff', + runId: 'run-non-cancel-late', + }, + ], + measuredRunCount: 1, + notApplicableMeasuredRuns: [], + validForBenchmark: false, + validForComparison: false, + validMeasuredRunCount: 0, + }); +}); + +test('the benchmark preserves raw artifacts before rejecting an invalid formal comparison', () => { + const source = readFileSync( + new URL('./m3u-refresh-cancellation.benchmark.ts', import.meta.url), + 'utf8' + ); + const summaryWrite = source.indexOf( + "writeJson(join(config.outputDirectory, 'summary.json'), summary)" + ); + const validityGuard = source.indexOf( + 'summary.validity.databaseWorkerPostGc.validForBenchmark' + ); + const cancellationEffectGuard = source.indexOf( + 'summary.cancellationEffectRate !== 1' + ); + + assert.ok(summaryWrite >= 0, 'summary artifact write must exist'); + assert.ok( + cancellationEffectGuard >= 0, + 'formal cancellation-effect guard must exist' + ); + assert.ok(validityGuard >= 0, 'formal validity guard must exist'); + assert.ok( + summaryWrite < cancellationEffectGuard && summaryWrite < validityGuard, + 'raw summary must be durable before either formal guard fails' + ); + assert.match(source, /Cancellation effect was not observed/); + assert.match(source, /Database worker post-GC capture is invalid/); +}); + +test('the benchmark persists the real seed DB lifecycle and atomically arms the measured generation', () => { + const source = readFileSync( + new URL('./m3u-refresh-cancellation.benchmark.ts', import.meta.url), + 'utf8' + ); + const seedCaptureStart = source.indexOf( + 'await startMainCapture(app.electronApp', + source.indexOf('const rendererWindowIdentity') + ); + const seedPlaylist = source.indexOf('await seedPlaylist'); + const seedSettlement = source.indexOf('await waitForSeedSettlement'); + const seedRollover = source.indexOf( + 'await rolloverMainCapture(app.electronApp', + seedSettlement + ); + const seedArtifactWrite = source.indexOf( + "'seed-main-capture.json'", + seedRollover + ); + + assert.ok(seedCaptureStart >= 0 && seedCaptureStart < seedPlaylist); + assert.ok(seedPlaylist < seedSettlement); + assert.ok(seedSettlement < seedRollover); + assert.ok(seedRollover < seedArtifactWrite); + assert.doesNotMatch( + source.slice(seedRollover + 1, source.indexOf('const renderer =')), + /await startMainCapture\(app\.electronApp/ + ); + assert.match(source, /seedRollover\.nextCaptureStarted/); + assert.match(source, /seedRollover\.nextCaptureUnavailableReason/); + assert.match( + source, + /status\.databaseRequests\s*>\s*0[\s\S]*status\.databaseUpsertsCompleted\s*>\s*0[\s\S]*status\.databasePending\s*===\s*0/ + ); +}); diff --git a/apps/electron-backend-e2e/src/performance/database-worker-post-gc-validity.ts b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-validity.ts new file mode 100644 index 000000000..43345e718 --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/database-worker-post-gc-validity.ts @@ -0,0 +1,141 @@ +import { + type DatabaseWorkerPostGcValidity, + type InvalidDatabaseWorkerPostGcMeasuredRun, + WORKER_POST_GC_HEAP_UNAVAILABLE_REASON, +} from './m3u-refresh-cancellation-contract'; + +export interface DatabaseWorkerPostGcValidityWorker { + readonly kind: string; + readonly postGcHeapUnavailableReason: unknown; + readonly postGcHeapUsedBytes: unknown; +} + +export interface DatabaseWorkerPostGcValidityIteration { + readonly cancellationEffectObserved?: boolean; + readonly kind: string; + readonly main: { + readonly timeline?: readonly { + readonly type: string; + }[]; + readonly workers: readonly DatabaseWorkerPostGcValidityWorker[]; + }; + readonly runId: string; +} + +const UNAVAILABLE_REASONS = new Set( + Object.values(WORKER_POST_GC_HEAP_UNAVAILABLE_REASON) +); + +export function assessDatabaseWorkerPostGcValidity( + iterations: readonly DatabaseWorkerPostGcValidityIteration[] +): DatabaseWorkerPostGcValidity { + const measured = iterations.filter( + (iteration) => iteration.kind === 'measured' + ); + const invalidMeasuredRuns: InvalidDatabaseWorkerPostGcMeasuredRun[] = []; + const notApplicableMeasuredRuns: { + readonly reason: 'operation-cancelled-before-database-phase'; + readonly runId: string; + }[] = []; + let applicableMeasuredRunCount = 0; + let validMeasuredRunCount = 0; + + for (const iteration of measured) { + const databaseWorkers = iteration.main.workers.filter( + (worker) => worker.kind === 'database.worker' + ); + const timeline = iteration.main.timeline ?? []; + const lateDatabaseRequest = timeline.some( + (record) => record.type === 'db-request-after-capture-cutoff' + ); + const databasePhaseObserved = + lateDatabaseRequest || + timeline.some((record) => record.type === 'db-request'); + + if (lateDatabaseRequest) { + applicableMeasuredRunCount += 1; + invalidMeasuredRuns.push({ + databaseWorkerCount: databaseWorkers.length, + reason: WORKER_POST_GC_HEAP_UNAVAILABLE_REASON.DATABASE_WORKER_ACTIVITY_AFTER_CUTOFF, + runId: iteration.runId, + }); + continue; + } + + if (iteration.cancellationEffectObserved === true) { + if (!databasePhaseObserved && databaseWorkers.length === 0) { + notApplicableMeasuredRuns.push({ + reason: 'operation-cancelled-before-database-phase', + runId: iteration.runId, + }); + continue; + } + + applicableMeasuredRunCount += 1; + invalidMeasuredRuns.push({ + databaseWorkerCount: databaseWorkers.length, + reason: WORKER_POST_GC_HEAP_UNAVAILABLE_REASON.DATABASE_WORKER_UNEXPECTED_ACTIVITY, + runId: iteration.runId, + }); + continue; + } + + if (databaseWorkers.length === 0) { + applicableMeasuredRunCount += 1; + invalidMeasuredRuns.push({ + databaseWorkerCount: 0, + reason: WORKER_POST_GC_HEAP_UNAVAILABLE_REASON.DATABASE_WORKER_MISSING, + runId: iteration.runId, + }); + continue; + } + + applicableMeasuredRunCount += 1; + if (databaseWorkers.length !== 1) { + invalidMeasuredRuns.push({ + databaseWorkerCount: databaseWorkers.length, + reason: WORKER_POST_GC_HEAP_UNAVAILABLE_REASON.MULTIPLE_DATABASE_WORKERS, + runId: iteration.runId, + }); + continue; + } + + const worker = databaseWorkers[0] as DatabaseWorkerPostGcValidityWorker; + const heap = worker.postGcHeapUsedBytes; + const reason = worker.postGcHeapUnavailableReason; + if ( + Number.isSafeInteger(heap) && + Number(heap) >= 0 && + reason === null + ) { + validMeasuredRunCount += 1; + continue; + } + + invalidMeasuredRuns.push({ + databaseWorkerCount: 1, + reason: + heap === null && + typeof reason === 'string' && + UNAVAILABLE_REASONS.has(reason) + ? reason + : WORKER_POST_GC_HEAP_UNAVAILABLE_REASON.INVALID_CAPTURE, + runId: iteration.runId, + }); + } + + return Object.freeze({ + applicableMeasuredRunCount, + invalidMeasuredRuns: Object.freeze(invalidMeasuredRuns), + measuredRunCount: measured.length, + notApplicableMeasuredRuns: Object.freeze(notApplicableMeasuredRuns), + validForBenchmark: + measured.length > 0 && invalidMeasuredRuns.length === 0, + validForComparison: + measured.length > 0 && + applicableMeasuredRunCount === measured.length && + validMeasuredRunCount === measured.length && + invalidMeasuredRuns.length === 0, + validMeasuredRunCount, + }); +} diff --git a/apps/electron-backend-e2e/src/performance/diagnostic-playlist-worker-capture.spec.ts b/apps/electron-backend-e2e/src/performance/diagnostic-playlist-worker-capture.spec.ts new file mode 100644 index 000000000..f8fdc7434 --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/diagnostic-playlist-worker-capture.spec.ts @@ -0,0 +1,193 @@ +/* eslint-disable playwright/expect-expect -- This is a Node assertion-based performance contract test. */ +import assert from 'node:assert/strict'; +import { readFileSync } from 'node:fs'; +import { mkdtemp, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import test from 'node:test'; + +import type { MainCaptureMetrics } from './m3u-refresh-cancellation-contract'; + +interface DiagnosticCaptureModule { + validateDiagnosticPlaylistWorkerCapture?: ( + main: MainCaptureMetrics, + iterationDirectory: string + ) => Promise; +} + +const diagnosticCaptureModulePromise = import( + new URL('./diagnostic-playlist-worker-capture.ts', import.meta.url).href +) + .then((module) => module as DiagnosticCaptureModule) + .catch(() => null); + +function mainCapture(input: { + readonly profilePath: string | null; + readonly snapshotPath?: string | null; + readonly timelineType?: string; +}): MainCaptureMetrics { + return { + timeline: input.timelineType + ? [{ epochMs: 1, type: input.timelineType }] + : [], + workers: [ + { + kind: 'playlist-refresh.worker', + ordinal: 2, + postGcHeapUnavailableReason: + 'worker-force-terminated-before-gc', + postGcHeapUsedBytes: null, + profilePath: input.profilePath, + snapshotPath: input.snapshotPath ?? null, + terminatedEpochMs: 2, + }, + ], + } as unknown as MainCaptureMetrics; +} + +test('accepts one parseable diagnostic playlist-worker profile inside the iteration directory', async () => { + const module = await diagnosticCaptureModulePromise; + assert.ok( + module, + 'diagnostic playlist worker capture validator must exist' + ); + const validate = module.validateDiagnosticPlaylistWorkerCapture; + assert.equal(typeof validate, 'function'); + const directory = await mkdtemp( + join(tmpdir(), 'iptvnator-diagnostic-worker-') + ); + + try { + const profilePath = join( + directory, + 'playlist-refresh.worker-2.cpuprofile' + ); + await writeFile( + profilePath, + JSON.stringify({ nodes: [{ id: 1 }], samples: [1] }) + ); + await assert.doesNotReject(() => + validate?.(mainCapture({ profilePath }), directory) + ); + } finally { + await rm(directory, { force: true, recursive: true }); + } +}); + +test('fails closed for missing, escaped, malformed, or contaminated diagnostic artifacts', async () => { + const module = await diagnosticCaptureModulePromise; + assert.ok(module); + const validate = module.validateDiagnosticPlaylistWorkerCapture; + assert.equal(typeof validate, 'function'); + const directory = await mkdtemp( + join(tmpdir(), 'iptvnator-diagnostic-worker-invalid-') + ); + + try { + await assert.rejects( + () => validate?.(mainCapture({ profilePath: null }), directory), + /diagnostic-playlist-worker-profile-missing/ + ); + await assert.rejects( + () => + validate?.( + mainCapture({ + profilePath: join( + directory, + '..', + 'escaped.cpuprofile' + ), + }), + directory + ), + /diagnostic-playlist-worker-profile-path-invalid/ + ); + const malformedPath = join( + directory, + 'playlist-refresh.worker-2.cpuprofile' + ); + await writeFile(malformedPath, '{"nodes":[]}'); + await assert.rejects( + () => + validate?.( + mainCapture({ profilePath: malformedPath }), + directory + ), + /diagnostic-playlist-worker-profile-invalid/ + ); + await assert.rejects( + () => + validate?.( + mainCapture({ + profilePath: malformedPath, + timelineType: + 'worker-artifact-error:profile-stop:failed', + }), + directory + ), + /diagnostic-playlist-worker-artifact-error/ + ); + await assert.rejects( + () => + validate?.( + mainCapture({ + profilePath: malformedPath, + snapshotPath: join( + directory, + 'playlist-refresh.worker-2.heapsnapshot' + ), + }), + directory + ), + /diagnostic-playlist-worker-snapshot-unexpected/ + ); + } finally { + await rm(directory, { force: true, recursive: true }); + } +}); + +test('the benchmark persists diagnostic raw metrics before requiring a parseable playlist-worker profile', () => { + const source = readFileSync( + new URL('./m3u-refresh-cancellation.benchmark.ts', import.meta.url), + 'utf8' + ); + const resultWrite = source.indexOf( + "writeJson(join(iterationDirectory, 'result.json'), result)" + ); + const summaryWrite = source.indexOf( + "writeJson(join(config.outputDirectory, 'summary.json'), summary)" + ); + const diagnosticValidation = source.indexOf( + 'await validateDiagnosticPlaylistWorkerCapture(' + ); + + assert.ok(resultWrite >= 0); + assert.ok(summaryWrite >= 0); + assert.ok(diagnosticValidation > summaryWrite); +}); + +test('the benchmark requires the diagnostic worker profile to come from an observed cancellation', () => { + const source = readFileSync( + new URL('./m3u-refresh-cancellation.benchmark.ts', import.meta.url), + 'utf8' + ); + const summaryWrite = source.indexOf( + "writeJson(join(config.outputDirectory, 'summary.json'), summary)" + ); + const diagnosticCancellationGuard = source.indexOf( + 'if (!diagnosticIteration.cancellationEffectObserved)' + ); + const diagnosticValidation = source.indexOf( + 'await validateDiagnosticPlaylistWorkerCapture(' + ); + + assert.ok(summaryWrite >= 0); + assert.ok( + diagnosticCancellationGuard > summaryWrite, + 'raw manifest and summary must remain durable when the diagnostic cancellation is invalid' + ); + assert.ok( + diagnosticCancellationGuard < diagnosticValidation, + 'the cancellation guard must reject a normal-completion worker profile before artifact validation' + ); +}); diff --git a/apps/electron-backend-e2e/src/performance/diagnostic-playlist-worker-capture.ts b/apps/electron-backend-e2e/src/performance/diagnostic-playlist-worker-capture.ts new file mode 100644 index 000000000..855931ed4 --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/diagnostic-playlist-worker-capture.ts @@ -0,0 +1,72 @@ +import { readFile } from 'node:fs/promises'; +import { join, resolve } from 'node:path'; + +import type { MainCaptureMetrics } from './m3u-refresh-cancellation-contract'; + +interface CpuProfile { + readonly nodes?: unknown; + readonly samples?: unknown; +} + +export async function validateDiagnosticPlaylistWorkerCapture( + main: MainCaptureMetrics, + iterationDirectory: string +): Promise { + if ( + main.timeline.some((record) => + record.type.startsWith('worker-artifact-error:') + ) + ) { + throw new Error('diagnostic-playlist-worker-artifact-error'); + } + + const workers = main.workers.filter( + (worker) => worker.kind === 'playlist-refresh.worker' + ); + if (workers.length !== 1) { + throw new Error('diagnostic-playlist-worker-cardinality-invalid'); + } + const worker = workers[0]; + if (!worker) { + throw new Error('diagnostic-playlist-worker-cardinality-invalid'); + } + if (worker.snapshotPath !== null) { + throw new Error('diagnostic-playlist-worker-snapshot-unexpected'); + } + if ( + worker.postGcHeapUsedBytes !== null || + worker.postGcHeapUnavailableReason !== + 'worker-force-terminated-before-gc' || + worker.terminatedEpochMs === null + ) { + throw new Error('diagnostic-playlist-worker-termination-invalid'); + } + if (worker.profilePath === null) { + throw new Error('diagnostic-playlist-worker-profile-missing'); + } + + const expectedProfilePath = join( + resolve(iterationDirectory), + `playlist-refresh.worker-${worker.ordinal}.cpuprofile` + ); + if (resolve(worker.profilePath) !== expectedProfilePath) { + throw new Error('diagnostic-playlist-worker-profile-path-invalid'); + } + + let profile: CpuProfile; + try { + profile = JSON.parse( + await readFile(expectedProfilePath, 'utf8') + ) as CpuProfile; + } catch { + throw new Error('diagnostic-playlist-worker-profile-invalid'); + } + if ( + !Array.isArray(profile.nodes) || + profile.nodes.length === 0 || + !Array.isArray(profile.samples) || + profile.samples.length === 0 + ) { + throw new Error('diagnostic-playlist-worker-profile-invalid'); + } +} 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 daf609acb..ca5b02a0b 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 @@ -17,6 +17,40 @@ export const PERFORMANCE_WORKER_KIND = { export type PerformanceWorkerKind = (typeof PERFORMANCE_WORKER_KIND)[keyof typeof PERFORMANCE_WORKER_KIND]; +export const WORKER_POST_GC_HEAP_UNAVAILABLE_REASON = { + CAPTURE_FAILED: 'capture-failed', + DATABASE_WORKER_ACTIVITY_AFTER_CUTOFF: + 'database-worker-activity-after-cutoff', + DATABASE_WORKER_MISSING: 'database-worker-missing', + DATABASE_WORKER_NOT_IDLE: 'database-worker-not-idle', + DATABASE_WORKER_UNEXPECTED_ACTIVITY: 'database-worker-unexpected-activity', + GC_UNAVAILABLE: 'gc-unavailable', + INVALID_CAPTURE: 'post-gc-capture-invalid', + MULTIPLE_DATABASE_WORKERS: 'multiple-database-workers', + PROBE_INVALID_RESPONSE: 'post-gc-probe-invalid-response', + PROBE_MESSAGE_ERROR: 'post-gc-probe-message-error', + PROBE_NOT_RUN: 'post-gc-probe-not-run', + PROBE_PORT_CLOSED: 'post-gc-probe-port-closed', + PROBE_POST_FAILED: 'post-gc-probe-post-failed', + PROBE_TIMEOUT: 'post-gc-probe-timeout', + PROFILING_DISABLED: 'profiling-disabled', + WORKER_BUSY: 'worker-busy', + WORKER_FORCE_TERMINATED: 'worker-force-terminated-before-gc', +} as const; + +export type WorkerPostGcHeapUnavailableReason = + (typeof WORKER_POST_GC_HEAP_UNAVAILABLE_REASON)[keyof typeof WORKER_POST_GC_HEAP_UNAVAILABLE_REASON]; + +export type WorkerPostGcHeapCapture = + | { + readonly postGcHeapUnavailableReason: null; + readonly postGcHeapUsedBytes: number; + } + | { + readonly postGcHeapUnavailableReason: WorkerPostGcHeapUnavailableReason; + readonly postGcHeapUsedBytes: null; + }; + export interface NumericDistribution { readonly count: number; readonly max: number | null; @@ -73,7 +107,7 @@ export interface WorkerRequestPerformanceMetrics { readonly workStartedEpochMs: number | null; } -export interface WorkerCaptureMetrics { +export type WorkerCaptureMetrics = { readonly cancelPostedEpochMs: number | null; readonly cpuSystemMicros: number | null; readonly cpuUserMicros: number | null; @@ -82,16 +116,16 @@ export interface WorkerCaptureMetrics { readonly eventLoopUtilization: number | null; readonly kind: PerformanceWorkerKind; readonly operationId: string | null; + readonly ordinal: number; readonly peakExternalBytes: number; readonly peakHeapUsedBytes: number; 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; -} +} & WorkerPostGcHeapCapture; export interface MainCaptureMetrics { readonly cpuProfilePath: string | null; @@ -199,6 +233,27 @@ export interface CancellationBenchmarkManifest { readonly warmupRuns: number; } +export interface InvalidDatabaseWorkerPostGcMeasuredRun { + readonly databaseWorkerCount: number; + readonly reason: string; + readonly runId: string; +} + +export interface NotApplicableDatabaseWorkerPostGcMeasuredRun { + readonly reason: 'operation-cancelled-before-database-phase'; + readonly runId: string; +} + +export interface DatabaseWorkerPostGcValidity { + readonly applicableMeasuredRunCount: number; + readonly invalidMeasuredRuns: readonly InvalidDatabaseWorkerPostGcMeasuredRun[]; + readonly measuredRunCount: number; + readonly notApplicableMeasuredRuns: readonly NotApplicableDatabaseWorkerPostGcMeasuredRun[]; + readonly validForBenchmark: boolean; + readonly validForComparison: boolean; + readonly validMeasuredRunCount: number; +} + export interface CancellationBenchmarkSummary { readonly cancellationEffectRate: number; readonly iterations: readonly CancellationIterationResult[]; @@ -252,4 +307,7 @@ export interface CancellationBenchmarkSummary { readonly unresponsiveEvents: number; readonly visibleTotalMs: NumericDistribution; }; + readonly validity: { + readonly databaseWorkerPostGc: DatabaseWorkerPostGcValidity; + }; } 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 5d5542024..a384d27de 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 @@ -19,6 +19,7 @@ import { nonNegativeDifference, summarizeNumbers, } from './performance-statistics'; +import { assessDatabaseWorkerPostGcValidity } from './database-worker-post-gc-validity'; export function createCancellationIterationResult(input: { readonly kind: CancellationIterationResult['kind']; @@ -306,6 +307,10 @@ export function createCancellationBenchmarkSummary( measured.map((iteration) => iteration.phases.visibleTotalMs) ), }), + validity: Object.freeze({ + databaseWorkerPostGc: + assessDatabaseWorkerPostGcValidity(iterations), + }), }); } diff --git a/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation.benchmark.ts b/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation.benchmark.ts index cf0858f9b..56b8a1a04 100644 --- a/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation.benchmark.ts +++ b/apps/electron-backend-e2e/src/performance/m3u-refresh-cancellation.benchmark.ts @@ -35,10 +35,12 @@ import { createCancellationBenchmarkSummary, createCancellationIterationResult, } from './m3u-refresh-cancellation-report'; +import { validateDiagnosticPlaylistWorkerCapture } from './diagnostic-playlist-worker-capture'; import { assertPerformanceArtifactCapacity } from './performance-artifact-preflight'; import { installMainCapture, readMainCaptureStatus, + rolloverMainCapture, startMainCapture, stopMainCapture, } from './m3u-refresh-main-capture'; @@ -135,6 +137,33 @@ export async function runM3uRefreshCancellationBenchmark(): Promise { const summary = createCancellationBenchmarkSummary(manifest, iterations); await writeJson(join(config.outputDirectory, 'manifest.json'), manifest); await writeJson(join(config.outputDirectory, 'summary.json'), summary); + const diagnosticIteration = iterations.find( + (iteration) => iteration.kind === PERFORMANCE_ITERATION_KIND.DIAGNOSTIC + ); + if (!diagnosticIteration) { + throw new Error('Diagnostic iteration is missing'); + } + if (!diagnosticIteration.cancellationEffectObserved) { + throw new Error( + 'Cancellation effect was not observed in the diagnostic run' + ); + } + await validateDiagnosticPlaylistWorkerCapture( + diagnosticIteration.main, + join(config.outputDirectory, diagnosticIteration.runId) + ); + if (summary.cancellationEffectRate !== 1) { + throw new Error( + `Cancellation effect was not observed in every measured run: ${summary.cancellationEffectRate}` + ); + } + if (!summary.validity.databaseWorkerPostGc.validForBenchmark) { + throw new Error( + `Database worker post-GC capture is invalid: ${JSON.stringify( + summary.validity.databaseWorkerPostGc.invalidMeasuredRuns + )}` + ); + } console.log( `[performance] completed: ${join( config.outputDirectory, @@ -179,6 +208,12 @@ async function runIteration( await assertRendererCdpTarget(app.mainWindow); app.mainWindow.setDefaultTimeout(120_000); await installMainCapture(app.electronApp); + const rendererWindowIdentity = await resolveRendererWindowIdentity(app); + await startMainCapture(app.electronApp, { + diagnostic: false, + outputDirectory: iterationDirectory, + rendererWindowIdentity, + }); await seedPlaylist(app.mainWindow, server.resourceUrl); await waitForSeedSettlement(app); server.serveLargeFixture(); @@ -195,12 +230,20 @@ async function runIteration( const diagnostic = definition.kind === PERFORMANCE_ITERATION_KIND.DIAGNOSTIC; - const rendererWindowIdentity = await resolveRendererWindowIdentity(app); - await startMainCapture(app.electronApp, { + const seedRollover = await rolloverMainCapture(app.electronApp, { diagnostic, outputDirectory: iterationDirectory, rendererWindowIdentity, }); + await writeJson( + join(iterationDirectory, 'seed-main-capture.json'), + seedRollover.completedCapture + ); + if (!seedRollover.nextCaptureStarted) { + throw new Error( + `Seed capture could not atomically arm the measured generation: ${seedRollover.nextCaptureUnavailableReason}` + ); + } const renderer = await startRendererCapture(app.mainWindow, { diagnostic, outputDirectory: iterationDirectory, @@ -287,11 +330,17 @@ async function seedPlaylist(page: Page, playlistUrl: string): Promise { async function waitForSeedSettlement(app: LaunchedElectronApp): Promise { await expect .poll( - async () => - (await readMainCaptureStatus(app.electronApp)).databasePending, + async () => { + const status = await readMainCaptureStatus(app.electronApp); + return ( + status.databaseRequests > 0 && + status.databaseUpsertsCompleted > 0 && + status.databasePending === 0 + ); + }, { timeout: 120_000 } ) - .toBe(0); + .toBe(true); await waitForTwoAnimationFrames(app.mainWindow); } 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 9d8a6323f..0931ac31a 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 @@ -7,11 +7,36 @@ import { type DatabaseRequestIdentity, type DatabaseRequestIdentityCaptureApi, } from './database-request-identity-capture'; -import type { MainCaptureMetrics } from './m3u-refresh-cancellation-contract'; +import { + createDatabaseWorkerPostGcCutoffApi, + type DatabaseWorkerPostGcCutoffApi, +} from './database-worker-post-gc-cutoff'; +import { + createDatabaseWorkerPostGcFinalizationApi, + type DatabaseWorkerPostGcAncillaryFailureStage, + type DatabaseWorkerPostGcFinalizationApi, +} from './database-worker-post-gc-finalization'; +import { + createDatabaseWorkerPostGcProbeApi, + type DatabaseWorkerPostGcProbeApi, +} from './database-worker-post-gc-probe'; +import { + createDatabaseWorkerPostGcSelectionApi, + type DatabaseWorkerPostGcSelectionApi, + type DatabaseWorkerPostGcUnavailableReason, +} from './database-worker-post-gc-selection'; +import type { + MainCaptureMetrics, + WorkerPostGcHeapUnavailableReason, +} from './m3u-refresh-cancellation-contract'; import { selectMainCaptureGeneration, type MainCaptureGenerationTransport, } from './worker-request-performance'; +import { + createWorkerTerminationGenerationApi, + type WorkerTerminationGenerationApi, +} from './worker-termination-generation'; import { createRendererProcessRssCaptureApi, type RendererProcessRssCaptureApi, @@ -42,6 +67,21 @@ export interface MainCaptureStatus { readonly playlistResponsesSucceeded: number; } +export interface MainCaptureRolloverResult { + readonly completedCapture: MainCaptureMetrics; + readonly nextCaptureStarted: boolean; + readonly nextCaptureUnavailableReason: string | null; +} + +interface MainCaptureRolloverStatus { + readonly nextCaptureStarted: boolean; + readonly nextCaptureUnavailableReason: string | null; +} + +type MainCaptureStopTransport = MainCaptureGenerationTransport & { + readonly rollover: MainCaptureRolloverStatus | null; +}; + export async function installMainCapture( electronApp: ElectronApplication ): Promise { @@ -49,11 +89,21 @@ export async function installMainCapture( const captureStateKeys = { databaseRequestIdentityStateKey: DATABASE_REQUEST_IDENTITY_CAPTURE_STATE_KEY, + databaseWorkerPostGcCutoffApiFactorySource: + createDatabaseWorkerPostGcCutoffApi.toString(), + databaseWorkerPostGcFinalizationApiFactorySource: + createDatabaseWorkerPostGcFinalizationApi.toString(), + databaseWorkerPostGcProbeApiFactorySource: + createDatabaseWorkerPostGcProbeApi.toString(), + databaseWorkerPostGcSelectionApiFactorySource: + createDatabaseWorkerPostGcSelectionApi.toString(), rendererProcessRssApiFactorySource: createRendererProcessRssCaptureApi.toString(), rendererWindowRssSessionApiFactorySource: createRendererWindowRssSessionApi.toString(), stateKey: MAIN_CAPTURE_STATE_KEY, + workerTerminationGenerationApiFactorySource: + createWorkerTerminationGenerationApi.toString(), }; await electronApp.evaluate(async ({ app, BrowserWindow }, input) => { type JsonRecord = Record; @@ -110,17 +160,25 @@ export async function installMainCapture( eluStart: WorkerElu | null; externalPeak: number; finalized: boolean; + finalizationKey: object; + finalizationTimedOut: boolean; finalizing: Promise | null; heapPeak: number; kind: 'database.worker' | 'playlist-refresh.worker'; operationId: string | null; + ordinal: number; + pendingCount: number; playlistId: string | null; + postGcHeapUnavailableReason: WorkerPostGcHeapUnavailableReason | null; postGcHeapUsed: number | null; profileHandle: Promise | null; + profileCaptureKey: object; profilePath: string | null; + profileResult: unknown | null; + resolvedProfileHandle: CpuProfileHandle | null; requestPerformance: WorkerRequestPerformance[]; responseEpochMs: number | null; - sampleBusy: boolean; + samplePromise: Promise | null; sampleTimer: NodeJS.Timeout | null; snapshotPath: string | null; terminatedEpochMs: number | null; @@ -149,6 +207,22 @@ export async function installMainCapture( )() as () => T; return factory(); }; + const databaseWorkerPostGcCutoffApi = + restoreFactory( + input.databaseWorkerPostGcCutoffApiFactorySource + ); + const databaseWorkerPostGcFinalizationApi = + restoreFactory( + input.databaseWorkerPostGcFinalizationApiFactorySource + ); + const databaseWorkerPostGcProbeApi = + restoreFactory( + input.databaseWorkerPostGcProbeApiFactorySource + ); + const databaseWorkerPostGcSelectionApi = + restoreFactory( + input.databaseWorkerPostGcSelectionApiFactorySource + ); const rendererProcessRssApi = restoreFactory( input.rendererProcessRssApiFactorySource @@ -157,6 +231,10 @@ export async function installMainCapture( restoreFactory( input.rendererWindowRssSessionApiFactorySource ); + const workerTerminationGenerationApi = + restoreFactory( + input.workerTerminationGenerationApiFactorySource + ); const runtimeProcess = process as typeof process & { getBuiltinModule(id: string): unknown; @@ -183,6 +261,7 @@ export async function installMainCapture( const originalPostMessage = WorkerClass.prototype.postMessage; const originalTerminate = WorkerClass.prototype.terminate; const records = new Map(); + let nextWorkerOrdinal = 1; const operationWorkers = new Map(); const dbRequests = new Map< string, @@ -215,6 +294,7 @@ export async function installMainCapture( postGcRss: null as number | null, rendererWindowSession: null as RendererWindowRssSession | null, sampleTimer: null as NodeJS.Timeout | null, + stopping: false, timeline: [] as TimelineRecord[], }; @@ -223,7 +303,7 @@ export async function installMainCapture( const recordTimeline = ( record: Omit ): void => { - if (state.active) { + if (state.active || state.stopping) { state.timeline.push({ epochMs: nowEpochMs(), ...record }); } }; @@ -270,22 +350,31 @@ export async function installMainCapture( eluStart: null, externalPeak: 0, finalized: false, + finalizationKey: {}, + finalizationTimedOut: false, finalizing: null, heapPeak: 0, kind, operationId: null, + ordinal: nextWorkerOrdinal, + pendingCount: 0, playlistId: null, + postGcHeapUnavailableReason: 'post-gc-probe-not-run', postGcHeapUsed: null, profileHandle: null, + profileCaptureKey: {}, profilePath: null, + profileResult: null, + resolvedProfileHandle: null, requestPerformance: [], responseEpochMs: null, - sampleBusy: false, + samplePromise: null, sampleTimer: null, snapshotPath: null, terminatedEpochMs: null, worker, }; + nextWorkerOrdinal += 1; records.set(worker, record); worker.on('message', (incoming) => { if (typeof incoming !== 'object' || incoming === null) { @@ -336,6 +425,10 @@ export async function installMainCapture( ) { const request = dbRequests.get(message['requestId']); if (request) { + request.record.pendingCount = Math.max( + 0, + request.record.pendingCount - 1 + ); record.requestPerformance.push({ identity: request.identity, operation: request.operation, @@ -362,39 +455,44 @@ export async function installMainCapture( }); return record; }; - const sampleWorker = async (record: WorkerRecord): Promise => { - if (record.sampleBusy || record.finalized) { - return; + const sampleWorker = (record: WorkerRecord): Promise => { + if (record.finalized) { + return Promise.resolve(); } - record.sampleBusy = true; - try { - const stats = await record.worker.getHeapStatistics?.(); - if (stats) { - record.heapPeak = Math.max( - record.heapPeak, - Number(stats.used_heap_size ?? 0) - ); - record.externalPeak = Math.max( - record.externalPeak, - Number(stats.external_memory ?? 0) - ); - } - const cpu = await record.worker.cpuUsage?.(); - if (cpu) { - record.cpuFirst ??= cpu; - record.cpuLast = cpu; - } - const elu = record.worker.performance?.eventLoopUtilization( - record.eluStart ?? undefined - ); - if (elu) { - record.elu = elu.utilization; - } - } catch { - // A one-shot worker may terminate between sampling calls. - } finally { - record.sampleBusy = false; + if (record.samplePromise) { + return record.samplePromise; } + record.samplePromise = (async () => { + try { + const stats = await record.worker.getHeapStatistics?.(); + if (stats) { + record.heapPeak = Math.max( + record.heapPeak, + Number(stats.used_heap_size ?? 0) + ); + record.externalPeak = Math.max( + record.externalPeak, + Number(stats.external_memory ?? 0) + ); + } + const cpu = await record.worker.cpuUsage?.(); + if (cpu) { + record.cpuFirst ??= cpu; + record.cpuLast = cpu; + } + const elu = record.worker.performance?.eventLoopUtilization( + record.eluStart ?? undefined + ); + if (elu) { + record.elu = elu.utilization; + } + } catch { + // A one-shot worker may terminate between sampling calls. + } + })().finally(() => { + record.samplePromise = null; + }); + return record.samplePromise; }; const resetWorkerForCapture = (record: WorkerRecord): void => { if (record.sampleTimer) { @@ -408,16 +506,23 @@ export async function installMainCapture( record.eluStart = null; record.externalPeak = 0; record.finalized = false; + record.finalizationKey = {}; + record.finalizationTimedOut = false; record.finalizing = null; record.heapPeak = 0; record.operationId = null; + record.pendingCount = 0; record.playlistId = null; + record.postGcHeapUnavailableReason = 'post-gc-probe-not-run'; record.postGcHeapUsed = null; record.profileHandle = null; + record.profileCaptureKey = {}; record.profilePath = null; + record.profileResult = null; + record.resolvedProfileHandle = null; record.requestPerformance = []; record.responseEpochMs = null; - record.sampleBusy = false; + record.samplePromise = null; record.sampleTimer = null; record.snapshotPath = null; record.terminatedEpochMs = null; @@ -446,58 +551,242 @@ export async function installMainCapture( ) { record.profilePath = path.join( state.outputDirectory, - `${record.kind}.cpuprofile` + `${record.kind}-${record.ordinal}.cpuprofile` + ); + const profileHandle = record.worker.startCpuProfile(); + const profileCaptureKey = record.profileCaptureKey; + record.profileHandle = profileHandle; + void profileHandle.then( + (handle) => { + if ( + record.profileCaptureKey === profileCaptureKey && + record.profileHandle === profileHandle + ) { + record.resolvedProfileHandle = handle; + } + }, + () => undefined ); - record.profileHandle = record.worker.startCpuProfile(); } }; - const finalizeWorker = (record: WorkerRecord): Promise => { - record.finalizing ??= (async () => { - if (record.sampleTimer) { - clearInterval(record.sampleTimer); - record.sampleTimer = null; + const stopWorkerSampling = (record: WorkerRecord): void => { + if (record.sampleTimer) { + clearInterval(record.sampleTimer); + record.sampleTimer = null; + } + }; + const joinFinalWorkerSample = async ( + record: WorkerRecord + ): Promise => { + if (record.samplePromise) { + await record.samplePromise; + } + await sampleWorker(record); + }; + const stopWorkerProfile = async ( + record: WorkerRecord, + waitForHandle = true + ): Promise => { + const profileCaptureKey = record.profileCaptureKey; + const profileHandle = record.profileHandle; + if (!profileHandle || !record.profilePath) { + return; + } + const handle = + record.resolvedProfileHandle ?? + (waitForHandle ? await profileHandle : null); + if (!handle) { + throw new Error( + 'cpu-profile-handle-not-ready-before-worker-termination' + ); + } + const profileResult = await handle.stop(); + if ( + record.profileCaptureKey !== profileCaptureKey || + record.finalizationTimedOut + ) { + return; + } + record.profileResult = profileResult; + }; + const writeWorkerProfile = (record: WorkerRecord): void => { + if (!record.profileHandle || !record.profilePath) { + return; + } + if (record.profileResult === null) { + throw new Error('worker-cpu-profile-result-missing'); + } + fs.writeFileSync( + record.profilePath, + JSON.stringify(normalizeWorkerCpuProfile(record.profileResult)) + ); + record.profileHandle = null; + record.profileResult = null; + record.resolvedProfileHandle = null; + }; + const flushWorkerProfiles = (workerRecords: WorkerRecord[]): void => { + for (const record of workerRecords) { + try { + writeWorkerProfile(record); + } catch (error: unknown) { + reportWorkerArtifactFailure(record, 'profile-write', error); } - await sampleWorker(record); - if (record.profileHandle && record.profilePath) { - const handle = await record.profileHandle; - const profile = await handle.stop(); - fs.writeFileSync( - record.profilePath, - JSON.stringify(normalizeWorkerCpuProfile(profile)) - ); - } - if ( - state.diagnostic && - typeof record.worker.getHeapSnapshot === 'function' - ) { - record.snapshotPath = path.join( - state.outputDirectory, - `${record.kind}.heapsnapshot` - ); - const snapshot = await record.worker.getHeapSnapshot(); - await streamPromises.pipeline( - snapshot, - fs.createWriteStream(record.snapshotPath) - ); - const postSnapshot = - await record.worker.getHeapStatistics?.(); - record.postGcHeapUsed = Number( - postSnapshot?.used_heap_size ?? 0 - ); - } - record.finalized = true; - })().catch((error: unknown) => { - recordTimeline({ - type: `worker-profile-error:${ - error instanceof Error - ? error.message.slice(0, 160) - : String(error).slice(0, 160) - }`, - }); - record.finalized = true; + } + }; + const takeWorkerHeapSnapshot = async ( + record: WorkerRecord + ): Promise => { + if (typeof record.worker.getHeapSnapshot !== 'function') { + return; + } + record.snapshotPath = path.join( + state.outputDirectory, + `${record.kind}-${record.ordinal}.heapsnapshot` + ); + const snapshot = await record.worker.getHeapSnapshot(); + await streamPromises.pipeline( + snapshot, + fs.createWriteStream(record.snapshotPath) + ); + }; + const reportWorkerArtifactFailure = ( + record: WorkerRecord, + stage: DatabaseWorkerPostGcAncillaryFailureStage, + error: unknown + ): void => { + if (stage === 'heap-snapshot') { + record.snapshotPath = null; + } else if (stage === 'profile-stop' || stage === 'profile-write') { + record.profilePath = null; + record.profileResult = null; + record.profileHandle = null; + record.resolvedProfileHandle = null; + } + recordTimeline({ + type: `worker-artifact-error:${stage}:${ + error instanceof Error + ? error.message.slice(0, 160) + : String(error).slice(0, 160) + }`, }); + }; + const finalizeDatabaseWorker = ( + record: WorkerRecord, + selectionUnavailableReason: DatabaseWorkerPostGcUnavailableReason | null + ): Promise => { + record.finalizing ??= databaseWorkerPostGcFinalizationApi + .finalize({ + finalizationKey: record.finalizationKey, + joinFinalSample: () => joinFinalWorkerSample(record), + probePostGc: () => + selectionUnavailableReason === null + ? databaseWorkerPostGcProbeApi.probe({ + createMessageChannel: () => + new workerThreads.MessageChannel(), + worker: { + postMessage( + message: unknown, + transferList: readonly unknown[] + ): void { + record.worker.postMessage( + message, + transferList as readonly WorkerTransferable[] + ); + }, + }, + }) + : Promise.resolve({ + postGcHeapUsedBytes: null, + unavailableReason: selectionUnavailableReason, + }), + reportAncillaryFailure: (stage, error) => + reportWorkerArtifactFailure(record, stage, error), + stopProfile: () => stopWorkerProfile(record), + stopSampling: () => stopWorkerSampling(record), + ...(state.diagnostic + ? { + takeHeapSnapshot: () => + takeWorkerHeapSnapshot(record), + } + : {}), + }) + .then((outcome) => { + record.postGcHeapUsed = outcome.postGcHeapUsedBytes; + record.postGcHeapUnavailableReason = + outcome.unavailableReason; + record.finalized = true; + }) + .catch((error: unknown) => { + record.postGcHeapUsed = null; + record.postGcHeapUnavailableReason = 'capture-failed'; + reportWorkerArtifactFailure(record, 'post-gc-probe', error); + record.finalized = true; + }); return record.finalizing; }; + const finalizeTerminatingWorker = ( + record: WorkerRecord, + waitForProfileHandle: boolean + ): Promise => { + if (record.finalizing) { + return record.finalizing; + } + const finalizationKey = record.finalizationKey; + record.finalizing = (async () => { + stopWorkerSampling(record); + record.postGcHeapUsed = null; + record.postGcHeapUnavailableReason = + 'worker-force-terminated-before-gc'; + try { + await stopWorkerProfile(record, waitForProfileHandle); + } catch (error: unknown) { + if ( + record.finalizationKey === finalizationKey && + !record.finalizationTimedOut + ) { + reportWorkerArtifactFailure( + record, + 'profile-stop', + error + ); + } + } + if (record.finalizationKey === finalizationKey) { + record.finalized = true; + } + })(); + return record.finalizing; + }; + const waitForTerminatingWorkerFinalization = async ( + record: WorkerRecord + ): Promise => { + if (!record.finalizing || record.finalizationTimedOut) { + return; + } + let timeout: NodeJS.Timeout | null = null; + let timedOut = false; + await Promise.race([ + record.finalizing, + new Promise((resolve) => { + timeout = setTimeout(() => { + timedOut = true; + resolve(); + }, 5_000); + }), + ]); + if (timeout) { + clearTimeout(timeout); + } + if (timedOut) { + record.finalizationTimedOut = true; + record.profileCaptureKey = {}; + reportWorkerArtifactFailure( + record, + 'profile-stop', + new Error('worker-profile-finalization-timeout') + ); + } + }; WorkerClass.prototype.postMessage = function ( this: InstrumentedWorker, @@ -509,6 +798,32 @@ export async function installMainCapture( if (value['type'] === 'request') { const kind = classifyRequest(value); if (kind) { + const requestDisposition = + kind === 'database.worker' + ? databaseWorkerPostGcCutoffApi.observeDatabaseRequest() + : null; + if (requestDisposition === 'after-cutoff') { + state.timeline.push({ + epochMs: nowEpochMs(), + operation: + typeof value['operation'] === 'string' + ? value['operation'] + : undefined, + playlistId: + readDatabasePlaylistId(value) ?? undefined, + requestId: + typeof value['requestId'] === 'string' + ? value['requestId'] + : undefined, + type: 'db-request-after-capture-cutoff', + }); + originalPostMessage.call( + this, + message, + transferList + ); + return; + } const record = createWorkerRecord(this, kind); startWorker(record); if ( @@ -551,6 +866,7 @@ export async function installMainCapture( operationId: payloadOperationId, operationIdUnavailableReason: null, }; + record.pendingCount += 1; dbRequests.set(value['requestId'], { identity, operation: value['operation'], @@ -590,7 +906,17 @@ export async function installMainCapture( if (!record || !isCurrentCaptureRecord(record)) { return originalTerminate.call(this); } + const terminationGeneration = record.captureGeneration; const markTerminated = (code: number): number => { + if ( + !workerTerminationGenerationApi.isCurrent({ + capturedGeneration: terminationGeneration, + currentGeneration: state.captureGeneration, + recordGeneration: record.captureGeneration, + }) + ) { + return code; + } record.terminatedEpochMs = nowEpochMs(); record.finalized = true; recordTimeline({ @@ -600,16 +926,17 @@ export async function installMainCapture( }); return code; }; - if (!state.diagnostic) { - if (record.sampleTimer) { - clearInterval(record.sampleTimer); - record.sampleTimer = null; - } - return originalTerminate.call(this).then(markTerminated); + if (state.diagnostic) { + const diagnosticTermination = (async (): Promise => { + void finalizeTerminatingWorker(record, true); + await waitForTerminatingWorkerFinalization(record); + return originalTerminate.call(this); + })(); + return diagnosticTermination.then(markTerminated); } - return finalizeWorker(record) - .then(() => originalTerminate.call(this)) - .then(markTerminated); + const termination = originalTerminate.call(this); + void finalizeTerminatingWorker(record, false); + return termination.then(markTerminated); }; const inspectorPost = ( @@ -669,6 +996,59 @@ export async function installMainCapture( state.mainPeakRss = Math.max(state.mainPeakRss, memory.rss); state.rendererWindowSession?.sample(); }; + const startCapture = async ( + options: MainCaptureStartOptions + ): Promise => { + state.rendererWindowSession?.detach(); + state.rendererWindowSession = null; + const rendererWindowSession = rendererWindowRssSessionApi.create({ + browserWindowFromId: (browserWindowId) => + BrowserWindow.fromId(browserWindowId), + browserWindowId: options.rendererWindowIdentity.browserWindowId, + getAppMetrics: () => app.getAppMetrics(), + rendererRssApi: rendererProcessRssApi, + webContentsId: options.rendererWindowIdentity.webContentsId, + }); + state.stopping = false; + state.captureGeneration += 1; + state.active = true; + databaseRequestIdentityCapture.start(); + dbRequests.clear(); + operationWorkers.clear(); + state.diagnostic = options.diagnostic; + state.outputDirectory = options.outputDirectory; + state.timeline = []; + state.mainPeakHeap = 0; + state.mainPeakRss = 0; + state.mainProfilePath = null; + state.mainSnapshotPath = null; + state.postGcHeap = null; + state.postGcRss = null; + state.rendererWindowSession = rendererWindowSession; + state.cpuStart = process.cpuUsage(); + state.eventLoopStart = perfHooks.performance.eventLoopUtilization(); + state.eventLoopDelay = perfHooks.monitorEventLoopDelay({ + resolution: 1, + }); + state.eventLoopDelay.enable(); + sampleMain(); + state.sampleTimer = setInterval(sampleMain, 20); + if (state.diagnostic) { + const session = new inspector.Session(); + session.connect(); + state.inspectorSession = session; + await inspectorPost(session, 'Profiler.enable'); + await inspectorPost(session, 'Profiler.start'); + state.mainProfilePath = path.join( + state.outputDirectory, + 'main.cpuprofile' + ); + state.mainSnapshotPath = path.join( + state.outputDirectory, + 'main.heapsnapshot' + ); + } + }; const api = { status: (): MainCaptureStatus => ({ @@ -703,62 +1083,27 @@ export async function installMainCapture( ).length, }), start: async (options: MainCaptureStartOptions): Promise => { - state.rendererWindowSession?.detach(); - state.rendererWindowSession = null; - const rendererWindowSession = - rendererWindowRssSessionApi.create({ - browserWindowFromId: (browserWindowId) => - BrowserWindow.fromId(browserWindowId), - browserWindowId: - options.rendererWindowIdentity.browserWindowId, - getAppMetrics: () => app.getAppMetrics(), - rendererRssApi: rendererProcessRssApi, - webContentsId: - options.rendererWindowIdentity.webContentsId, - }); - state.captureGeneration += 1; - state.active = true; - databaseRequestIdentityCapture.start(); - dbRequests.clear(); - operationWorkers.clear(); - state.diagnostic = options.diagnostic; - state.outputDirectory = options.outputDirectory; - state.timeline = []; - state.mainPeakHeap = 0; - state.mainPeakRss = 0; - state.postGcHeap = null; - state.postGcRss = null; - state.rendererWindowSession = rendererWindowSession; - state.cpuStart = process.cpuUsage(); - state.eventLoopStart = - perfHooks.performance.eventLoopUtilization(); - state.eventLoopDelay = perfHooks.monitorEventLoopDelay({ - resolution: 1, - }); - state.eventLoopDelay.enable(); - sampleMain(); - state.sampleTimer = setInterval(sampleMain, 20); - if (state.diagnostic) { - const session = new inspector.Session(); - session.connect(); - state.inspectorSession = session; - await inspectorPost(session, 'Profiler.enable'); - await inspectorPost(session, 'Profiler.start'); - state.mainProfilePath = path.join( - state.outputDirectory, - 'main.cpuprofile' - ); - state.mainSnapshotPath = path.join( - state.outputDirectory, - 'main.heapsnapshot' - ); - } + databaseWorkerPostGcCutoffApi.beginCapture(); + await startCapture(options); }, - stop: async (): Promise => { + stop: async ( + nextOptions?: MainCaptureStartOptions + ): Promise => { + databaseWorkerPostGcCutoffApi.beginStop(); + state.stopping = true; + state.active = false; + databaseRequestIdentityCapture.stop(); if (state.sampleTimer) { clearInterval(state.sampleTimer); state.sampleTimer = null; } + const currentWorkerRecords = [...records.values()].filter( + (record) => + record.captureGeneration === state.captureGeneration + ); + for (const record of currentWorkerRecords) { + stopWorkerSampling(record); + } state.eventLoopDelay?.disable(); sampleMain(); const rendererWindow: RendererWindowRssSessionMetrics | null = @@ -782,15 +1127,6 @@ export async function installMainCapture( ? null : elu.utilization; const delay = state.eventLoopDelay; - const databaseRecords = [...records.values()].filter( - (record) => - record.captureGeneration === state.captureGeneration && - record.kind === 'database.worker' && - record.sampleTimer !== null - ); - await Promise.all( - databaseRecords.map((record) => finalizeWorker(record)) - ); const session = state.inspectorSession ?? new inspector.Session(); if (!state.inspectorSession) { @@ -806,6 +1142,41 @@ export async function installMainCapture( JSON.stringify(result['profile']) ); } + await Promise.all( + currentWorkerRecords + .filter( + (record) => + record.kind === 'playlist-refresh.worker' && + record.finalizing !== null + ) + .map((record) => + waitForTerminatingWorkerFinalization(record) + ) + ); + const databaseSelection = + databaseWorkerPostGcSelectionApi.select( + currentWorkerRecords, + state.captureGeneration + ); + const currentDatabaseRecords = currentWorkerRecords.filter( + (record) => record.kind === 'database.worker' + ); + if (databaseSelection.selected) { + await finalizeDatabaseWorker( + databaseSelection.selected, + null + ); + } else if (databaseSelection.unavailableReason !== null) { + await Promise.all( + currentDatabaseRecords.map((record) => + finalizeDatabaseWorker( + record, + databaseSelection.unavailableReason + ) + ) + ); + } + flushWorkerProfiles(currentWorkerRecords); await inspectorPost(session, 'HeapProfiler.enable'); await inspectorPost(session, 'HeapProfiler.collectGarbage'); const postGc = process.memoryUsage(); @@ -830,15 +1201,20 @@ export async function installMainCapture( } session.disconnect(); state.inspectorSession = null; - state.active = false; - databaseRequestIdentityCapture.stop(); + const cutoff = databaseWorkerPostGcCutoffApi.snapshot(); + if (cutoff.lateRequestCount > 0) { + for (const record of currentDatabaseRecords) { + record.postGcHeapUsed = null; + record.postGcHeapUnavailableReason = + 'database-worker-activity-after-cutoff'; + } + } - const workers = [...records.values()] + const workers = currentWorkerRecords .filter( (record) => - record.kind === 'playlist-refresh.worker' || - record.heapPeak > 0 || - record.requestPerformance.length > 0 + record.kind === 'database.worker' || + record.kind === 'playlist-refresh.worker' ) .map((record) => ({ captureGeneration: record.captureGeneration, @@ -856,9 +1232,12 @@ export async function installMainCapture( eventLoopUtilization: record.elu, kind: record.kind, operationId: record.operationId, + ordinal: record.ordinal, peakExternalBytes: record.externalPeak, peakHeapUsedBytes: record.heapPeak, playlistId: record.playlistId, + postGcHeapUnavailableReason: + record.postGcHeapUnavailableReason, postGcHeapUsedBytes: record.postGcHeapUsed, profilePath: record.profilePath, responseEpochMs: record.responseEpochMs, @@ -877,7 +1256,7 @@ export async function installMainCapture( success: request.success, })), })); - return { + const transport: MainCaptureGenerationTransport = { captureGeneration: state.captureGeneration, metrics: { cpuProfilePath: state.mainProfilePath, @@ -910,6 +1289,52 @@ export async function installMainCapture( }, workers, }; + let rollover: MainCaptureRolloverStatus | null = null; + if (nextOptions) { + const databaseRecord = + currentDatabaseRecords.length === 1 + ? currentDatabaseRecords[0] + : null; + const nextCaptureUnavailableReason = + cutoff.lateRequestCount > 0 + ? 'database-worker-activity-after-cutoff' + : dbRequests.size > 0 + ? 'database-worker-not-idle' + : currentDatabaseRecords.length === 0 + ? 'database-worker-missing' + : currentDatabaseRecords.length > 1 + ? 'multiple-database-workers' + : !databaseRecord || + !Number.isSafeInteger( + databaseRecord.postGcHeapUsed + ) || + Number(databaseRecord.postGcHeapUsed) < + 0 || + databaseRecord.postGcHeapUnavailableReason !== + null + ? (databaseRecord?.postGcHeapUnavailableReason ?? + 'post-gc-capture-invalid') + : null; + if (nextCaptureUnavailableReason === null) { + databaseWorkerPostGcCutoffApi.rolloverCapture(); + await startCapture(nextOptions); + rollover = { + nextCaptureStarted: true, + nextCaptureUnavailableReason: null, + }; + } else { + databaseWorkerPostGcCutoffApi.finishStop(); + state.stopping = false; + rollover = { + nextCaptureStarted: false, + nextCaptureUnavailableReason, + }; + } + } else { + databaseWorkerPostGcCutoffApi.finishStop(); + state.stopping = false; + } + return { ...transport, rollover }; }, }; target[input.stateKey] = api; @@ -944,6 +1369,33 @@ export async function readMainCaptureStatus( }, MAIN_CAPTURE_STATE_KEY); } +export async function rolloverMainCapture( + electronApp: ElectronApplication, + options: MainCaptureStartOptions +): Promise { + const transport = await electronApp.evaluate( + async (_electron, input) => { + const target = globalThis as unknown as Record; + const api = target[input.stateKey] as { + stop( + nextOptions?: MainCaptureStartOptions + ): Promise; + }; + return api.stop(input.options); + }, + { options, stateKey: MAIN_CAPTURE_STATE_KEY } + ); + if (transport.rollover === null) { + throw new Error('main-capture-rollover-status-missing'); + } + return Object.freeze({ + completedCapture: selectMainCaptureGeneration(transport), + nextCaptureStarted: transport.rollover.nextCaptureStarted, + nextCaptureUnavailableReason: + transport.rollover.nextCaptureUnavailableReason, + }); +} + export async function stopMainCapture( electronApp: ElectronApplication ): Promise { @@ -951,7 +1403,7 @@ export async function stopMainCapture( async (_electron, stateKey) => { const target = globalThis as unknown as Record; const api = target[stateKey] as { - stop(): Promise; + stop(): Promise; }; return api.stop(); }, 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 index a1cfb2a64..c96ec7f03 100644 --- a/apps/electron-backend-e2e/src/performance/worker-request-performance.spec.ts +++ b/apps/electron-backend-e2e/src/performance/worker-request-performance.spec.ts @@ -12,6 +12,7 @@ import { type WorkerRequestPerformanceMetrics, } from './m3u-refresh-cancellation-contract'; import { + type MainCaptureGenerationTransport, normalizeWorkerRequestPerformanceOutcome, selectMainCaptureGeneration, } from './worker-request-performance'; @@ -217,7 +218,6 @@ test('main capture retains raw request identity, timestamps, metrics, and reason ); 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/); @@ -294,6 +294,53 @@ test('capture-generation selection excludes a pre-start seed worker', () => { assert.equal(selected.workers[0]?.terminatedEpochMs, 2_100); }); +test('fails closed for incoherent worker post-GC metrics at the Electron boundary', () => { + const invalidOutcomes = [ + { + postGcHeapUnavailableReason: 'gc-unavailable', + postGcHeapUsedBytes: 100, + }, + { + postGcHeapUnavailableReason: null, + postGcHeapUsedBytes: null, + }, + { + postGcHeapUnavailableReason: null, + postGcHeapUsedBytes: -1, + }, + { + postGcHeapUnavailableReason: 'unknown-reason', + postGcHeapUsedBytes: null, + }, + ] as const; + + for (const outcome of invalidOutcomes) { + const selected = selectMainCaptureGeneration({ + captureGeneration: 3, + metrics: measuredIteration([]).main, + workers: [ + { + captureGeneration: 3, + metrics: { + ...databaseWorker([]), + ...outcome, + }, + requests: [], + }, + ], + } as unknown as MainCaptureGenerationTransport); + const worker = selected.workers[0] as WorkerCaptureMetrics & { + readonly postGcHeapUnavailableReason: string | null; + }; + + assert.equal(worker.postGcHeapUsedBytes, null); + assert.equal( + worker.postGcHeapUnavailableReason, + 'post-gc-capture-invalid' + ); + } +}); + test('raw outcomes retain missing and malformed captures while summaries exclude their null metrics', () => { const valid = normalizeWorkerRequestPerformanceOutcome( requestTransport(requestPerformance('request-1', 5, 100)) diff --git a/apps/electron-backend-e2e/src/performance/worker-request-performance.ts b/apps/electron-backend-e2e/src/performance/worker-request-performance.ts index edf157282..d348cd71c 100644 --- a/apps/electron-backend-e2e/src/performance/worker-request-performance.ts +++ b/apps/electron-backend-e2e/src/performance/worker-request-performance.ts @@ -2,8 +2,10 @@ import type { EventLoopDelayMetrics, MainCaptureMetrics, WorkerCaptureMetrics, + WorkerPostGcHeapCapture, WorkerRequestPerformanceMetrics, } from './m3u-refresh-cancellation-contract'; +import { WORKER_POST_GC_HEAP_UNAVAILABLE_REASON } from './m3u-refresh-cancellation-contract'; export interface WorkerRequestPerformanceOutcomeTransport { readonly operation: string | null; @@ -62,6 +64,9 @@ const THREAD_CPU_REASONS = new Set([ 'thread-cpu-usage-invalid', 'thread-cpu-usage-unavailable', ]); +const WORKER_POST_GC_HEAP_UNAVAILABLE_REASONS = new Set( + Object.values(WORKER_POST_GC_HEAP_UNAVAILABLE_REASON) +); function isRecord(input: unknown): input is Record { return typeof input === 'object' && input !== null && !Array.isArray(input); @@ -84,6 +89,38 @@ function isNullableReason( return value === null || (typeof value === 'string' && allowed.has(value)); } +function normalizeWorkerPostGcHeapCapture( + input: Record +): WorkerPostGcHeapCapture { + const postGcHeapUnavailableReason = input['postGcHeapUnavailableReason']; + const postGcHeapUsedBytes = input['postGcHeapUsedBytes']; + if ( + Number.isSafeInteger(postGcHeapUsedBytes) && + Number(postGcHeapUsedBytes) >= 0 && + postGcHeapUnavailableReason === null + ) { + return { + postGcHeapUnavailableReason: null, + postGcHeapUsedBytes: Number(postGcHeapUsedBytes), + }; + } + if ( + postGcHeapUsedBytes === null && + typeof postGcHeapUnavailableReason === 'string' && + WORKER_POST_GC_HEAP_UNAVAILABLE_REASONS.has(postGcHeapUnavailableReason) + ) { + return { + postGcHeapUnavailableReason, + postGcHeapUsedBytes: null, + } as WorkerPostGcHeapCapture; + } + return { + postGcHeapUnavailableReason: + WORKER_POST_GC_HEAP_UNAVAILABLE_REASON.INVALID_CAPTURE, + postGcHeapUsedBytes: null, + }; +} + function parseEventLoopDelay( input: unknown ): EventLoopDelayMetrics | null | undefined { @@ -280,8 +317,12 @@ export function selectMainCaptureGeneration( ); const onlyRequest = requests.length === 1 ? (requests[0] ?? null) : null; + const postGcHeap = normalizeWorkerPostGcHeapCapture( + worker.metrics as unknown as Record + ); return { ...worker.metrics, + ...postGcHeap, eventLoopDelay: onlyRequest?.eventLoopDelay ?? null, eventLoopDelayUnavailableReason: onlyRequest ? (onlyRequest.eventLoopDelayUnavailableReason ?? diff --git a/apps/electron-backend-e2e/src/performance/worker-termination-generation.spec.ts b/apps/electron-backend-e2e/src/performance/worker-termination-generation.spec.ts new file mode 100644 index 000000000..f1ac4d93a --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/worker-termination-generation.spec.ts @@ -0,0 +1,102 @@ +/* eslint-disable playwright/expect-expect -- This is a Node assertion-based performance contract test. */ +import assert from 'node:assert/strict'; +import { readFileSync } from 'node:fs'; +import test from 'node:test'; + +interface WorkerTerminationGenerationApi { + isCurrent(input: { + readonly capturedGeneration: number | null; + readonly currentGeneration: number; + readonly recordGeneration: number | null; + }): boolean; +} + +interface WorkerTerminationGenerationModule { + createWorkerTerminationGenerationApi?: () => WorkerTerminationGenerationApi; +} + +const modulePromise = import( + new URL('./worker-termination-generation.ts', import.meta.url).href +) + .then((module) => module as WorkerTerminationGenerationModule) + .catch(() => null); + +test('rejects a delayed termination completion after capture rollover or record reuse', async () => { + const module = await modulePromise; + assert.ok(module, 'worker termination generation guard must exist'); + const api = module.createWorkerTerminationGenerationApi?.(); + assert.ok(api); + + assert.equal( + api.isCurrent({ + capturedGeneration: 4, + currentGeneration: 4, + recordGeneration: 4, + }), + true + ); + assert.equal( + api.isCurrent({ + capturedGeneration: 4, + currentGeneration: 5, + recordGeneration: 4, + }), + false, + 'a late completion must not write into the next global capture generation' + ); + assert.equal( + api.isCurrent({ + capturedGeneration: 4, + currentGeneration: 5, + recordGeneration: 5, + }), + false, + 'a reused record must not be finalized by the prior termination promise' + ); +}); + +test('fails closed for missing or invalid generation identities', async () => { + const module = await modulePromise; + assert.ok(module); + const api = module.createWorkerTerminationGenerationApi?.(); + assert.ok(api); + + for (const capturedGeneration of [null, 0, -1, 1.5, Number.NaN]) { + assert.equal( + api.isCurrent({ + capturedGeneration, + currentGeneration: 1, + recordGeneration: 1, + }), + false + ); + } +}); + +test('the Electron termination callback applies the generation guard before mutating capture state', () => { + const source = readFileSync( + new URL('./m3u-refresh-main-capture.ts', import.meta.url), + 'utf8' + ); + const terminateStart = source.indexOf( + 'WorkerClass.prototype.terminate = function' + ); + const terminateEnd = source.indexOf('const inspectorPost', terminateStart); + const terminateBlock = source.slice(terminateStart, terminateEnd); + const terminationGeneration = terminateBlock.indexOf( + 'const terminationGeneration = record.captureGeneration' + ); + const generationGuard = terminateBlock.indexOf( + 'workerTerminationGenerationApi.isCurrent' + ); + const terminationMutation = terminateBlock.indexOf( + 'record.terminatedEpochMs = nowEpochMs()' + ); + + assert.ok( + terminationGeneration >= 0 && + terminationGeneration < generationGuard && + generationGuard < terminationMutation, + 'late termination completion must be generation-gated before mutating the record or timeline' + ); +}); diff --git a/apps/electron-backend-e2e/src/performance/worker-termination-generation.ts b/apps/electron-backend-e2e/src/performance/worker-termination-generation.ts new file mode 100644 index 000000000..68f77efbd --- /dev/null +++ b/apps/electron-backend-e2e/src/performance/worker-termination-generation.ts @@ -0,0 +1,24 @@ +export interface WorkerTerminationGenerationInput { + readonly capturedGeneration: number | null; + readonly currentGeneration: number; + readonly recordGeneration: number | null; +} + +export interface WorkerTerminationGenerationApi { + isCurrent(input: WorkerTerminationGenerationInput): boolean; +} + +export function createWorkerTerminationGenerationApi(): WorkerTerminationGenerationApi { + const isCaptureGeneration = (value: number | null): value is number => + Number.isSafeInteger(value) && Number(value) > 0; + + return Object.freeze({ + isCurrent(input: WorkerTerminationGenerationInput): boolean { + return ( + isCaptureGeneration(input.capturedGeneration) && + input.capturedGeneration === input.currentGeneration && + input.capturedGeneration === input.recordGeneration + ); + }, + }); +} diff --git a/docs/architecture/m3u-playlist-module.md b/docs/architecture/m3u-playlist-module.md index a6bb2cd23..6456e626c 100644 --- a/docs/architecture/m3u-playlist-module.md +++ b/docs/architecture/m3u-playlist-module.md @@ -137,8 +137,59 @@ 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 +Database-worker post-GC heap is not inferred from a heap snapshot. The harness +selects exactly one current-generation database worker independently of its +sampling timer, waits for any in-flight sample plus one final sample, stops the +worker CPU profile, sends an explicit-GC one-shot probe over a transferred +`MessagePort`, and takes an optional diagnostic snapshot only afterward. +Capture stop first closes a synchronous request cutoff and stops the main CPU +profile, so worker profile writes and heap snapshots cannot appear as +application stacks in `main.cpuprofile`. A database request observed after that +cutoff cannot restart worker sampling or profiling; it records +`database-worker-activity-after-cutoff` and invalidates the post-GC value. +Worker artifacts include a stable isolate ordinal in their filename so an +invalid multiple-worker capture cannot overwrite another isolate's profile. +The raw heap and unavailability reason form a strict XOR; playlist workers +terminated during cancellation report `worker-force-terminated-before-gc` +instead of a snapshot-derived heap. Warm-up and measured cancellation invoke +the real worker termination before best-effort finalization so profiling cannot +delay headline acknowledgement metrics. The separate diagnostic iteration +stops its CPU profile before termination; its timings are excluded from +headline distributions. That pre-termination drain is bounded; timeout +invalidates the profile and worker termination still proceeds. Late +termination completions are generation-gated so they cannot mutate a reused +record or append timeline events after an atomic capture rollover. The +benchmark fails closed unless the diagnostic iteration observed cancellation, +its profile is inside the iteration directory and parses with non-empty nodes +and samples, and the dying worker has no heap snapshot. + +Headline worker memory distributions include every matching isolate and +exclude only nullable unavailable values. This aggregation is not used as a +validity shortcut. Runs that reach database persistence must independently +contain exactly one database worker with a coherent numeric post-GC heap. +Parsing cancellation happens before the renderer dispatches persistence, so a +run with no database request is explicitly `N/A: +operation-cancelled-before-database-phase`, not a missing capture. Any database +activity in a run that observed the cancellation effect is unrelated +contamination and invalidates the run. Duplicate, busy, late, timed-out, or +malformed captures are listed under +`summary.validity.databaseWorkerPostGc` and invalidate the benchmark. The +benchmark writes `manifest.json` and `summary.json` first, then fails the formal +run so the raw evidence remains available without supporting a before/after +claim. A formal cancellation run also requires the cancellation effect in every +measured iteration. Database-worker values are comparable only when every +measured iteration reaches that phase and has one valid capture. + +Seed setup runs inside its own capture generation. The harness requires an +observed seed database request, a completed playlist upsert, and no pending +request, then persists that generation as `seed-main-capture.json`. Main +capture finalization validates one idle database worker with a coherent +explicit-GC result and atomically rolls the clean cutoff into the measured +generation before yielding. A request before rollover remains late and rejects +the transition; a request after rollover belongs to the measured generation. +This prevents seed writes or an unobserved stop/start gap from crossing into +measured main RSS, worker peaks, or profiles; generation selection also +excludes the seed worker record. 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 diff --git a/docs/architecture/sqlite-db-worker.md b/docs/architecture/sqlite-db-worker.md index 52bdbb860..fc18d14d2 100644 --- a/docs/architecture/sqlite-db-worker.md +++ b/docs/architecture/sqlite-db-worker.md @@ -147,6 +147,40 @@ profiling queue. If captures overlap, every overlapping response carries CPU, ELU, and event-loop-delay values are `null`. This avoids assigning shared worker activity to one request while preserving normal worker concurrency. +### Opt-in post-GC heap capture + +The same `IPTVNATOR_PERF_WORKER_PROFILING=1` opt-in enables a development/test +one-shot `performance:collect-post-gc-heap` request. The benchmark transfers a +dedicated `MessagePort`; the database worker accepts the request only while no +request performance capture is active, calls exposed `globalThis.gc()`, reads +its own V8 isolate through `v8.getHeapStatistics()`, posts one result, and +closes the port. Production launches do not expose GC or send this request. + +The response is a strict XOR: either a non-negative +`postGcHeapUsedBytes` with a `null` reason, or a `null` heap with a fixed +unavailability reason. Disabled profiling, a busy worker, unavailable GC, and +capture failure all fail closed without changing the database operation or +terminating the worker. The performance launcher supplies +`--js-flags=--expose-gc` to Electron itself; worker `execArgv` remains +untouched. + +The main-process benchmark selects exactly one current-generation database +worker, waits for its final sampling call, stops its CPU profile, performs the +explicit-GC probe, and only then takes an optional diagnostic heap snapshot. +Profile or snapshot failure cannot overwrite an already captured post-GC +value. Capture stop closes a synchronous cutoff before awaiting profile or +snapshot work. A database request after that cutoff cannot restart sampling; +it records `database-worker-activity-after-cutoff` and invalidates the result. +Missing, multiple, busy, timed-out, or malformed worker captures remain raw +nullable outcomes and make a database-applicable measured run invalid for +comparison. A scenario cancelled before its database phase reports the worker +metric as not applicable rather than manufacturing an idle database request. +Before a measured generation, the M3U benchmark persists and validates its seed +capture, then atomically rolls a clean stopping cutoff into the next active +generation. This closes the otherwise unobservable gap between separate +stop/start calls: pre-rollover requests remain late and fatal, while +post-rollover requests are attributed to the new generation. + ### Progress event contract The worker now emits request-scoped events with: @@ -564,7 +598,8 @@ executed inside a synchronous transaction callback: ```ts // favorites is playlist-scoped: filter by (contentId, playlistId), otherwise // a same-contentId favorite in another playlist gets rewritten too. -const stmt = db.update(schema.favorites) +const stmt = db + .update(schema.favorites) .set({ position: sql`${sql.placeholder('position')}` }) .where( and(