test(stalker): add replay listeners and control plane

This commit is contained in:
4gray committed 2026-07-27 09:05:40 +02:00
1 parent e43d94ded4
commit af0f210ea5
10 files changed
+2698 -11

No files matched your search

@@ -0,0 +1,572 @@
/* eslint-disable max-lines -- The control-plane security boundary is clearer as explicit wire-level cases. */
import {
mkdtemp,
mkdir,
rm,
symlink,
writeFile,
} from 'node:fs/promises';
import http from 'node:http';
import os from 'node:os';
import path from 'node:path';
import {
createReplayControlPlaneCapability,
InMemoryReplayFixtureRepository,
RepositoryReplayFixtureRepository,
ReplayFixtureRepositoryError,
REPLAY_CONTROL_CAPABILITY_HEADER,
REPLAY_CONTROL_MAX_BODY_BYTES,
startReplayControlPlane,
} from './replay-control-plane.js';
import { ReplayFixtureV1 } from './replay.types.js';
const TEST_CAPABILITY = 'test-capability-0123456789abcdef';
function replayFixture(): ReplayFixtureV1 {
return {
schemaVersion: 1,
scenarioId: 'control-contract',
origins: { portal: {}, auth: {} },
entry: { origin: 'portal', path: '/entry' },
expectedEndpoint: { origin: 'portal', path: '/entry' },
initialState: 'start',
terminalState: 'complete',
failOnUnexpectedRequest: true,
symbols: [
{ kind: 'generate', symbol: 'mac', valueKind: 'mac' },
{
kind: 'generate',
symbol: 'username',
valueKind: 'credential',
},
{
kind: 'generate',
symbol: 'password',
valueKind: 'credential',
},
{ kind: 'generate', symbol: 'token', valueKind: 'token' },
{ kind: 'generate', symbol: 'random', valueKind: 'random' },
{ kind: 'generate', symbol: 'cookie', valueKind: 'cookie' },
{
kind: 'generate',
symbol: 'correlation',
valueKind: 'correlation',
},
],
phases: [
{
name: 'request',
state: 'start',
nextState: 'complete',
mode: 'ordered',
expectations: [
{
id: 'entry',
operation: 'entry',
origin: 'portal',
method: 'GET',
path: '/entry',
request: {
query: {
exact: {},
present: [],
absent: [],
},
headers: {
exact: {},
present: [],
absent: [],
},
cookies: {
exact: {},
present: [],
absent: [],
attributes: {},
},
body: { kind: 'absent' },
},
response: {
status: 200,
headers: {
'content-type': ['application/json'],
},
body: {
kind: 'json',
value: { js: true },
},
},
cardinality: { min: 1, max: 1 },
},
],
},
],
};
}
interface ControlResponse {
status: number;
headers: http.IncomingHttpHeaders;
body: unknown;
}
function controlRequest(
controlUrl: string,
action: 'create' | 'finalize' | 'dispose',
body: unknown,
options: {
capability?: string;
host?: string;
method?: string;
contentType?: string;
rawBody?: string;
origin?: string;
} = {}
): Promise<ControlResponse> {
const url = new URL(`/v1/replay/${action}`, controlUrl);
const serialized = options.rawBody ?? JSON.stringify(body);
const headers: http.OutgoingHttpHeaders = {
host: options.host ?? url.host,
'content-type': options.contentType ?? 'application/json',
'content-length': Buffer.byteLength(serialized),
};
if (options.capability !== undefined) {
headers[REPLAY_CONTROL_CAPABILITY_HEADER] = options.capability;
}
if (options.origin !== undefined) {
headers['origin'] = options.origin;
}
return new Promise((resolve, reject) => {
const request = http.request(
url,
{
method: options.method ?? 'POST',
headers,
},
(response) => {
const chunks: Buffer[] = [];
response.on('data', (chunk: Buffer) => chunks.push(chunk));
response.on('end', () => {
const text = Buffer.concat(chunks).toString('utf8');
resolve({
status: response.statusCode ?? 0,
headers: response.headers,
body: text.length === 0 ? null : JSON.parse(text),
});
});
}
);
request.on('error', reject);
request.end(serialized);
});
}
function authorizedRequest(
controlUrl: string,
action: 'create' | 'finalize' | 'dispose',
body: unknown
): Promise<ControlResponse> {
return controlRequest(controlUrl, action, body, {
capability: TEST_CAPABILITY,
});
}
function streamingControlRequest(controlUrl: string): {
request: http.ClientRequest;
response: Promise<ControlResponse>;
} {
const url = new URL('/v1/replay/create', controlUrl);
let requestResult: http.ClientRequest | undefined;
const response = new Promise<ControlResponse>((resolve, reject) => {
let responseStarted = false;
const request = http.request(
url,
{
method: 'POST',
headers: {
host: url.host,
'content-type': 'application/json',
[REPLAY_CONTROL_CAPABILITY_HEADER]: TEST_CAPABILITY,
},
},
(incoming) => {
responseStarted = true;
const chunks: Buffer[] = [];
incoming.on('data', (chunk: Buffer) => chunks.push(chunk));
incoming.on('end', () => {
const text = Buffer.concat(chunks).toString('utf8');
resolve({
status: incoming.statusCode ?? 0,
headers: incoming.headers,
body: text.length === 0 ? null : JSON.parse(text),
});
});
}
);
request.on('error', (error) => {
if (!responseStarted) {
reject(error);
}
});
requestResult = request;
});
if (requestResult === undefined) {
throw new Error('Streaming control request was not initialized.');
}
return { request: requestResult, response };
}
describe('replay control-plane capability and wire boundary', () => {
it('creates cryptographically-sized process-local capabilities', () => {
const first = createReplayControlPlaneCapability();
const second = createReplayControlPlaneCapability();
expect(first).toMatch(/^[A-Za-z0-9_-]{43}$/);
expect(second).toMatch(/^[A-Za-z0-9_-]{43}$/);
expect(first).not.toBe(second);
});
it('binds only loopback, validates the exact Host and requires the capability on every action', async () => {
const control = await startReplayControlPlane({
capability: TEST_CAPABILITY,
repository: new InMemoryReplayFixtureRepository({
'authentication/ready': replayFixture(),
}),
});
try {
const url = new URL(control.url);
expect(url.hostname).toBe('127.0.0.1');
expect(Number(url.port)).toBeGreaterThan(0);
await expect(
controlRequest(
control.url,
'create',
{ fixtureId: 'authentication/ready' },
{
capability: TEST_CAPABILITY,
host: `localhost:${url.port}`,
}
)
).resolves.toMatchObject({
status: 400,
body: { error: { code: 'invalid-host' } },
});
await expect(
controlRequest(control.url, 'create', {
fixtureId: 'authentication/ready',
})
).resolves.toMatchObject({
status: 403,
body: { error: { code: 'invalid-capability' } },
});
await expect(
controlRequest(
control.url,
'create',
{ fixtureId: 'authentication/ready' },
{ capability: `${TEST_CAPABILITY}-wrong` }
)
).resolves.toMatchObject({
status: 403,
body: { error: { code: 'invalid-capability' } },
});
const created = await authorizedRequest(
control.url,
'create',
{ fixtureId: 'authentication/ready' }
);
const runId = (created.body as { runId: string }).runId;
await expect(
controlRequest(control.url, 'finalize', { runId })
).resolves.toMatchObject({
status: 403,
body: { error: { code: 'invalid-capability' } },
});
await expect(
controlRequest(control.url, 'dispose', { runId })
).resolves.toMatchObject({
status: 403,
body: { error: { code: 'invalid-capability' } },
});
await authorizedRequest(control.url, 'dispose', { runId });
} finally {
await control.close();
}
});
it('never enables CORS and accepts only POST application/json', async () => {
const control = await startReplayControlPlane({
capability: TEST_CAPABILITY,
repository: new InMemoryReplayFixtureRepository({
ready: replayFixture(),
}),
});
try {
const preflight = await controlRequest(
control.url,
'create',
{},
{
capability: TEST_CAPABILITY,
method: 'OPTIONS',
origin: 'https://example.invalid',
}
);
expect(preflight.status).toBe(405);
expect(preflight.headers).not.toHaveProperty(
'access-control-allow-origin'
);
expect(preflight.headers).not.toHaveProperty(
'access-control-allow-headers'
);
await expect(
controlRequest(
control.url,
'create',
{ fixtureId: 'ready' },
{
capability: TEST_CAPABILITY,
contentType: 'text/plain',
}
)
).resolves.toMatchObject({
status: 415,
body: { error: { code: 'json-required' } },
});
} finally {
await control.close();
}
});
it('enforces the 64 KiB raw cap before attempting JSON parsing', async () => {
const control = await startReplayControlPlane({
capability: TEST_CAPABILITY,
repository: new InMemoryReplayFixtureRepository({}),
});
try {
const oversizedInvalidJson = 'x'.repeat(
REPLAY_CONTROL_MAX_BODY_BYTES + 1
);
await expect(
controlRequest(
control.url,
'create',
{},
{
capability: TEST_CAPABILITY,
rawBody: oversizedInvalidJson,
}
)
).resolves.toMatchObject({
status: 413,
body: { error: { code: 'control-body-too-large' } },
});
const stream = streamingControlRequest(control.url);
stream.request.write(
Buffer.alloc(REPLAY_CONTROL_MAX_BODY_BYTES + 1, 120)
);
await expect(stream.response).resolves.toMatchObject({
status: 413,
body: { error: { code: 'control-body-too-large' } },
});
stream.request.destroy();
} finally {
await control.close();
}
});
it('times out a control body that never ends', async () => {
const control = await startReplayControlPlane({
capability: TEST_CAPABILITY,
repository: new InMemoryReplayFixtureRepository({}),
requestBodyTimeoutMs: 25,
});
try {
const stream = streamingControlRequest(control.url);
stream.request.write('{"fixtureId":');
await expect(stream.response).resolves.toMatchObject({
status: 408,
body: { error: { code: 'control-body-timeout' } },
});
stream.request.destroy();
} finally {
await control.close();
}
});
});
describe('replay fixture repository allowlist', () => {
let temporaryRoot: string;
beforeEach(async () => {
temporaryRoot = await mkdtemp(
path.join(os.tmpdir(), 'iptvnator-replay-repository-')
);
});
afterEach(async () => {
await rm(temporaryRoot, { recursive: true, force: true });
});
it('loads only regular JSON fixtures beneath the configured repository root', async () => {
const repositoryRoot = path.join(temporaryRoot, 'fixtures');
const outsidePath = path.join(temporaryRoot, 'outside.json');
await mkdir(path.join(repositoryRoot, 'authentication'), {
recursive: true,
});
await writeFile(
path.join(repositoryRoot, 'authentication', 'ready.json'),
JSON.stringify(replayFixture())
);
await writeFile(outsidePath, JSON.stringify(replayFixture()));
await symlink(
outsidePath,
path.join(repositoryRoot, 'authentication', 'escape.json')
);
const repository = new RepositoryReplayFixtureRepository(
repositoryRoot
);
await expect(
repository.loadFixture('authentication/ready')
).resolves.toMatchObject({ scenarioId: 'control-contract' });
await expect(
repository.loadFixture('../outside')
).rejects.toMatchObject({
code: 'fixture-not-allowed',
});
await expect(
repository.loadFixture('/absolute')
).rejects.toBeInstanceOf(ReplayFixtureRepositoryError);
await expect(
repository.loadFixture('authentication/escape')
).rejects.toMatchObject({
code: 'fixture-not-allowed',
});
await expect(
repository.loadFixture('authentication/missing')
).rejects.toMatchObject({
code: 'fixture-not-allowed',
});
});
});
describe('replay control-plane lifecycle and public response', () => {
it('returns only opaque routing and filtered public inputs, then finalizes and disposes the run', async () => {
const control = await startReplayControlPlane({
capability: TEST_CAPABILITY,
repository: new InMemoryReplayFixtureRepository({
'authentication/ready': replayFixture(),
}),
});
try {
const created = await authorizedRequest(
control.url,
'create',
{ fixtureId: 'authentication/ready' }
);
expect(created.status).toBe(201);
expect(created.headers).not.toHaveProperty(
'access-control-allow-origin'
);
expect(created.body).toEqual({
runId: expect.stringMatching(/^run-[a-f0-9]{32}$/),
entryUrl: expect.stringMatching(
/^http:\/\/127\.0\.0\.1:\d+\//
),
origins: {
auth: expect.stringMatching(
/^http:\/\/127\.0\.0\.1:\d+\//
),
portal: expect.stringMatching(
/^http:\/\/127\.0\.0\.1:\d+\//
),
},
inputs: {
mac: expect.stringMatching(/^02(?::[a-f0-9]{2}){5}$/),
password: expect.stringMatching(/^test-credential-/),
username: expect.stringMatching(/^test-credential-/),
},
});
const createdText = JSON.stringify(created.body);
expect(createdText).not.toContain(TEST_CAPABILITY);
expect(createdText).not.toContain('test-token-');
expect(createdText).not.toContain('test-random-');
expect(createdText).not.toContain('test-cookie-');
expect(createdText).not.toContain('test-correlation-');
const publicResult = created.body as {
runId: string;
entryUrl: string;
};
await expect(fetch(publicResult.entryUrl)).resolves.toMatchObject({
status: 200,
});
await expect(
authorizedRequest(control.url, 'finalize', {
runId: publicResult.runId,
})
).resolves.toMatchObject({
status: 200,
body: {
ok: true,
ledger: {
operationCounts: { entry: 1 },
mismatchCounts: {},
terminalState: 'complete',
},
},
});
await expect(
authorizedRequest(control.url, 'dispose', {
runId: publicResult.runId,
})
).resolves.toMatchObject({
status: 200,
body: { disposed: true },
});
await expect(fetch(publicResult.entryUrl)).rejects.toThrow();
await expect(
authorizedRequest(control.url, 'finalize', {
runId: publicResult.runId,
})
).resolves.toMatchObject({
status: 404,
body: { error: { code: 'run-not-found' } },
});
} finally {
await control.close();
}
});
it('rejects unknown and malformed fixture IDs with one sanitized allowlist code', async () => {
const control = await startReplayControlPlane({
capability: TEST_CAPABILITY,
repository: new InMemoryReplayFixtureRepository({}),
});
try {
for (const fixtureId of [
'missing',
'../outside',
'/absolute',
'authentication//ready',
'authentication/../../outside',
]) {
await expect(
authorizedRequest(control.url, 'create', { fixtureId })
).resolves.toMatchObject({
status: 404,
body: { error: { code: 'fixture-not-allowed' } },
});
}
} finally {
await control.close();
}
});
});
@@ -0,0 +1,640 @@
/* eslint-disable max-lines -- Keep the loopback control-plane validation and lifecycle in one auditable boundary. */
import { randomBytes, timingSafeEqual } from 'node:crypto';
import { constants as fsConstants } from 'node:fs';
import { lstat, open, realpath } from 'node:fs/promises';
import type { FileHandle } from 'node:fs/promises';
import http, {
IncomingMessage,
Server,
ServerResponse,
} from 'node:http';
import { AddressInfo } from 'node:net';
import path from 'node:path';
import { REPLAY_MAX_FIXTURE_BYTES } from './replay.constants.js';
import {
createReplayServerRun,
ReplayServerRun,
} from './replay-server.js';
import {
parseReplayFixture,
parseReplayFixtureText,
ReplaySchemaError,
} from './replay-schema.js';
import { ReplayFixtureV1 } from './replay.types.js';
const LOOPBACK_ADDRESS = '127.0.0.1';
const CREATE_PATH = '/v1/replay/create';
const FINALIZE_PATH = '/v1/replay/finalize';
const DISPOSE_PATH = '/v1/replay/dispose';
const SAFE_FIXTURE_ID =
/^[a-z0-9][a-z0-9._-]{0,63}(?:\/[a-z0-9][a-z0-9._-]{0,63}){0,7}$/;
const SAFE_RUN_ID = /^run-[a-f0-9]{32}$/;
const DEFAULT_REQUEST_BODY_TIMEOUT_MS = 5000;
export const REPLAY_CONTROL_CAPABILITY_HEADER =
'x-iptvnator-replay-capability';
export const REPLAY_CONTROL_MAX_BODY_BYTES = 64 * 1024;
const CONTROL_ERROR_CODE = {
BODY_TOO_LARGE: 'control-body-too-large',
BODY_TIMEOUT: 'control-body-timeout',
ENDPOINT_NOT_FOUND: 'endpoint-not-found',
FIXTURE_INVALID: 'fixture-invalid',
FIXTURE_NOT_ALLOWED: 'fixture-not-allowed',
INTERNAL: 'control-internal-error',
INVALID_CAPABILITY: 'invalid-capability',
INVALID_HOST: 'invalid-host',
INVALID_JSON: 'invalid-json',
INVALID_REQUEST: 'invalid-request',
JSON_REQUIRED: 'json-required',
METHOD_NOT_ALLOWED: 'method-not-allowed',
RUN_NOT_FOUND: 'run-not-found',
} as const;
type ReplayControlErrorCode =
(typeof CONTROL_ERROR_CODE)[keyof typeof CONTROL_ERROR_CODE];
class ReplayControlError extends Error {
constructor(
readonly code: ReplayControlErrorCode,
readonly status: number
) {
super(`Replay control request rejected: ${code}.`);
this.name = 'ReplayControlError';
}
}
export type ReplayFixtureRepositoryErrorCode =
| 'fixture-not-allowed'
| 'fixture-invalid';
export class ReplayFixtureRepositoryError extends Error {
constructor(readonly code: ReplayFixtureRepositoryErrorCode) {
super(`Replay fixture repository rejected: ${code}.`);
this.name = 'ReplayFixtureRepositoryError';
}
}
export interface ReplayFixtureRepository {
loadFixture(fixtureId: string): Promise<ReplayFixtureV1>;
}
function assertSafeFixtureId(fixtureId: string): void {
if (!SAFE_FIXTURE_ID.test(fixtureId)) {
throw new ReplayFixtureRepositoryError('fixture-not-allowed');
}
}
export class InMemoryReplayFixtureRepository
implements ReplayFixtureRepository
{
private readonly fixtures: Readonly<Record<string, unknown>>;
constructor(fixtures: Readonly<Record<string, unknown>>) {
this.fixtures = Object.freeze({ ...fixtures });
}
async loadFixture(fixtureId: string): Promise<ReplayFixtureV1> {
assertSafeFixtureId(fixtureId);
if (!Object.prototype.hasOwnProperty.call(this.fixtures, fixtureId)) {
throw new ReplayFixtureRepositoryError('fixture-not-allowed');
}
try {
return parseReplayFixture(this.fixtures[fixtureId]);
} catch (error) {
if (error instanceof ReplaySchemaError) {
throw new ReplayFixtureRepositoryError('fixture-invalid');
}
throw error;
}
}
}
function isWithinRoot(root: string, candidate: string): boolean {
const relative = path.relative(root, candidate);
return (
relative.length > 0 &&
!relative.startsWith(`..${path.sep}`) &&
relative !== '..' &&
!path.isAbsolute(relative)
);
}
async function readBoundedFixture(handle: FileHandle): Promise<string> {
const chunks: Buffer[] = [];
const readBuffer = Buffer.alloc(64 * 1024);
let position = 0;
while (true) {
const { bytesRead } = await handle.read(
readBuffer,
0,
readBuffer.byteLength,
position
);
if (bytesRead === 0) {
break;
}
position += bytesRead;
if (position > REPLAY_MAX_FIXTURE_BYTES) {
throw new ReplayFixtureRepositoryError(
'fixture-not-allowed'
);
}
chunks.push(Buffer.from(readBuffer.subarray(0, bytesRead)));
}
return Buffer.concat(chunks).toString('utf8');
}
export class RepositoryReplayFixtureRepository
implements ReplayFixtureRepository
{
constructor(private readonly repositoryRoot: string) {}
async loadFixture(fixtureId: string): Promise<ReplayFixtureV1> {
assertSafeFixtureId(fixtureId);
try {
const root = await realpath(this.repositoryRoot);
const requestedPath = path.join(
root,
...fixtureId.split('/')
).concat('.json');
const requestedStats = await lstat(requestedPath);
if (
requestedStats.isSymbolicLink() ||
!requestedStats.isFile()
) {
throw new ReplayFixtureRepositoryError(
'fixture-not-allowed'
);
}
const resolvedPath = await realpath(requestedPath);
if (!isWithinRoot(root, resolvedPath)) {
throw new ReplayFixtureRepositoryError(
'fixture-not-allowed'
);
}
const handle = await open(
requestedPath,
fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW
);
try {
const stats = await handle.stat();
if (
!stats.isFile() ||
stats.size > REPLAY_MAX_FIXTURE_BYTES ||
stats.dev !== requestedStats.dev ||
stats.ino !== requestedStats.ino
) {
throw new ReplayFixtureRepositoryError(
'fixture-not-allowed'
);
}
const text = await readBoundedFixture(handle);
return parseReplayFixtureText(text);
} finally {
await handle.close();
}
} catch (error) {
if (error instanceof ReplayFixtureRepositoryError) {
throw error;
}
if (error instanceof ReplaySchemaError) {
throw new ReplayFixtureRepositoryError('fixture-invalid');
}
throw new ReplayFixtureRepositoryError('fixture-not-allowed');
}
}
}
export function createRepositoryReplayFixtureRepository(): ReplayFixtureRepository {
return new RepositoryReplayFixtureRepository(
path.resolve(
process.cwd(),
'apps/stalker-mock-server/fixtures/replay'
)
);
}
export function createReplayControlPlaneCapability(): string {
return randomBytes(32).toString('base64url');
}
function capabilityMatches(actual: string | undefined, expected: string): boolean {
if (actual === undefined) {
return false;
}
const actualBytes = Buffer.from(actual);
const expectedBytes = Buffer.from(expected);
return (
actualBytes.byteLength === expectedBytes.byteLength &&
timingSafeEqual(actualBytes, expectedBytes)
);
}
function sendJson(
response: ServerResponse,
status: number,
value: unknown
): void {
const body = Buffer.from(JSON.stringify(value));
response.writeHead(status, {
'content-type': 'application/json',
'content-length': String(body.byteLength),
'cache-control': 'no-store',
connection: 'close',
});
response.end(body);
}
function sendControlError(
response: ServerResponse,
error: ReplayControlError
): void {
sendJson(response, error.status, { error: { code: error.code } });
}
function requestPath(request: IncomingMessage): string {
if (request.url === undefined) {
throw new ReplayControlError(
CONTROL_ERROR_CODE.ENDPOINT_NOT_FOUND,
404
);
}
const url = new URL(request.url, 'http://control.invalid');
if (url.search.length > 0) {
throw new ReplayControlError(
CONTROL_ERROR_CODE.ENDPOINT_NOT_FOUND,
404
);
}
return url.pathname;
}
function contentType(request: IncomingMessage): string {
const header = request.headers['content-type'];
const value = Array.isArray(header) ? header[0] : header;
return (value ?? '').split(';', 1)[0]?.trim().toLowerCase() ?? '';
}
function readControlBody(
request: IncomingMessage,
timeoutMs: number
): Promise<Buffer> {
const declaredLength = request.headers['content-length'];
if (
declaredLength !== undefined &&
(!/^\d+$/.test(declaredLength) ||
Number(declaredLength) > REPLAY_CONTROL_MAX_BODY_BYTES)
) {
request.resume();
throw new ReplayControlError(
CONTROL_ERROR_CODE.BODY_TOO_LARGE,
413
);
}
return new Promise<Buffer>((resolve, reject) => {
let total = 0;
const chunks: Buffer[] = [];
let settled = false;
const cleanup = () => {
clearTimeout(timer);
request.off('data', onData);
request.off('end', onEnd);
request.off('aborted', onAborted);
request.off('error', onError);
};
const rejectOnce = (error: ReplayControlError) => {
if (settled) {
return;
}
settled = true;
cleanup();
request.resume();
reject(error);
};
const onData = (chunk: Buffer | string) => {
const buffer = Buffer.isBuffer(chunk)
? chunk
: Buffer.from(chunk);
total += buffer.byteLength;
if (total > REPLAY_CONTROL_MAX_BODY_BYTES) {
rejectOnce(
new ReplayControlError(
CONTROL_ERROR_CODE.BODY_TOO_LARGE,
413
)
);
return;
}
chunks.push(buffer);
};
const onEnd = () => {
if (settled) {
return;
}
settled = true;
cleanup();
resolve(Buffer.concat(chunks));
};
const onAborted = () =>
rejectOnce(
new ReplayControlError(
CONTROL_ERROR_CODE.INVALID_REQUEST,
400
)
);
const onError = () => onAborted();
const timer = setTimeout(
() =>
rejectOnce(
new ReplayControlError(
CONTROL_ERROR_CODE.BODY_TIMEOUT,
408
)
),
timeoutMs
);
timer.unref();
request.on('data', onData);
request.once('end', onEnd);
request.once('aborted', onAborted);
request.once('error', onError);
});
}
function isRecord(value: unknown): value is Record<string, unknown> {
return (
value !== null &&
typeof value === 'object' &&
!Array.isArray(value) &&
(Object.getPrototypeOf(value) === Object.prototype ||
Object.getPrototypeOf(value) === null)
);
}
async function parseJsonBody(
request: IncomingMessage,
requestBodyTimeoutMs: number
): Promise<Record<string, unknown>> {
if (contentType(request) !== 'application/json') {
throw new ReplayControlError(
CONTROL_ERROR_CODE.JSON_REQUIRED,
415
);
}
const rawBody = await readControlBody(
request,
requestBodyTimeoutMs
);
try {
const value: unknown = JSON.parse(rawBody.toString('utf8'));
if (!isRecord(value)) {
throw new ReplayControlError(
CONTROL_ERROR_CODE.INVALID_REQUEST,
400
);
}
return value;
} catch (error) {
if (error instanceof ReplayControlError) {
throw error;
}
throw new ReplayControlError(CONTROL_ERROR_CODE.INVALID_JSON, 400);
}
}
function hasExactKeys(
value: Record<string, unknown>,
keys: readonly string[]
): boolean {
return (
Object.keys(value).length === keys.length &&
keys.every((key) => Object.prototype.hasOwnProperty.call(value, key))
);
}
function readFixtureId(body: Record<string, unknown>): string {
if (
!hasExactKeys(body, ['fixtureId']) ||
typeof body['fixtureId'] !== 'string' ||
!SAFE_FIXTURE_ID.test(body['fixtureId'])
) {
throw new ReplayControlError(
CONTROL_ERROR_CODE.FIXTURE_NOT_ALLOWED,
404
);
}
return body['fixtureId'];
}
function readRunId(body: Record<string, unknown>): string {
if (
!hasExactKeys(body, ['runId']) ||
typeof body['runId'] !== 'string' ||
!SAFE_RUN_ID.test(body['runId'])
) {
throw new ReplayControlError(
CONTROL_ERROR_CODE.RUN_NOT_FOUND,
404
);
}
return body['runId'];
}
async function closeServer(server: Server): Promise<void> {
if (!server.listening) {
return;
}
await new Promise<void>((resolve, reject) => {
server.close((error) =>
error === undefined ? resolve() : reject(error)
);
});
}
async function listen(server: Server): Promise<number> {
await new Promise<void>((resolve, reject) => {
const onError = (error: Error) => {
server.off('listening', onListening);
reject(error);
};
const onListening = () => {
server.off('error', onError);
resolve();
};
server.once('error', onError);
server.once('listening', onListening);
server.listen(0, LOOPBACK_ADDRESS);
});
const address = server.address();
if (address === null || typeof address === 'string') {
throw new Error('Replay control listener has no TCP address.');
}
return (address as AddressInfo).port;
}
export interface StartReplayControlPlaneOptions {
capability?: string;
repository?: ReplayFixtureRepository;
requestBodyTimeoutMs?: number;
}
export interface ReplayControlPlane {
readonly url: string;
close(): Promise<void>;
}
export async function startReplayControlPlane(
options: StartReplayControlPlaneOptions = {}
): Promise<ReplayControlPlane> {
const capability =
options.capability ?? createReplayControlPlaneCapability();
const repository =
options.repository ?? createRepositoryReplayFixtureRepository();
const requestBodyTimeoutMs =
options.requestBodyTimeoutMs ?? DEFAULT_REQUEST_BODY_TIMEOUT_MS;
if (
!Number.isInteger(requestBodyTimeoutMs) ||
requestBodyTimeoutMs <= 0
) {
throw new Error('Replay control body timeout must be positive.');
}
const runs = new Map<string, ReplayServerRun>();
let expectedHost = '';
const server = http.createServer(async (request, response) => {
try {
if (request.headers.host !== expectedHost) {
throw new ReplayControlError(
CONTROL_ERROR_CODE.INVALID_HOST,
400
);
}
const capabilityHeader =
request.headers[REPLAY_CONTROL_CAPABILITY_HEADER];
if (
Array.isArray(capabilityHeader) ||
!capabilityMatches(capabilityHeader, capability)
) {
throw new ReplayControlError(
CONTROL_ERROR_CODE.INVALID_CAPABILITY,
403
);
}
const actionPath = requestPath(request);
if (
actionPath !== CREATE_PATH &&
actionPath !== FINALIZE_PATH &&
actionPath !== DISPOSE_PATH
) {
throw new ReplayControlError(
CONTROL_ERROR_CODE.ENDPOINT_NOT_FOUND,
404
);
}
if (request.method !== 'POST') {
throw new ReplayControlError(
CONTROL_ERROR_CODE.METHOD_NOT_ALLOWED,
405
);
}
const body = await parseJsonBody(
request,
requestBodyTimeoutMs
);
if (actionPath === CREATE_PATH) {
const fixtureId = readFixtureId(body);
let fixture: ReplayFixtureV1;
try {
fixture = await repository.loadFixture(fixtureId);
} catch (error) {
if (
error instanceof ReplayFixtureRepositoryError &&
error.code === 'fixture-not-allowed'
) {
throw new ReplayControlError(
CONTROL_ERROR_CODE.FIXTURE_NOT_ALLOWED,
404
);
}
if (
error instanceof ReplayFixtureRepositoryError &&
error.code === 'fixture-invalid'
) {
throw new ReplayControlError(
CONTROL_ERROR_CODE.FIXTURE_INVALID,
422
);
}
throw error;
}
const run = await createReplayServerRun(fixture);
if (runs.has(run.runId)) {
await run.dispose();
throw new ReplayControlError(
CONTROL_ERROR_CODE.INTERNAL,
500
);
}
runs.set(run.runId, run);
sendJson(response, 201, {
runId: run.runId,
entryUrl: run.entryUrl,
origins: run.originUrls,
inputs: run.generatedInputs,
});
return;
}
const runId = readRunId(body);
const run = runs.get(runId);
if (run === undefined) {
throw new ReplayControlError(
CONTROL_ERROR_CODE.RUN_NOT_FOUND,
404
);
}
if (actionPath === FINALIZE_PATH) {
sendJson(response, 200, run.finalize());
return;
}
await run.dispose();
runs.delete(runId);
sendJson(response, 200, { disposed: true });
} catch (error) {
if (!request.complete) {
request.resume();
}
if (error instanceof ReplayControlError) {
sendControlError(response, error);
} else {
sendControlError(
response,
new ReplayControlError(CONTROL_ERROR_CODE.INTERNAL, 500)
);
}
}
});
const port = await listen(server);
expectedHost = `${LOOPBACK_ADDRESS}:${port}`;
const url = `http://${expectedHost}`;
let closePromise: Promise<void> | undefined;
return {
url,
close: () => {
closePromise ??= (async () => {
await Promise.all(
[...runs.values()].map((run) => run.dispose())
);
runs.clear();
await closeServer(server);
})();
return closePromise;
},
};
}
@@ -1,17 +1,52 @@
import { ReplaySymbolTable } from './replay-symbols.js';
import {
ReplayOriginUrlNode,
ReplayRenderedResponse,
ReplayResponseDefinition,
ReplayResponseHeaderValue,
} from './replay.types.js';
export type ReplayOriginUrlMap = Readonly<Record<string, string>>;
export class ReplayResponseError extends Error {
constructor(readonly code: 'unknown-origin-url') {
super(`Replay response rejected: ${code}.`);
this.name = 'ReplayResponseError';
}
}
function isOriginUrlNode(
value: ReplayResponseHeaderValue
): value is ReplayOriginUrlNode {
return typeof value !== 'string' && value.kind === 'origin-url';
}
function renderHeaderValue(
value: ReplayResponseHeaderValue,
symbols: ReplaySymbolTable,
originUrls: ReplayOriginUrlMap
): string {
if (!isOriginUrlNode(value)) {
return symbols.resolveString(value);
}
const originUrl = originUrls[value.origin];
if (originUrl === undefined) {
throw new ReplayResponseError('unknown-origin-url');
}
return `${originUrl}${value.path}`;
}
export function renderReplayResponse(
response: ReplayResponseDefinition,
symbols: ReplaySymbolTable
symbols: ReplaySymbolTable,
originUrls: ReplayOriginUrlMap = {}
): ReplayRenderedResponse {
const headers = Object.fromEntries(
Object.entries(response.headers).map(([name, values]) => [
name,
values.map((value) => symbols.resolveString(value)),
values.map((value) =>
renderHeaderValue(value, symbols, originUrls)
),
])
);
@@ -63,6 +63,7 @@ export type ReplayFinalizeResult =
export interface CreateReplayRunOptions {
runId?: string;
now?: () => number;
originUrls?: Readonly<Record<string, string>>;
}
interface PendingBarrierResponse {
@@ -84,6 +85,7 @@ export class ReplayRun {
private readonly symbols: ReplaySymbolTable;
private readonly now: () => number;
private readonly originUrls: Readonly<Record<string, string>>;
private readonly createdAt: number;
private lastActivityAt: number;
private phaseIndex = 0;
@@ -106,6 +108,7 @@ export class ReplayRun {
) {
this.runId = options.runId ?? createReplayRunId();
this.now = options.now ?? Date.now;
this.originUrls = options.originUrls ?? {};
this.createdAt = this.now();
this.lastActivityAt = this.createdAt;
this.state = fixture.initialState;
@@ -175,7 +178,8 @@ export class ReplayRun {
const response = renderReplayResponse(
match.expectation.response,
this.symbols
this.symbols,
this.originUrls
);
const phase = this.fixture.phases[this.phaseIndex];
const barrier = phase?.barrier;
@@ -204,6 +204,79 @@ describe('replay fixture schema v1', () => {
});
});
it('accepts repeated exact query values and a declared typed named-origin redirect', () => {
const value = fixture();
const firstExpectation = (
(value['phases'] as Record<string, unknown>[])[0]![
'expectations'
] as Record<string, unknown>[]
)[0]!;
const request = firstExpectation['request'] as Record<string, unknown>;
request['query'] = {
exact: { repeated: ['one', 'two'] },
present: [],
absent: [],
};
const response = firstExpectation['response'] as Record<
string,
unknown
>;
response['headers'] = {
location: [
{
kind: 'origin-url',
origin: 'auth',
path: '/login',
},
],
};
const parsed = parseReplayFixture(value);
expect(
parsed.phases[0]?.expectations[0]?.request.query.exact['repeated']
).toEqual(['one', 'two']);
expect(
parsed.phases[0]?.expectations[0]?.response.headers['location']
).toEqual([
{
kind: 'origin-url',
origin: 'auth',
path: '/login',
},
]);
});
it('rejects unknown or free-form absolute redirect origins', () => {
const unknownOrigin = fixture();
const unknownResponse = (
(unknownOrigin['phases'] as Record<string, unknown>[])[0]![
'expectations'
] as Record<string, unknown>[]
)[0]!['response'] as Record<string, unknown>;
unknownResponse['headers'] = {
location: [
{
kind: 'origin-url',
origin: 'undeclared',
path: '/login',
},
],
};
expectSchemaCode(unknownOrigin, 'unknown-origin');
const literalOrigin = fixture();
const literalResponse = (
(literalOrigin['phases'] as Record<string, unknown>[])[0]![
'expectations'
] as Record<string, unknown>[]
)[0]!['response'] as Record<string, unknown>;
literalResponse['headers'] = {
location: ['http://127.0.0.1:1234/login'],
};
expectSchemaCode(literalOrigin, 'invalid-redirect-location');
});
it.each([
['unknown response', { kind: 'binary', value: 'x' }],
['empty response with data', { kind: 'empty', value: null }],
@@ -706,7 +706,8 @@ function validateRequestBody(
function validateResponseHeaders(
value: unknown,
path: string,
context: ValidationContext
context: ValidationContext,
origins: Set<string>
): boolean {
if (!isRecord(value)) {
issue(
@@ -735,10 +736,64 @@ function validateResponseHeaders(
continue;
}
values.forEach((headerValue, index) => {
const headerPath = `${path}.${name}[${index}]`;
if (
isRecord(headerValue) &&
headerValue['kind'] === 'origin-url'
) {
if (
name !== 'location' ||
!hasExactKeys(headerValue, ['kind', 'origin', 'path'])
) {
issue(
context,
'invalid-redirect-location',
headerPath,
'Named-origin URL nodes are valid only as Location values.'
);
valid = false;
return;
}
if (
typeof headerValue['origin'] !== 'string' ||
!origins.has(headerValue['origin'])
) {
issue(
context,
'unknown-origin',
`${headerPath}.origin`,
'Redirect origin must reference a declared named origin.'
);
valid = false;
}
valid =
validatePath(
headerValue['path'],
`${headerPath}.path`,
context
) && valid;
return;
}
if (name === 'location') {
if (
typeof headerValue !== 'string' ||
/^[A-Za-z][A-Za-z0-9+.-]*:/.test(headerValue) ||
headerValue.startsWith('//')
) {
issue(
context,
'invalid-redirect-location',
headerPath,
'Cross-origin redirects require a typed named-origin URL node.'
);
valid = false;
return;
}
}
valid =
validateTemplateString(
headerValue,
`${path}.${name}[${index}]`,
headerPath,
context,
true
) && valid;
@@ -843,7 +898,8 @@ function validateResponseBody(
function validateResponse(
value: unknown,
path: string,
context: ValidationContext
context: ValidationContext,
origins: Set<string>
): boolean {
if (
!isRecord(value) ||
@@ -872,7 +928,12 @@ function validateResponse(
valid = false;
}
return (
validateResponseHeaders(value['headers'], `${path}.headers`, context) &&
validateResponseHeaders(
value['headers'],
`${path}.headers`,
context,
origins
) &&
validateResponseBody(value['body'], `${path}.body`, context) &&
valid
);
@@ -967,7 +1028,8 @@ function validateExpectation(
validateFieldMatchers(
value['request']['query'],
`${path}.request.query`,
context
context,
{ allowArrays: true }
);
validateFieldMatchers(
value['request']['headers'],
@@ -987,7 +1049,12 @@ function validateExpectation(
);
}
validateResponse(value['response'], `${path}.response`, context);
validateResponse(
value['response'],
`${path}.response`,
context,
origins
);
let cardinalityMinimum = 0;
if (
@@ -0,0 +1,742 @@
/* eslint-disable max-lines, @typescript-eslint/no-non-null-assertion -- Network contract scenarios stay explicit so every synthetic origin and wire value is auditable. */
import http from 'node:http';
import {
createReplayServerRun,
ReplayServerRun,
} from './replay-server.js';
import { REPLAY_MAX_REQUEST_BODY_BYTES } from './replay.constants.js';
import {
ReplayExpectation,
ReplayFixtureV1,
ReplayRequestBodyMatcher,
ReplayResponseDefinition,
} from './replay.types.js';
const EMPTY_FIELDS = {
exact: {},
present: [],
absent: [],
};
function expectation(
id: string,
options: {
origin?: string;
method?: ReplayExpectation['method'];
path?: string;
query?: ReplayExpectation['request']['query'];
headers?: ReplayExpectation['request']['headers'];
cookies?: ReplayExpectation['request']['cookies'];
body?: ReplayRequestBodyMatcher;
response?: ReplayResponseDefinition;
} = {}
): ReplayExpectation {
return {
id,
operation: id,
origin: options.origin ?? 'portal',
method: options.method ?? 'GET',
path: options.path ?? `/${id}`,
request: {
query: options.query ?? EMPTY_FIELDS,
headers: options.headers ?? EMPTY_FIELDS,
cookies: options.cookies ?? {
...EMPTY_FIELDS,
attributes: {},
},
body: options.body ?? { kind: 'absent' },
},
response: options.response ?? {
status: 200,
headers: {
'content-type': ['application/json'],
},
body: { kind: 'json', value: { js: true } },
},
cardinality: { min: 1, max: 1 },
};
}
function fixture(
phases: ReplayFixtureV1['phases'],
options: {
origins?: ReplayFixtureV1['origins'];
entry?: ReplayFixtureV1['entry'];
expectedEndpoint?: ReplayFixtureV1['expectedEndpoint'];
symbols?: ReplayFixtureV1['symbols'];
} = {}
): ReplayFixtureV1 {
return {
schemaVersion: 1,
scenarioId: 'network-contract',
origins: options.origins ?? { portal: {} },
entry: options.entry ?? { origin: 'portal', path: '/entry' },
expectedEndpoint:
options.expectedEndpoint ?? {
origin: 'portal',
path: phases.at(-1)!.expectations.at(-1)!.path,
},
initialState: phases[0]!.state,
terminalState: phases.at(-1)!.nextState,
failOnUnexpectedRequest: true,
symbols: options.symbols ?? [],
phases,
};
}
function orderedPhases(
expectations: ReplayExpectation[]
): ReplayFixtureV1['phases'] {
return expectations.map((item, index) => ({
name: `phase-${index}`,
state: `state-${index}`,
nextState: `state-${index + 1}`,
mode: 'ordered',
expectations: [item],
}));
}
interface RawResponse {
status: number;
headers: string[];
body: string;
}
function rawRequest(
url: string,
options: {
method?: string;
headers?: http.OutgoingHttpHeaders;
body?: string;
} = {}
): Promise<RawResponse> {
return new Promise((resolve, reject) => {
const request = http.request(
url,
{
method: options.method ?? 'GET',
headers: options.headers,
},
(response) => {
const chunks: Buffer[] = [];
response.on('data', (chunk: Buffer) => chunks.push(chunk));
response.on('end', () =>
resolve({
status: response.statusCode ?? 0,
headers: response.rawHeaders,
body: Buffer.concat(chunks).toString('utf8'),
})
);
}
);
request.on('error', reject);
if (options.body !== undefined) {
request.end(options.body);
} else {
request.end();
}
});
}
function streamingRequest(
url: string,
headers: http.OutgoingHttpHeaders
): {
request: http.ClientRequest;
response: Promise<RawResponse>;
} {
let responseStarted = false;
let streamingResult: http.ClientRequest | undefined;
const response = new Promise<RawResponse>((resolve, reject) => {
const request = http.request(
url,
{
method: 'POST',
headers,
},
(incoming) => {
responseStarted = true;
const chunks: Buffer[] = [];
incoming.on('data', (chunk: Buffer) => chunks.push(chunk));
incoming.on('end', () =>
resolve({
status: incoming.statusCode ?? 0,
headers: incoming.rawHeaders,
body: Buffer.concat(chunks).toString('utf8'),
})
);
}
);
request.on('error', (error) => {
if (!responseStarted) {
reject(error);
}
});
streamingResult = request;
});
if (streamingResult === undefined) {
throw new Error('Streaming request was not initialized.');
}
return { request: streamingResult, response };
}
function headerValues(rawHeaders: string[], name: string): string[] {
const values: string[] = [];
for (let index = 0; index < rawHeaders.length; index += 2) {
if (rawHeaders[index]?.toLowerCase() === name.toLowerCase()) {
values.push(rawHeaders[index + 1]!);
}
}
return values;
}
describe('ReplayServerRun listeners and routing', () => {
const activeRuns: ReplayServerRun[] = [];
afterEach(async () => {
await Promise.all(activeRuns.splice(0).map((run) => run.dispose()));
});
it('binds one distinct ephemeral 127.0.0.1 listener per named origin and strips the run route before matching', async () => {
const first = expectation('landing', {
origin: 'portal',
path: '/c/',
});
const second = expectation('auth', {
origin: 'auth',
path: '/server/load.php',
});
const run = await createReplayServerRun(
fixture(orderedPhases([first, second]), {
origins: { portal: {}, auth: {} },
entry: { origin: 'portal', path: '/c/' },
expectedEndpoint: {
origin: 'auth',
path: '/server/load.php',
},
}),
{ runId: 'run-listeners' }
);
activeRuns.push(run);
const portal = new URL(run.originUrls['portal']!);
const auth = new URL(run.originUrls['auth']!);
expect(portal.hostname).toBe('127.0.0.1');
expect(auth.hostname).toBe('127.0.0.1');
expect(Number(portal.port)).toBeGreaterThan(0);
expect(Number(auth.port)).toBeGreaterThan(0);
expect(portal.port).not.toBe(auth.port);
expect(portal.pathname).toContain('/run-listeners');
expect(run.entryUrl).toBe(`${run.originUrls['portal']}/c/`);
await expect(fetch(run.entryUrl)).resolves.toMatchObject({
status: 200,
});
await expect(
fetch(`${run.originUrls['auth']}/server/load.php`)
).resolves.toMatchObject({ status: 200 });
expect(run.finalize().ok).toBe(true);
});
it('keeps simultaneous runs and their generated inputs isolated', async () => {
const scenario = fixture(
orderedPhases([
expectation('entry', {
path: '/entry',
}),
]),
{
symbols: [
{ kind: 'generate', symbol: 'mac', valueKind: 'mac' },
],
}
);
const first = await createReplayServerRun(scenario, {
runId: 'run-one',
});
const second = await createReplayServerRun(scenario, {
runId: 'run-two',
});
activeRuns.push(first, second);
expect(first.entryUrl).not.toBe(second.entryUrl);
expect(first.generatedInputs['mac']).not.toBe(
second.generatedInputs['mac']
);
await fetch(first.entryUrl);
expect(first.finalize().ok).toBe(true);
expect(second.finalize()).toMatchObject({
ok: false,
issues: expect.arrayContaining([
expect.objectContaining({ code: 'unmet-cardinality' }),
]),
});
});
});
describe('ReplayServerRun redirects and wire preservation', () => {
const activeRuns: ReplayServerRun[] = [];
afterEach(async () => {
await Promise.all(activeRuns.splice(0).map((run) => run.dispose()));
});
it('keeps an absolute-path redirect on the same synthetic run origin', async () => {
const redirect = expectation('redirect', {
path: '/entry',
response: {
status: 302,
headers: { location: ['/landing'] },
body: { kind: 'empty' },
},
});
const landing = expectation('landing', { path: '/landing' });
const run = await createReplayServerRun(
fixture(orderedPhases([redirect, landing]), {
entry: { origin: 'portal', path: '/entry' },
expectedEndpoint: { origin: 'portal', path: '/landing' },
}),
{ runId: 'run-relative' }
);
activeRuns.push(run);
const first = await fetch(run.entryUrl, { redirect: 'manual' });
const location = first.headers.get('location');
expect(location).toBe(`${run.originUrls['portal']}/landing`);
expect(new URL(location!).origin).toBe(
new URL(run.entryUrl).origin
);
await fetch(location!);
expect(run.finalize().ok).toBe(true);
});
it('renders a typed named-origin redirect as a real cross-origin absolute URL', async () => {
const redirect = expectation('redirect', {
path: '/entry',
response: {
status: 302,
headers: {
location: [
{
kind: 'origin-url',
origin: 'auth',
path: '/login',
},
],
},
body: { kind: 'empty' },
},
});
const login = expectation('login', {
origin: 'auth',
path: '/login',
});
const run = await createReplayServerRun(
fixture(orderedPhases([redirect, login]), {
origins: { portal: {}, auth: {} },
entry: { origin: 'portal', path: '/entry' },
expectedEndpoint: { origin: 'auth', path: '/login' },
}),
{ runId: 'run-cross-origin' }
);
activeRuns.push(run);
const first = await fetch(run.entryUrl, { redirect: 'manual' });
const location = first.headers.get('location');
expect(location).toBe(`${run.originUrls['auth']}/login`);
expect(new URL(location!).origin).not.toBe(
new URL(run.entryUrl).origin
);
await fetch(location!);
expect(run.finalize().ok).toBe(true);
});
it('parses repeated query/header values, cookies and form bodies and preserves repeated response headers', async () => {
const request = expectation('request', {
method: 'POST',
path: '/portal.php',
query: {
exact: { tag: ['one', 'two'] },
present: [],
absent: [],
},
headers: {
exact: { 'x-repeated': ['alpha', 'beta'] },
present: [],
absent: [],
},
cookies: {
exact: { session: 'synthetic-cookie' },
present: ['theme'],
absent: ['forbidden'],
attributes: {},
},
body: {
kind: 'form',
exact: { username: 'synthetic-user' },
present: ['password'],
absent: ['token'],
},
response: {
status: 200,
headers: {
'content-type': ['text/plain'],
'set-cookie': [
'session=rotated; Path=/; HttpOnly',
'theme=dark; Path=/',
],
'x-repeated': ['one', 'two'],
},
body: { kind: 'text', value: 'accepted' },
},
});
const run = await createReplayServerRun(
fixture(orderedPhases([request]), {
entry: { origin: 'portal', path: '/portal.php' },
expectedEndpoint: {
origin: 'portal',
path: '/portal.php',
},
}),
{ runId: 'run-wire' }
);
activeRuns.push(run);
const response = await rawRequest(
`${run.entryUrl}?tag=one&tag=two`,
{
method: 'POST',
headers: {
'content-type': 'application/x-www-form-urlencoded',
'x-repeated': ['alpha', 'beta'],
cookie: 'session=synthetic-cookie; theme=dark',
},
body: 'username=synthetic-user&password=synthetic-password',
}
);
expect(response).toMatchObject({ status: 200, body: 'accepted' });
expect(headerValues(response.headers, 'set-cookie')).toEqual([
'session=rotated; Path=/; HttpOnly',
'theme=dark; Path=/',
]);
expect(headerValues(response.headers, 'x-repeated')).toEqual([
'one',
'two',
]);
expect(run.finalize().ok).toBe(true);
});
it.each([
[
'json',
'application/json',
JSON.stringify({ username: 'synthetic-user', count: 2 }),
{
kind: 'json',
exact: { username: 'synthetic-user', count: 2 },
present: [],
absent: [],
} satisfies ReplayRequestBodyMatcher,
],
[
'text',
'text/plain',
'synthetic text',
{
kind: 'text',
exact: 'synthetic text',
} satisfies ReplayRequestBodyMatcher,
],
])(
'parses a bounded %s request body',
async (label, contentType, body, matcher) => {
const request = expectation(label, {
method: 'POST',
path: '/request',
body: matcher,
});
const run = await createReplayServerRun(
fixture(orderedPhases([request]), {
entry: { origin: 'portal', path: '/request' },
expectedEndpoint: {
origin: 'portal',
path: '/request',
},
}),
{ runId: `run-body-${label}` }
);
activeRuns.push(run);
const response = await rawRequest(run.entryUrl, {
method: 'POST',
headers: { 'content-type': contentType },
body,
});
expect(response.status).toBe(200);
expect(run.finalize().ok).toBe(true);
}
);
it('rejects JSON nesting beyond the bounded matcher depth with a stable request error', async () => {
const request = expectation('deep-json', {
method: 'POST',
path: '/request',
body: {
kind: 'json',
exact: {},
present: [],
absent: [],
},
});
const run = await createReplayServerRun(
fixture(orderedPhases([request]), {
entry: { origin: 'portal', path: '/request' },
expectedEndpoint: {
origin: 'portal',
path: '/request',
},
}),
{ runId: 'run-deep-json' }
);
activeRuns.push(run);
let nestedJson = 'null';
for (let depth = 0; depth < 66; depth += 1) {
nestedJson = `{"value":${nestedJson}}`;
}
const response = await rawRequest(run.entryUrl, {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: nestedJson,
});
expect(response).toMatchObject({
status: 400,
body: JSON.stringify({
error: { code: 'invalid-request-body' },
}),
});
});
it('requires duplicate cookie names to be matched through the raw repeated-header boundary', async () => {
const request = expectation('duplicate-cookie', {
path: '/request',
headers: {
exact: {
cookie: ['session=first; session=second'],
},
present: [],
absent: [],
},
cookies: {
exact: {},
present: [],
absent: ['session'],
attributes: {},
},
});
const run = await createReplayServerRun(
fixture(orderedPhases([request]), {
entry: { origin: 'portal', path: '/request' },
expectedEndpoint: {
origin: 'portal',
path: '/request',
},
}),
{ runId: 'run-duplicate-cookie' }
);
activeRuns.push(run);
const response = await rawRequest(run.entryUrl, {
headers: { cookie: 'session=first; session=second' },
});
expect(response.status).toBe(200);
expect(run.finalize().ok).toBe(true);
});
it('rejects a request body beyond the network cap without parsing it', async () => {
const request = expectation('oversized', {
method: 'POST',
path: '/request',
body: { kind: 'json', exact: {}, present: [], absent: [] },
});
const run = await createReplayServerRun(
fixture(orderedPhases([request]), {
entry: { origin: 'portal', path: '/request' },
expectedEndpoint: {
origin: 'portal',
path: '/request',
},
}),
{ runId: 'run-oversized' }
);
activeRuns.push(run);
const response = await rawRequest(run.entryUrl, {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: JSON.stringify({ value: 'x'.repeat(1024 * 1024) }),
});
expect(response.status).toBe(413);
expect(JSON.parse(response.body)).toEqual({
error: { code: 'request-body-too-large' },
});
});
it('rejects a chunked body as soon as it crosses the cap without waiting for request end', async () => {
const request = expectation('oversized-stream', {
method: 'POST',
path: '/request',
body: { kind: 'text', exact: 'never-matched' },
});
const run = await createReplayServerRun(
fixture(orderedPhases([request]), {
entry: { origin: 'portal', path: '/request' },
expectedEndpoint: {
origin: 'portal',
path: '/request',
},
}),
{ runId: 'run-oversized-stream' }
);
activeRuns.push(run);
const stream = streamingRequest(run.entryUrl, {
'content-type': 'text/plain',
});
stream.request.write(
Buffer.alloc(REPLAY_MAX_REQUEST_BODY_BYTES + 1, 97)
);
await expect(stream.response).resolves.toMatchObject({
status: 413,
body: JSON.stringify({
error: { code: 'request-body-too-large' },
}),
});
stream.request.destroy();
});
it('times out a chunked body that never ends', async () => {
const request = expectation('slow-stream', {
method: 'POST',
path: '/request',
body: { kind: 'text', exact: 'never-matched' },
});
const run = await createReplayServerRun(
fixture(orderedPhases([request]), {
entry: { origin: 'portal', path: '/request' },
expectedEndpoint: {
origin: 'portal',
path: '/request',
},
}),
{
runId: 'run-slow-stream',
requestBodyTimeoutMs: 25,
}
);
activeRuns.push(run);
const stream = streamingRequest(run.entryUrl, {
'content-type': 'text/plain',
});
stream.request.write('partial');
await expect(stream.response).resolves.toMatchObject({
status: 408,
body: JSON.stringify({
error: { code: 'request-body-timeout' },
}),
});
stream.request.destroy();
});
});
describe('ReplayServerRun lifecycle', () => {
it('disposal rejects held barriers and closes every listener', async () => {
const held = expectation('held', { path: '/held' });
const never = expectation('never', { path: '/never' });
const run = await createReplayServerRun(
fixture(
[
{
name: 'barrier',
state: 'start',
nextState: 'complete',
mode: 'unordered',
barrier: {
name: 'two-requests',
releaseWhenMatched: 2,
},
expectations: [held, never],
},
],
{
entry: { origin: 'portal', path: '/held' },
expectedEndpoint: {
origin: 'portal',
path: '/held',
},
}
),
{ runId: 'run-dispose' }
);
const heldResponse = rawRequest(run.entryUrl, {
headers: { connection: 'close' },
});
for (let attempt = 0; attempt < 100; attempt += 1) {
if (run.getLedger().operationCounts['held'] === 1) {
break;
}
await new Promise<void>((resolve) => setImmediate(resolve));
}
expect(run.getLedger().operationCounts['held']).toBe(1);
await run.dispose();
await expect(heldResponse).resolves.toMatchObject({
status: 410,
body: JSON.stringify({ error: { code: 'run-disposed' } }),
});
await expect(fetch(run.entryUrl)).rejects.toThrow();
await expect(run.dispose()).resolves.toBeUndefined();
});
it('exposes the sanitized nonterminal finalization result without closing before dispose', async () => {
const run = await createReplayServerRun(
fixture(
orderedPhases([
expectation('expected', { path: '/expected' }),
]),
{
entry: { origin: 'portal', path: '/expected' },
}
),
{ runId: 'run-finalize' }
);
expect(run.finalize()).toMatchObject({
ok: false,
issues: expect.arrayContaining([
expect.objectContaining({ code: 'unmet-cardinality' }),
expect.objectContaining({ code: 'nonterminal-state' }),
]),
ledger: {
operationCounts: {},
mismatchCounts: {},
terminalState: null,
},
});
await expect(fetch(run.entryUrl)).resolves.toMatchObject({
status: 409,
});
await run.dispose();
});
});
@@ -0,0 +1,541 @@
/* eslint-disable max-lines -- Request parsing, listener routing, and lifecycle stay together as one auditable network boundary. */
import http, {
IncomingHttpHeaders,
IncomingMessage,
Server,
ServerResponse,
} from 'node:http';
import { AddressInfo } from 'node:net';
import {
REPLAY_HTTP_METHODS,
REPLAY_MAX_REQUEST_BODY_BYTES,
} from './replay.constants.js';
import {
createReplayRun,
CreateReplayRunOptions,
ReplayFinalizeResult,
ReplayLedger,
ReplayRun,
ReplayRunError,
} from './replay-run.js';
import { parseReplayFixture } from './replay-schema.js';
import { createReplayRunId } from './replay-symbols.js';
import {
ReplayFixtureV1,
ReplayHttpMethod,
ReplayJsonValue,
ReplayObservedCookie,
ReplayObservedRequest,
ReplayObservedRequestBody,
ReplayRenderedResponse,
} from './replay.types.js';
const REPLAY_ROUTE_ROOT = '/__iptvnator_stalker_replay__';
const LOOPBACK_ADDRESS = '127.0.0.1';
const DEFAULT_REQUEST_BODY_TIMEOUT_MS = 5000;
const MAX_JSON_DEPTH = 64;
const MAX_JSON_NODES = 100_000;
const NETWORK_ERROR_CODE = {
INTERNAL: 'replay-internal-error',
INVALID_BODY: 'invalid-request-body',
INVALID_METHOD: 'invalid-request-method',
INVALID_ROUTE: 'invalid-replay-route',
REQUEST_BODY_TOO_LARGE: 'request-body-too-large',
REQUEST_BODY_TIMEOUT: 'request-body-timeout',
} as const;
type ReplayNetworkErrorCode =
(typeof NETWORK_ERROR_CODE)[keyof typeof NETWORK_ERROR_CODE];
class ReplayNetworkError extends Error {
constructor(
readonly code: ReplayNetworkErrorCode,
readonly status: number
) {
super(`Replay network request rejected: ${code}.`);
this.name = 'ReplayNetworkError';
}
}
export type ReplayServerRunOptions = Pick<
CreateReplayRunOptions,
'runId' | 'now'
> & {
requestBodyTimeoutMs?: number;
};
export interface ReplayServerRun {
readonly runId: string;
readonly entryUrl: string;
readonly originUrls: Readonly<Record<string, string>>;
readonly generatedInputs: Readonly<Record<string, string>>;
getLedger(): ReplayLedger;
finalize(): ReplayFinalizeResult;
dispose(): Promise<void>;
}
interface ReplayListener {
server: Server;
}
function sendError(
response: ServerResponse,
status: number,
code: string
): void {
const body = Buffer.from(JSON.stringify({ error: { code } }));
response.writeHead(status, {
'content-type': 'application/json',
'content-length': String(body.byteLength),
'cache-control': 'no-store',
connection: 'close',
});
response.end(body);
}
function statusForRunError(error: ReplayRunError): number {
if (error.code === 'run-disposed' || error.code === 'run-expired') {
return 410;
}
return 409;
}
function asHttpMethod(method: string | undefined): ReplayHttpMethod {
if (
method === undefined ||
!REPLAY_HTTP_METHODS.includes(method as ReplayHttpMethod)
) {
throw new ReplayNetworkError(
NETWORK_ERROR_CODE.INVALID_METHOD,
405
);
}
return method as ReplayHttpMethod;
}
function requestPath(
requestUrl: string | undefined,
routePrefix: string
): { path: string; url: URL } {
if (requestUrl === undefined) {
throw new ReplayNetworkError(NETWORK_ERROR_CODE.INVALID_ROUTE, 404);
}
const url = new URL(requestUrl, 'http://replay.invalid');
const path =
url.pathname === routePrefix
? '/'
: url.pathname.startsWith(`${routePrefix}/`)
? url.pathname.slice(routePrefix.length)
: undefined;
if (path === undefined || path.length === 0) {
throw new ReplayNetworkError(NETWORK_ERROR_CODE.INVALID_ROUTE, 404);
}
return { path, url };
}
function repeatedQuery(url: URL): Record<string, string[]> {
const query: Record<string, string[]> = {};
for (const [name, value] of url.searchParams) {
(query[name] ??= []).push(value);
}
return query;
}
function repeatedHeaders(request: IncomingMessage): Record<string, string[]> {
const headers: Record<string, string[]> = {};
for (let index = 0; index < request.rawHeaders.length; index += 2) {
const name = request.rawHeaders[index]?.toLowerCase();
const value = request.rawHeaders[index + 1];
if (name !== undefined && value !== undefined) {
(headers[name] ??= []).push(value);
}
}
return headers;
}
function parseCookies(
headers: IncomingHttpHeaders
): Record<string, ReplayObservedCookie> {
const cookieHeaders = Array.isArray(headers.cookie)
? headers.cookie
: headers.cookie === undefined
? []
: [headers.cookie];
const cookies: Record<string, ReplayObservedCookie> = {};
const duplicateNames = new Set<string>();
for (const header of cookieHeaders) {
for (const pair of header.split(';')) {
const separator = pair.indexOf('=');
if (separator <= 0) {
continue;
}
const name = pair.slice(0, separator).trim();
if (name.length === 0) {
continue;
}
if (duplicateNames.has(name)) {
continue;
}
if (Object.prototype.hasOwnProperty.call(cookies, name)) {
delete cookies[name];
duplicateNames.add(name);
continue;
}
cookies[name] = {
value: pair.slice(separator + 1).trim(),
};
}
}
return cookies;
}
function readBoundedBody(
request: IncomingMessage,
timeoutMs: number
): Promise<Buffer> {
const declaredLength = request.headers['content-length'];
if (
declaredLength !== undefined &&
(!/^\d+$/.test(declaredLength) ||
Number(declaredLength) > REPLAY_MAX_REQUEST_BODY_BYTES)
) {
request.resume();
throw new ReplayNetworkError(
NETWORK_ERROR_CODE.REQUEST_BODY_TOO_LARGE,
413
);
}
return new Promise<Buffer>((resolve, reject) => {
let byteLength = 0;
const chunks: Buffer[] = [];
let settled = false;
const cleanup = () => {
clearTimeout(timer);
request.off('data', onData);
request.off('end', onEnd);
request.off('aborted', onAborted);
request.off('error', onError);
};
const rejectOnce = (error: ReplayNetworkError) => {
if (settled) {
return;
}
settled = true;
cleanup();
request.resume();
reject(error);
};
const onData = (chunk: Buffer | string) => {
const buffer = Buffer.isBuffer(chunk)
? chunk
: Buffer.from(chunk);
byteLength += buffer.byteLength;
if (byteLength > REPLAY_MAX_REQUEST_BODY_BYTES) {
rejectOnce(
new ReplayNetworkError(
NETWORK_ERROR_CODE.REQUEST_BODY_TOO_LARGE,
413
)
);
return;
}
chunks.push(buffer);
};
const onEnd = () => {
if (settled) {
return;
}
settled = true;
cleanup();
resolve(Buffer.concat(chunks));
};
const onAborted = () =>
rejectOnce(
new ReplayNetworkError(
NETWORK_ERROR_CODE.INVALID_BODY,
400
)
);
const onError = () => onAborted();
const timer = setTimeout(
() =>
rejectOnce(
new ReplayNetworkError(
NETWORK_ERROR_CODE.REQUEST_BODY_TIMEOUT,
408
)
),
timeoutMs
);
timer.unref();
request.on('data', onData);
request.once('end', onEnd);
request.once('aborted', onAborted);
request.once('error', onError);
});
}
function contentType(request: IncomingMessage): string {
const header = request.headers['content-type'];
const value = Array.isArray(header) ? header[0] : header;
return (value ?? '').split(';', 1)[0]?.trim().toLowerCase() ?? '';
}
function isReplayJsonValue(value: unknown): value is ReplayJsonValue {
const pending: Array<{ value: unknown; depth: number }> = [
{ value, depth: 0 },
];
let nodes = 0;
while (pending.length > 0) {
const item = pending.pop();
if (item === undefined) {
break;
}
nodes += 1;
if (nodes > MAX_JSON_NODES || item.depth > MAX_JSON_DEPTH) {
return false;
}
if (
item.value === null ||
typeof item.value === 'string' ||
typeof item.value === 'number' ||
typeof item.value === 'boolean'
) {
continue;
}
if (typeof item.value !== 'object') {
return false;
}
const children = Array.isArray(item.value)
? item.value
: Object.values(item.value);
for (const child of children) {
pending.push({ value: child, depth: item.depth + 1 });
}
}
return true;
}
function parseBody(
request: IncomingMessage,
body: Buffer
): ReplayObservedRequestBody {
if (body.byteLength === 0) {
return { kind: 'absent' };
}
const type = contentType(request);
const text = body.toString('utf8');
if (type === 'application/json' || type.endsWith('+json')) {
try {
const value: unknown = JSON.parse(text);
if (!isReplayJsonValue(value)) {
throw new Error('invalid JSON value');
}
return { kind: 'json', value };
} catch {
throw new ReplayNetworkError(
NETWORK_ERROR_CODE.INVALID_BODY,
400
);
}
}
if (type === 'application/x-www-form-urlencoded') {
const value: Record<string, string> = {};
for (const [name, fieldValue] of new URLSearchParams(text)) {
value[name] = fieldValue;
}
return { kind: 'form', value };
}
return { kind: 'text', value: text };
}
async function observeRequest(
request: IncomingMessage,
origin: string,
routePrefix: string,
requestBodyTimeoutMs: number
): Promise<ReplayObservedRequest> {
const { path, url } = requestPath(request.url, routePrefix);
const body = await readBoundedBody(request, requestBodyTimeoutMs);
return {
origin,
method: asHttpMethod(request.method),
path,
query: repeatedQuery(url),
headers: repeatedHeaders(request),
cookies: parseCookies(request.headers),
body: parseBody(request, body),
};
}
function routedHeaders(
response: ReplayRenderedResponse,
currentOriginUrl: string
): Record<string, string[]> {
return Object.fromEntries(
Object.entries(response.headers).map(([name, values]) => [
name,
name === 'location'
? values.map((value) =>
value.startsWith('/')
? `${currentOriginUrl}${value}`
: value
)
: values,
])
);
}
function sendReplayResponse(
response: ServerResponse,
replay: ReplayRenderedResponse,
currentOriginUrl: string
): void {
response.writeHead(
replay.status,
routedHeaders(replay, currentOriginUrl)
);
response.end(replay.body);
}
async function listen(server: Server): Promise<number> {
await new Promise<void>((resolve, reject) => {
const onError = (error: Error) => {
server.off('listening', onListening);
reject(error);
};
const onListening = () => {
server.off('error', onError);
resolve();
};
server.once('error', onError);
server.once('listening', onListening);
server.listen(0, LOOPBACK_ADDRESS);
});
const address = server.address();
if (address === null || typeof address === 'string') {
throw new Error('Replay listener did not expose a TCP address.');
}
return (address as AddressInfo).port;
}
async function closeServer(server: Server): Promise<void> {
if (!server.listening) {
return;
}
await new Promise<void>((resolve, reject) => {
server.close((error) =>
error === undefined ? resolve() : reject(error)
);
});
}
export async function createReplayServerRun(
fixtureInput: ReplayFixtureV1 | unknown,
options: ReplayServerRunOptions = {}
): Promise<ReplayServerRun> {
const fixture = parseReplayFixture(fixtureInput);
const runId = options.runId ?? createReplayRunId();
const requestBodyTimeoutMs =
options.requestBodyTimeoutMs ?? DEFAULT_REQUEST_BODY_TIMEOUT_MS;
if (
!Number.isInteger(requestBodyTimeoutMs) ||
requestBodyTimeoutMs <= 0
) {
throw new Error('Replay request body timeout must be positive.');
}
const routePrefix = `${REPLAY_ROUTE_ROOT}/${encodeURIComponent(
runId
)}`;
const listeners: ReplayListener[] = [];
const originUrls: Record<string, string> = {};
const runtime: { run?: ReplayRun } = {};
try {
for (const origin of Object.keys(fixture.origins)) {
const server = http.createServer(async (request, response) => {
try {
const activeRun = runtime.run;
if (activeRun === undefined) {
sendError(
response,
503,
NETWORK_ERROR_CODE.INTERNAL
);
return;
}
const observed = await observeRequest(
request,
origin,
routePrefix,
requestBodyTimeoutMs
);
const rendered = await activeRun.request(observed);
const currentOriginUrl = originUrls[origin];
if (currentOriginUrl === undefined) {
throw new Error(
'Replay origin listener is not registered.'
);
}
sendReplayResponse(
response,
rendered,
currentOriginUrl
);
} catch (error) {
if (error instanceof ReplayNetworkError) {
sendError(response, error.status, error.code);
} else if (error instanceof ReplayRunError) {
sendError(
response,
statusForRunError(error),
error.code
);
} else {
sendError(
response,
500,
NETWORK_ERROR_CODE.INTERNAL
);
}
}
});
const port = await listen(server);
listeners.push({ server });
originUrls[
origin
] = `http://${LOOPBACK_ADDRESS}:${port}${routePrefix}`;
}
} catch (error) {
await Promise.all(listeners.map(({ server }) => closeServer(server)));
throw error;
}
const run = createReplayRun(fixture, {
runId,
now: options.now,
originUrls,
});
runtime.run = run;
let disposePromise: Promise<void> | undefined;
return {
runId: run.runId,
entryUrl: `${originUrls[fixture.entry.origin]}${fixture.entry.path}`,
originUrls: Object.freeze({ ...originUrls }),
generatedInputs: run.getGeneratedInputs(),
getLedger: () => run.getLedger(),
finalize: () => run.finalize(),
dispose: () => {
disposePromise ??= (async () => {
run.dispose();
await Promise.all(
listeners.map(({ server }) => closeServer(server))
);
})();
return disposePromise;
},
};
}
@@ -4,6 +4,7 @@ export const REPLAY_MAX_FIXTURE_BYTES = 1024 * 1024;
export const REPLAY_MAX_PHASES = 128;
export const REPLAY_MAX_EXPECTATIONS = 512;
export const REPLAY_MAX_REQUESTS = 512;
export const REPLAY_MAX_REQUEST_BODY_BYTES = 1024 * 1024;
export const REPLAY_MAX_GENERATED_BODY_BYTES = 16 * 1024 * 1024;
export const REPLAY_RUN_HARD_LIFETIME_MS = 10 * 60 * 1000;
export const REPLAY_RUN_INACTIVITY_MS = 2 * 60 * 1000;
@@ -34,12 +34,22 @@ export interface ReplayPartsNode {
parts: Array<ReplayLiteralPart | ReplayRefNode>;
}
export interface ReplayOriginUrlNode {
kind: 'origin-url';
origin: string;
path: string;
}
export type ReplayTemplateString =
| string
| ReplayGenerateNode
| ReplayRefNode
| ReplayPartsNode;
export type ReplayResponseHeaderValue =
| ReplayTemplateString
| ReplayOriginUrlNode;
export type ReplayTemplateValue =
| ReplayJsonPrimitive
| ReplayGenerateNode
@@ -80,7 +90,9 @@ export interface ReplayCookieMatchers extends ReplayFieldMatchers<ReplayTemplate
}
export interface ReplayRequestMatcher {
query: ReplayFieldMatchers<ReplayTemplateString>;
query: ReplayFieldMatchers<
ReplayTemplateString | ReplayTemplateString[]
>;
headers: ReplayFieldMatchers<ReplayTemplateString | ReplayTemplateString[]>;
cookies: ReplayCookieMatchers;
body: ReplayRequestBodyMatcher;
@@ -99,7 +111,7 @@ export type ReplayResponseBody =
export interface ReplayResponseDefinition {
status: number;
headers: Record<string, ReplayTemplateString[]>;
headers: Record<string, ReplayResponseHeaderValue[]>;
body: ReplayResponseBody;
}