fix(perf): finalize exact worker capture safely

This commit is contained in:
4gray committed 2026-07-27 02:15:59 +02:00
1 parent d3fdbaa2ef
commit 554eecfee0
23 files changed
+2685 -181

No files matched your search

@@ -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<CutoffApi> {
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/);
});
@@ -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);
}
@@ -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<void>;
readonly probePostGc: () => Promise<PostGcOutcome>;
readonly reportAncillaryFailure: (
stage: AncillaryFailureStage,
error: unknown
) => void;
readonly stopProfile: () => Promise<void>;
readonly stopSampling: () => void;
readonly takeHeapSnapshot?: () => Promise<void>;
}
interface FinalizationApi {
finalize(input: FinalizationInput): Promise<PostGcOutcome>;
}
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<void>;
readonly resolve: () => void;
} {
let resolve!: () => void;
const promise = new Promise<void>((resolvePromise) => {
resolve = resolvePromise;
});
return { promise, resolve };
}
async function restoreSerializableApi(): Promise<FinalizationApi> {
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'/
);
});
@@ -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<void>;
readonly probePostGc: () => Promise<DatabaseWorkerPostGcFinalizationOutcome>;
readonly reportAncillaryFailure: (
stage: DatabaseWorkerPostGcAncillaryFailureStage,
error: unknown
) => void;
readonly stopProfile: () => Promise<void>;
readonly stopSampling: () => void;
readonly takeHeapSnapshot?: () => Promise<void>;
}
export interface DatabaseWorkerPostGcFinalizationApi {
finalize(
input: DatabaseWorkerPostGcFinalizationInput
): Promise<DatabaseWorkerPostGcFinalizationOutcome>;
}
export function createDatabaseWorkerPostGcFinalizationApi(): DatabaseWorkerPostGcFinalizationApi {
const finalizations = new WeakMap<
object,
Promise<DatabaseWorkerPostGcFinalizationOutcome>
>();
const helpers = {
async run(
input: DatabaseWorkerPostGcFinalizationInput
): Promise<DatabaseWorkerPostGcFinalizationOutcome> {
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<DatabaseWorkerPostGcFinalizationOutcome> {
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 });
}
@@ -63,11 +63,25 @@ interface ProbeModule {
createDatabaseWorkerPostGcProbeApi?: () => ProbeApi;
}
interface ProductionWorkerPostGcModule {
DATABASE_WORKER_POST_GC_HEAP_UNAVAILABLE_REASON?: Readonly<
Record<string, WorkerUnavailableReason>
>;
}
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<WorkerUnavailableReason>([
'capture-failed',
'gc-unavailable',
'profiling-disabled',
'worker-busy',
])
);
for (const unavailableReason of reasons) {
const harness = createProbeHarness();
@@ -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 & {
@@ -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/
);
});
@@ -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<string>(
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,
});
}
@@ -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<void>;
}
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'
);
});
@@ -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<void> {
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');
}
}
@@ -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;
};
}
@@ -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),
}),
});
}
@@ -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<void> {
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<void> {
async function waitForSeedSettlement(app: LaunchedElectronApp): Promise<void> {
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);
}
File diff suppressed because it is too large. Load diff
@@ -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))
@@ -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<string>(
Object.values(WORKER_POST_GC_HEAP_UNAVAILABLE_REASON)
);
function isRecord(input: unknown): input is Record<string, unknown> {
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<string, unknown>
): 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<string, unknown>
);
return {
...worker.metrics,
...postGcHeap,
eventLoopDelay: onlyRequest?.eventLoopDelay ?? null,
eventLoopDelayUnavailableReason: onlyRequest
? (onlyRequest.eventLoopDelayUnavailableReason ??
@@ -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'
);
});
@@ -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
);
},
});
}