test(performance): add end-to-end Xtream benchmark harness (#1300)

* docs(performance): plan Xtream benchmark

* feat(xtream-mock-server): add deterministic 100k fixture

* style(xtream-mock-server): apply repository formatting

* fix(xtream-mock-server): harden performance fixture data

* feat(xtream-mock-server): add performance control plane

* docs(performance): correct Xtream capture plan

* fix(xtream-mock-server): harden performance controls

* fix(xtream-mock-server): harden control lifecycle

* feat(performance): add Xtream preload markers

* feat(performance): trace Xtream main phases

* feat(performance): mark Xtream store publications

* feat(performance): trace Xtream database phases

* feat(performance): trace Xtream delete cancellation

* feat(performance): capture Xtream phase attribution

* feat(performance): mark Sources Xtream refresh

* test(performance): define Xtream benchmark evidence contracts

* test(performance): add Xtream benchmark runner

* test(performance): surface failure evidence writes

* test(performance): align database read clock

* test(performance): preserve capture failure contracts
This commit is contained in:
4gray authored and GitHub committed 2026-07-28 08:08:07 +02:00
1 parent 5e725a06f3
commit a2fafcfc08
256 files changed
+43760 -883

No files matched your search

+94 -12
View File
@@ -24,20 +24,99 @@ pnpm run serve:marketing-demo:web
---
## Local Performance Mode
The benchmark control plane is opt-in and binds only to an explicit loopback
address. Use a dedicated port; do not reuse the normal E2E server on `3211`.
```bash
HOST=127.0.0.1 \
PORT=3221 \
IPTVNATOR_XTREAM_MOCK_CONTROL=1 \
IPTVNATOR_XTREAM_MOCK_CONTROL_TOKEN=local-benchmark-token \
pnpm nx run xtream-mock-server:serve
```
When enabled, every `/__control/*` request requires the exact
`x-iptvnator-performance-token` header, including `OPTIONS` preflight
requests. The server refuses non-loopback binds and an empty token. The control
routes do not exist when the flag is absent or is any value other than `1`.
The Nx serve targets preserve an explicit shell `PORT`; when it is omitted, the
server parser still defaults to `3211`.
Prepare the fixed synthetic fixture before starting a capture:
```bash
curl -X POST http://127.0.0.1:3221/__control/prepare \
-H 'content-type: application/json' \
-H 'x-iptvnator-performance-token: local-benchmark-token' \
-d '{"scenario":"performance-100k"}'
```
The response contains only `epoch`, `scenario`, `seed`, `counts`, `bytes`, and
`catalogSha256`. `counts.categories` is `60/20/20` live/VOD/series (100 total)
and `counts.items` is `60000/20000/20000` (100,000 total). `bytes` and the
SHA-256 cover fixed-order UTF-8 JSON containing, in order, the live, VOD, and
series category arrays followed by the live, VOD, and series catalogs. The hash
input contains no credentials or server origin.
Control endpoints:
| Method | Path | Purpose |
| ------ | --------------------------------- | ---------------------------------------------------------------------- |
| POST | `/__control/prepare` | Materialize `performance-100k` and return its safe manifest |
| POST | `/__control/reset` | Reset `observations` or `all` state |
| POST | `/__control/barriers` | Add a one-shot, abort-aware request barrier |
| POST | `/__control/barriers/:id/release` | Release a request that reached a barrier |
| POST | `/__control/delays` | Add a one-shot delay of `0..5000` ms for control/smoke tests only |
| GET | `/__control/state` | Read bounded rules, held IDs, occurrences, lifecycle ledger, and epoch |
Rules match the exact tuple
`(epoch, scenario, transport, action, categoryId, occurrence)`. Occurrences are
counted independently per tuple without `occurrence`, so parallel category
arrival order cannot change which rule matches. Empty Xtream actions are
recorded as `get_account_info`; unknown actions are recorded only as `unknown`.
Numeric category IDs use canonical decimal spelling: signs, whitespace,
leading zeroes, non-decimal notation, and unsafe integers are rejected. The
state never includes the token, credentials, raw URL/query, response payloads,
titles, or catalog arrays.
`reset` with mode `observations` clears rules, occurrences, held requests, and
the ledger while preserving the prepared manifest and epoch. Mode `all` also
clears fixture caches and the prepared manifest, then increments the epoch.
Held clients are settled during either reset. The legacy unauthenticated
`POST /reset` returns `410` in performance mode so it cannot invalidate fixture
caches without also invalidating the prepared control manifest; use the
token-authenticated control reset instead.
JSON bodies are strict and limited to 16 KiB. A state epoch accepts at most 32
rules; the ledger retains 128 entries, and occurrence state has 512 slots for
the complete allowlisted identity domain without eviction or counter restart.
Duplicate IDs/matches, past occurrences, arbitrary scenarios/actions/category
IDs, unknown fields, and invalid delays are rejected.
Performance mode is local-only: `/playlist.m3u` and every performance-fixture
stream/timeshift URL return `410` without redirecting or contacting external
media. Barriers and delays are coordination tools, not timing inputs.
**Formal benchmark captures must start with zero barrier and delay rules.**
---
## Available Scenarios (credential pairs)
| Username | Password | Scenario | Live cats | VOD cats | Series cats | Items/cat | Status |
| ----------- | ----------- | ---------------------- | --------- | -------- | ----------- | --------- | -------- |
| `user1` | `pass1` | default | 8 | 8 | 8 | 40 | active |
| `large` | `large` | large catalog | 20 | 20 | 20 | 200 | active |
| `stress` | `stress` | stress catalog | 16 | 16 | 16 | 120 | active |
| `series` | `series` | series-heavy | 3 | 4 | 15 | 30 | active |
| `minimal` | `minimal` | minimal (edge cases) | 2 | 2 | 2 | 5 | active |
| `epg` | `epg` | EPG fixture | 2 | 1 | 1 | 3 | active |
| `emptyvod` | `emptyvod` | empty VOD metadata | 2 | 2 | 2 | 5 | active |
| `marketing` | `marketing` | fictional release demo | 4 | 4 | 4 | curated | active |
| `expired` | `expired` | expired account | 4 | 4 | 4 | 10 | Expired |
| `inactive` | `inactive` | disabled account | 4 | 4 | 4 | 10 | Disabled |
| Username | Password | Scenario | Live cats | VOD cats | Series cats | Items/cat | Status |
| ------------- | ------------- | ---------------------- | --------- | -------- | ----------- | --------- | -------- |
| `user1` | `pass1` | default | 8 | 8 | 8 | 40 | active |
| `large` | `large` | large catalog | 20 | 20 | 20 | 200 | active |
| `stress` | `stress` | stress catalog | 16 | 16 | 16 | 120 | active |
| `performance` | `performance` | performance-100k | 60 | 20 | 20 | 1,000 | active |
| `series` | `series` | series-heavy | 3 | 4 | 15 | 30 | active |
| `minimal` | `minimal` | minimal (edge cases) | 2 | 2 | 2 | 5 | active |
| `epg` | `epg` | EPG fixture | 2 | 1 | 1 | 3 | active |
| `emptyvod` | `emptyvod` | empty VOD metadata | 2 | 2 | 2 | 5 | active |
| `marketing` | `marketing` | fictional release demo | 4 | 4 | 4 | curated | active |
| `expired` | `expired` | expired account | 4 | 4 | 4 | 10 | Expired |
| `inactive` | `inactive` | disabled account | 4 | 4 | 4 | 10 | Disabled |
Any other credential pair is auto-generated using a hash of `username:password` as the faker seed (6 categories, 30 items each, active account).
@@ -145,6 +224,9 @@ application code.
- **EPG**: Titles and descriptions are base64-encoded (matches real Xtream API)
- **Dedicated EPG fixture**: `epg:epg` returns stable live channels plus deterministic `get_short_epg` and `get_simple_data_table` payloads for timezone-focused tests
- **Release screenshot fixture**: `marketing:marketing` returns fictional live, VOD, and series data with local generated artwork under `apps/xtream-mock-server/public/marketing`
- **Performance fixture**: `performance:performance` returns exactly 100,000
local-only catalog items from index-derived values; it does not use Faker,
`Date.now()`, `Math.random()`, external artwork, or external media URLs
- **Timestamp precedence coverage**: The `epg:epg` scenario intentionally shifts raw `start` / `end` strings away from `start_timestamp` / `stop_timestamp` so UI tests can prove timestamps drive rendering
- **Stream IDs**: Live 10,000+, VOD 20,000+, Series 30,000+
- **Category IDs**: Live 101+, VOD 201+, Series 301+
+13
View File
@@ -0,0 +1,13 @@
export default {
displayName: 'xtream-mock-server',
preset: '../../jest.preset.js',
testEnvironment: 'node',
transform: {
'^.+\\.[tj]s$': [
'ts-jest',
{ tsconfig: '<rootDir>/tsconfig.spec.json' },
],
},
moduleFileExtensions: ['ts', 'js', 'html'],
coverageDirectory: '../../coverage/apps/xtream-mock-server',
};
+8 -2
View File
@@ -13,7 +13,6 @@
"cwd": "{workspaceRoot}",
"color": true,
"env": {
"PORT": "3211",
"NODE_ENV": "development"
}
}
@@ -26,13 +25,20 @@
"cwd": "{workspaceRoot}",
"color": true,
"env": {
"PORT": "3211",
"NODE_ENV": "development"
}
}
},
"lint": {
"executor": "@nx/eslint:lint"
},
"test": {
"executor": "@nx/jest:jest",
"outputs": ["{workspaceRoot}/coverage/{projectRoot}"],
"options": {
"jestConfig": "apps/xtream-mock-server/jest.config.ts",
"tsConfig": "apps/xtream-mock-server/tsconfig.spec.json"
}
}
}
}
+114 -51
View File
@@ -8,14 +8,29 @@ import {
RawEpgListing,
RawLiveStream,
} from './generators/live.generator.js';
import { generateVodStreams, generateVodDetails, RawVodStream, RawVodDetails } from './generators/vod.generator.js';
import { generateSeriesItems, generateSeriesInfo, RawSeriesItem, RawSeriesInfo } from './generators/series.generator.js';
import {
generateVodStreams,
generateVodDetails,
RawVodStream,
RawVodDetails,
} from './generators/vod.generator.js';
import {
generateSeriesItems,
generateSeriesInfo,
RawSeriesItem,
RawSeriesInfo,
} from './generators/series.generator.js';
import { RawCategory } from './generators/categories.generator.js';
import {
buildMarketingPortalFixture,
buildMarketingSeriesInfo,
buildMarketingVodDetails,
} from './generators/marketing.generator.js';
import {
buildPerformancePortalFixture,
buildPerformanceSeriesInfo,
buildPerformanceVodDetails,
} from './generators/performance.generator.js';
import { getScenario, ScenarioConfig } from './scenarios.js';
export interface PortalData {
@@ -30,21 +45,31 @@ export interface PortalData {
}
const portalCache = new Map<string, PortalData>();
const vodDetailsCache = new Map<number, RawVodDetails>();
const seriesInfoCache = new Map<number, RawSeriesInfo>();
const vodDetailsCache = new Map<string, RawVodDetails>();
const seriesInfoCache = new Map<string, RawSeriesInfo>();
function generatePortalData(username: string, password: string): PortalData {
const scenario = getScenario(username, password);
if (scenario.performanceFixture === 'catalog-100k') {
return {
scenario,
...buildPerformancePortalFixture(),
};
}
faker.seed(scenario.seed);
const { categoryCount, itemsPerCategory, seasonsPerSeries, episodesPerSeason } = scenario;
const {
categoryCount,
itemsPerCategory,
seasonsPerSeries,
episodesPerSeason,
} = scenario;
if (scenario.marketingFixture) {
const marketingFixture = buildMarketingPortalFixture();
for (const s of marketingFixture.seriesItems) {
if (!seriesInfoCache.has(s.series_id)) {
const cacheKey = detailCacheKey(username, password, s.series_id);
if (!seriesInfoCache.has(cacheKey)) {
seriesInfoCache.set(
s.series_id,
cacheKey,
buildMarketingSeriesInfo(
s,
seasonsPerSeries,
@@ -53,22 +78,18 @@ function generatePortalData(username: string, password: string): PortalData {
);
}
}
return {
scenario,
...marketingFixture,
};
}
let liveCategories = generateCategories('live', categoryCount.live);
const vodCategories = generateCategories('vod', categoryCount.vod);
const seriesCategories = generateCategories('series', categoryCount.series);
const epgListingsByStreamId = new Map<number, RawEpgListing[]>();
let liveStreams = generateLiveStreams(liveCategories, itemsPerCategory);
const vodStreams = generateVodStreams(vodCategories, itemsPerCategory);
const seriesItems = generateSeriesItems(seriesCategories, itemsPerCategory);
if (scenario.epgFixture === 'timezone-focus') {
const timezoneFixture = buildTimezoneFixture();
liveCategories = timezoneFixture.liveCategories;
@@ -77,14 +98,18 @@ function generatePortalData(username: string, password: string): PortalData {
epgListingsByStreamId.set(streamId, listings);
});
}
// Pre-populate series info cache
for (const s of seriesItems) {
if (!seriesInfoCache.has(s.series_id)) {
seriesInfoCache.set(s.series_id, generateSeriesInfo(s, seasonsPerSeries, episodesPerSeason));
if (!scenario.deferSeriesDetails) {
for (const s of seriesItems) {
const cacheKey = detailCacheKey(username, password, s.series_id);
if (!seriesInfoCache.has(cacheKey)) {
seriesInfoCache.set(
cacheKey,
generateSeriesInfo(s, seasonsPerSeries, episodesPerSeason)
);
}
}
}
return {
scenario,
liveCategories,
@@ -105,47 +130,60 @@ export function getPortalData(username: string, password: string): PortalData {
return portalCache.get(key)!;
}
export function getVodDetails(username: string, password: string, vodId: number): RawVodDetails | null {
if (!vodDetailsCache.has(vodId)) {
export function getVodDetails(
username: string,
password: string,
vodId: number
): RawVodDetails | null {
const cacheKey = detailCacheKey(username, password, vodId);
if (!vodDetailsCache.has(cacheKey)) {
const data = getPortalData(username, password);
const stream = data.vodStreams.find(v => v.stream_id === vodId);
const stream = data.vodStreams.find((v) => v.stream_id === vodId);
if (!stream) return null;
if (data.scenario.vodDetailsFixture === 'empty-metadata') {
vodDetailsCache.set(vodId, { info: [] });
return vodDetailsCache.get(vodId) ?? null;
vodDetailsCache.set(cacheKey, { info: [] });
return vodDetailsCache.get(cacheKey) ?? null;
}
vodDetailsCache.set(
vodId,
data.scenario.marketingFixture
? buildMarketingVodDetails(stream)
: generateVodDetails(stream)
cacheKey,
data.scenario.performanceFixture === 'catalog-100k'
? buildPerformanceVodDetails(stream)
: data.scenario.marketingFixture
? buildMarketingVodDetails(stream)
: generateVodDetails(stream)
);
}
return vodDetailsCache.get(vodId) ?? null;
return vodDetailsCache.get(cacheKey) ?? null;
}
export function getSeriesInfo(username: string, password: string, seriesId: number): RawSeriesInfo | null {
if (!seriesInfoCache.has(seriesId)) {
export function getSeriesInfo(
username: string,
password: string,
seriesId: number
): RawSeriesInfo | null {
const cacheKey = detailCacheKey(username, password, seriesId);
if (!seriesInfoCache.has(cacheKey)) {
const data = getPortalData(username, password);
const series = data.seriesItems.find(s => s.series_id === seriesId);
const series = data.seriesItems.find((s) => s.series_id === seriesId);
if (!series) return null;
const scenario = getScenario(username, password);
seriesInfoCache.set(
seriesId,
scenario.marketingFixture
? buildMarketingSeriesInfo(
series,
scenario.seasonsPerSeries,
scenario.episodesPerSeason
)
: generateSeriesInfo(
series,
scenario.seasonsPerSeries,
scenario.episodesPerSeason
)
cacheKey,
data.scenario.performanceFixture === 'catalog-100k'
? buildPerformanceSeriesInfo(series)
: data.scenario.marketingFixture
? buildMarketingSeriesInfo(
series,
data.scenario.seasonsPerSeries,
data.scenario.episodesPerSeason
)
: generateSeriesInfo(
series,
data.scenario.seasonsPerSeries,
data.scenario.episodesPerSeason
)
);
}
return seriesInfoCache.get(seriesId) ?? null;
return seriesInfoCache.get(cacheKey) ?? null;
}
export function getEpgListings(
@@ -183,6 +221,25 @@ export function resetAll(): void {
seriesInfoCache.clear();
}
/** Test-only cache sizes; never exposes credential-bearing keys or fixture data. */
export function getDetailCacheCardinalityForTesting(): Readonly<{
vodDetails: number;
seriesInfo: number;
}> {
return {
vodDetails: vodDetailsCache.size,
seriesInfo: seriesInfoCache.size,
};
}
function detailCacheKey(
username: string,
password: string,
itemId: number
): string {
return `${username}:${password}:${itemId}`;
}
function buildTimezoneFixture(): Pick<
PortalData,
'liveCategories' | 'liveStreams' | 'epgListingsByStreamId'
@@ -258,42 +315,48 @@ function buildTimezoneNewsEpg(streamId: number): RawEpgListing[] {
{
id: `${streamId}-past`,
title: 'Earlier Bulletin',
description: 'Past schedule item used to anchor current-program detection.',
description:
'Past schedule item used to anchor current-program detection.',
startTimestamp: roundedNow - 75 * 60,
stopTimestamp: roundedNow - 15 * 60,
},
{
id: `${streamId}-current`,
title: 'Global Headlines',
description: 'Current program for list-row and detail EPG assertions.',
description:
'Current program for list-row and detail EPG assertions.',
startTimestamp: roundedNow - 15 * 60,
stopTimestamp: roundedNow + 15 * 60,
},
{
id: `${streamId}-next`,
title: 'Market Wrap',
description: 'Immediate next program for short-EPG preview assertions.',
description:
'Immediate next program for short-EPG preview assertions.',
startTimestamp: roundedNow + 15 * 60,
stopTimestamp: roundedNow + 45 * 60,
},
{
id: `${streamId}-later`,
title: 'Overnight Update',
description: 'Later same-day program for selected-channel list coverage.',
description:
'Later same-day program for selected-channel list coverage.',
startTimestamp: roundedNow + 45 * 60,
stopTimestamp: roundedNow + 75 * 60,
},
{
id: `${streamId}-boundary-1`,
title: 'Late Edition',
description: 'Boundary program spanning UTC midnight for timezone edge cases.',
description:
'Boundary program spanning UTC midnight for timezone edge cases.',
startTimestamp: nextUtcMidnight - 30 * 60,
stopTimestamp: nextUtcMidnight + 30 * 60,
},
{
id: `${streamId}-boundary-2`,
title: 'After Midnight',
description: 'Post-midnight program used for next-day navigation checks.',
description:
'Post-midnight program used for next-day navigation checks.',
startTimestamp: nextUtcMidnight + 30 * 60,
stopTimestamp: nextUtcMidnight + 90 * 60,
},
@@ -0,0 +1,388 @@
import { createHash } from 'node:crypto';
import { readFileSync } from 'node:fs';
import type { Request, Response } from 'express';
import * as dataStore from '../data-store.js';
import { dispatchAction } from '../routes/dispatch.js';
import { getScenario } from '../scenarios.js';
jest.mock('@faker-js/faker', () => {
let accessEnabled = false;
const fixedDate = new Date('2020-01-01T00:00:00.000Z');
const fixedText = 'Fixture value';
const deterministicFaker = {
seed: jest.fn(),
company: {
name: () => fixedText,
catchPhrase: () => fixedText,
},
date: {
past: () => fixedDate,
recent: () => fixedDate,
},
location: { country: () => fixedText },
lorem: {
paragraph: () => fixedText,
sentence: () => fixedText,
words: () => fixedText,
},
music: {
genre: () => fixedText,
songName: () => fixedText,
},
number: {
int: ({ min = 0 }: { min?: number }) => min,
},
person: { fullName: () => fixedText },
};
return {
faker: new Proxy(deterministicFaker, {
get(target, property, receiver) {
if (!accessEnabled) {
throw new Error(
`Performance fixture accessed Faker.${String(property)}`
);
}
return Reflect.get(target, property, receiver);
},
}),
setFakerAccessForTesting(enabled: boolean) {
accessEnabled = enabled;
},
};
});
const [USERNAME, PASSWORD] = ['performance', 'performance'];
const LOOPBACK_ORIGIN = 'http://127.0.0.1:3211';
jest.setTimeout(30_000);
type PortalData = ReturnType<typeof dataStore.getPortalData>;
type DispatchResult = {
body: unknown;
jsonCalls: number;
statusCode: number;
};
let initialPortal!: PortalData;
function setFakerAccess(enabled: boolean): void {
const fakerModule = jest.requireMock('@faker-js/faker') as {
setFakerAccessForTesting(value: boolean): void;
};
fakerModule.setFakerAccessForTesting(enabled);
}
function withGenericFaker(run: () => void): void {
setFakerAccess(true);
try {
run();
} finally {
dataStore.resetAll();
setFakerAccess(false);
}
}
function portalFingerprint(data: PortalData): {
byteCount: number;
sha256: string;
} {
const serialized = JSON.stringify({
scenario: data.scenario,
liveCategories: data.liveCategories,
vodCategories: data.vodCategories,
seriesCategories: data.seriesCategories,
liveStreams: data.liveStreams,
vodStreams: data.vodStreams,
seriesItems: data.seriesItems,
epgListingsByStreamId: [...data.epgListingsByStreamId.entries()],
});
return {
byteCount: Buffer.byteLength(serialized),
sha256: createHash('sha256').update(serialized).digest('hex'),
};
}
function dispatch(query: Record<string, string>): DispatchResult {
let body: unknown;
let jsonCalls = 0;
let statusCode = 200;
const response = {
json(value: unknown) {
body = value;
jsonCalls += 1;
return response;
},
status(value: number) {
statusCode = value;
return response;
},
};
dispatchAction(
{ query } as unknown as Request,
response as unknown as Response
);
return { body, jsonCalls, statusCode };
}
function expectLocalOrEmpty(values: string[]): void {
for (const value of values) {
expect(value === '' || value.startsWith(`${LOOPBACK_ORIGIN}/`)).toBe(
true
);
}
}
function expectOneThousandItemsPerCategory(
items: ReadonlyArray<{ category_id: number | string }>
): void {
const counts = new Map<string, number>();
for (const item of items) {
const categoryId = String(item.category_id);
counts.set(categoryId, (counts.get(categoryId) ?? 0) + 1);
}
expect([...counts.values()].every((count) => count === 1_000)).toBe(true);
}
describe('performance Xtream fixture', () => {
beforeAll(() => {
dataStore.resetAll();
initialPortal = dataStore.getPortalData(USERNAME, PASSWORD);
});
afterAll(() => {
dataStore.resetAll();
setFakerAccess(false);
});
it('selects the committed performance scenario', () => {
expect(getScenario(USERNAME, PASSWORD)).toMatchObject({
name: 'performance-100k',
description: 'Deterministic local-only 100k performance catalog',
seed: 91_001,
categoryCount: { live: 60, vod: 20, series: 20 },
itemsPerCategory: 1_000,
seasonsPerSeries: 1,
episodesPerSeason: 1,
accountStatus: 'Active',
expiryDate: '2099-12-31',
performanceFixture: 'catalog-100k',
deferSeriesDetails: true,
});
});
it('contains exactly 100 categories and 100,000 catalog items', () => {
expect([
initialPortal.liveCategories.length,
initialPortal.vodCategories.length,
initialPortal.seriesCategories.length,
]).toEqual([60, 20, 20]);
expect(initialPortal.liveStreams).toHaveLength(60_000);
expect(initialPortal.vodStreams).toHaveLength(20_000);
expect(initialPortal.seriesItems).toHaveLength(20_000);
expectOneThousandItemsPerCategory(initialPortal.liveStreams);
expectOneThousandItemsPerCategory(initialPortal.vodStreams);
expectOneThousandItemsPerCategory(initialPortal.seriesItems);
});
it('uses unique item IDs and valid category references', () => {
const itemIds = [
...initialPortal.liveStreams.map((item) => item.stream_id),
...initialPortal.vodStreams.map((item) => item.stream_id),
...initialPortal.seriesItems.map((item) => item.series_id),
];
const liveCategoryIds = new Set(
initialPortal.liveCategories.map((item) => item.category_id)
);
const vodCategoryIds = new Set(
initialPortal.vodCategories.map((item) => item.category_id)
);
const seriesCategoryIds = new Set(
initialPortal.seriesCategories.map((item) => item.category_id)
);
expect(new Set(itemIds).size).toBe(100_000);
expect(
initialPortal.liveStreams.every((item) =>
liveCategoryIds.has(item.category_id)
)
).toBe(true);
expect(
initialPortal.vodStreams.every((item) =>
vodCategoryIds.has(item.category_id)
)
).toBe(true);
expect(
initialPortal.seriesItems.every((item) =>
seriesCategoryIds.has(String(item.category_id))
)
).toBe(true);
});
it('rebuilds a distinct portal to the identical byte count and SHA-256', () => {
const source = readFileSync(
`${__dirname}/performance.generator.ts`,
'utf8'
);
expect(source).not.toMatch(
/@faker-js\/faker|\bfaker\b|Date\.now|Math\.random/
);
const dateNow = jest.spyOn(Date, 'now').mockImplementation(() => {
throw new Error('performance fixture used Date.now');
});
const mathRandom = jest.spyOn(Math, 'random').mockImplementation(() => {
throw new Error('performance fixture used Math.random');
});
try {
const firstFingerprint = portalFingerprint(initialPortal);
dataStore.resetAll();
const rebuiltPortal = dataStore.getPortalData(USERNAME, PASSWORD);
expect(rebuiltPortal).not.toBe(initialPortal);
expect(portalFingerprint(rebuiltPortal)).toEqual(firstFingerprint);
} finally {
dateNow.mockRestore();
mathRandom.mockRestore();
}
});
it('keeps catalog and detail URL fields local-only', () => {
const data = dataStore.getPortalData(USERNAME, PASSWORD);
const vodDetails = dataStore.getVodDetails(
USERNAME,
PASSWORD,
data.vodStreams[0].stream_id
);
const seriesInfo = dataStore.getSeriesInfo(
USERNAME,
PASSWORD,
data.seriesItems[0].series_id
);
expect(vodDetails).not.toBeNull();
expect(seriesInfo).not.toBeNull();
expectLocalOrEmpty([
...data.liveStreams.flatMap((item) => [
item.stream_icon,
item.direct_source,
]),
...data.vodStreams.flatMap((item) => [
item.stream_icon,
item.direct_source,
]),
...data.seriesItems.flatMap((item) => [
item.cover,
item.youtube_trailer,
...item.backdrop_path,
]),
...(vodDetails && !Array.isArray(vodDetails.info)
? [
vodDetails.info.kinopoisk_url,
vodDetails.info.cover_big,
vodDetails.info.movie_image,
vodDetails.info.youtube_trailer,
...vodDetails.info.backdrop_path,
vodDetails.movie_data?.direct_source ?? '',
]
: []),
...(seriesInfo
? [
seriesInfo.info.cover,
seriesInfo.info.youtube_trailer,
...seriesInfo.info.backdrop_path,
...seriesInfo.seasons.flatMap((season) => [
season.cover,
season.cover_big,
]),
...Object.values(seriesInfo.episodes).flatMap(
(episodes) =>
episodes.flatMap((episode) => [
episode.info.movie_image,
episode.direct_source,
])
),
]
: []),
]);
});
it('uses numeric fixed-epoch semantics for catalog timestamps', () => {
const firstSeries = initialPortal.seriesItems[0];
expect(initialPortal.liveStreams[0].added).toBe('1767225600');
expect(initialPortal.vodStreams[0].added).toBe('1767225600');
expect(
initialPortal.seriesItems.every((item) =>
/^\d+$/.test(item.last_modified)
)
).toBe(true);
expect(firstSeries.last_modified).toBe('1767225600');
expect(
new Date(Number(firstSeries.last_modified) * 1_000).toISOString()
).toBe('2026-01-01T00:00:00.000Z');
expect(firstSeries.releaseDate).toBe('2026-01-01');
});
it('keeps account-info lazy and materializes one series detail through dispatch', () => {
dataStore.resetAll();
const accountResult = dispatch({
action: 'get_account_info',
username: USERNAME,
password: PASSWORD,
});
expect([accountResult.statusCode, accountResult.jsonCalls]).toEqual([
200, 1,
]);
expect(
(accountResult.body as { user_info?: { username?: string } })
.user_info?.username
).toBe(USERNAME);
expect(dataStore.getDetailCacheCardinalityForTesting()).toEqual({
vodDetails: 0,
seriesInfo: 0,
});
const seriesId = dataStore.getPortalData(USERNAME, PASSWORD)
.seriesItems[0].series_id;
const seriesResult = dispatch({
action: 'get_series_info',
username: USERNAME,
password: PASSWORD,
series_id: String(seriesId),
});
const seriesBody = seriesResult.body as {
seasons?: unknown[];
episodes?: Record<string, unknown[]>;
};
expect([seriesResult.statusCode, seriesResult.jsonCalls]).toEqual([
200, 1,
]);
expect(seriesBody.seasons).toHaveLength(1);
expect(seriesBody.episodes?.['1']).toHaveLength(1);
expect(dataStore.getDetailCacheCardinalityForTesting()).toEqual({
vodDetails: 0,
seriesInfo: 1,
});
});
it('keys VOD and series details by both username and password', () => {
withGenericFaker(() => {
dataStore.resetAll();
const portals = [
['emptyvod', 'emptyvod'],
['emptyvod', 'alternate'],
['alternate', 'emptyvod'],
] as const;
const firstPortal = dataStore.getPortalData(...portals[0]);
const vodId = firstPortal.vodStreams[0].stream_id;
const seriesId = firstPortal.seriesItems[0].series_id;
const vodDetails = portals.map(([username, password]) =>
dataStore.getVodDetails(username, password, vodId)
);
const seriesDetails = portals.map(([username, password]) =>
dataStore.getSeriesInfo(username, password, seriesId)
);
expect(vodDetails[0]?.info).toEqual([]);
expect(
vodDetails.slice(1).every((item) => !Array.isArray(item?.info))
).toBe(true);
expect(seriesDetails.map((item) => item?.seasons.length)).toEqual([
2, 3, 3,
]);
expect(
seriesDetails.map((item) => item?.episodes['1'].length)
).toEqual([4, 8, 8]);
expect(
dataStore.getDetailCacheCardinalityForTesting().vodDetails
).toBe(3);
});
});
});
@@ -0,0 +1,294 @@
import type { RawCategory } from './categories.generator.js';
import type { RawEpgListing, RawLiveStream } from './live.generator.js';
import type {
RawEpisode,
RawSeriesInfo,
RawSeriesItem,
RawSeason,
} from './series.generator.js';
import type { RawVodDetails, RawVodStream } from './vod.generator.js';
interface PerformancePortalFixture {
liveCategories: RawCategory[];
vodCategories: RawCategory[];
seriesCategories: RawCategory[];
liveStreams: RawLiveStream[];
vodStreams: RawVodStream[];
seriesItems: RawSeriesItem[];
epgListingsByStreamId: Map<number, RawEpgListing[]>;
}
const FIXED_EPOCH_SECONDS = 1_767_225_600;
const DAY_SECONDS = 86_400;
const ITEMS_PER_CATEGORY = 1_000;
const LIVE_CATEGORY_COUNT = 60;
const VOD_CATEGORY_COUNT = 20;
const SERIES_CATEGORY_COUNT = 20;
const LIVE_CATEGORY_ID_BASE = 91_100;
const VOD_CATEGORY_ID_BASE = 92_100;
const SERIES_CATEGORY_ID_BASE = 93_100;
const LIVE_STREAM_ID_BASE = 1_000_000;
const VOD_STREAM_ID_BASE = 2_000_000;
const SERIES_ID_BASE = 3_000_000;
const EPISODE_ID_BASE = 4_000_000;
const CONTAINER_EXTENSIONS = ['mp4', 'mkv', 'avi'] as const;
export function buildPerformancePortalFixture(): PerformancePortalFixture {
const liveCategories = buildCategories(
'Live',
LIVE_CATEGORY_ID_BASE,
LIVE_CATEGORY_COUNT
);
const vodCategories = buildCategories(
'VOD',
VOD_CATEGORY_ID_BASE,
VOD_CATEGORY_COUNT
);
const seriesCategories = buildCategories(
'Series',
SERIES_CATEGORY_ID_BASE,
SERIES_CATEGORY_COUNT
);
return {
liveCategories,
vodCategories,
seriesCategories,
liveStreams: buildLiveStreams(liveCategories),
vodStreams: buildVodStreams(vodCategories),
seriesItems: buildSeriesItems(seriesCategories),
epgListingsByStreamId: new Map(),
};
}
export function buildPerformanceVodDetails(
stream: RawVodStream
): RawVodDetails {
const index = stream.stream_id - VOD_STREAM_ID_BASE;
const durationSeconds = 5_400 + (index % 1_800);
const hours = Math.floor(durationSeconds / 3_600);
const minutes = Math.floor((durationSeconds % 3_600) / 60);
return {
info: {
kinopoisk_url: '',
tmdb_id: stream.stream_id,
name: stream.name,
o_name: stream.name,
cover_big: '',
movie_image: '',
releasedate: dateFromIndex(index),
episode_run_time: durationSeconds,
youtube_trailer: '',
director: `Performance Director ${formatIndex(index)}`,
actors: `Performance Actor ${formatIndex(index)}`,
cast: `Performance Actor ${formatIndex(index)}`,
description: `Synthetic movie ${formatIndex(index)}.`,
plot: `Synthetic movie ${formatIndex(index)}.`,
age: '12',
mpaa_rating: 'PG-13',
rating_count_kinopoisk: 1_000 + index,
country: 'Fixtureland',
genre: `Performance Genre ${index % 10}`,
backdrop_path: [],
duration_secs: durationSeconds,
duration: `${hours}h ${minutes}min`,
video: ['H.264'],
audio: ['AAC'],
bitrate: 2_000 + (index % 4_000),
rating: stream.rating,
rating_kinopoisk: stream.rating_imdb,
rating_imdb: stream.rating_imdb,
},
movie_data: {
stream_id: stream.stream_id,
name: stream.name,
added: stream.added,
category_id: stream.category_id,
container_extension: stream.container_extension,
custom_sid: '',
direct_source: '',
},
};
}
export function buildPerformanceSeriesInfo(
series: RawSeriesItem
): RawSeriesInfo {
const index = series.series_id - SERIES_ID_BASE;
const releaseDate = dateFromIndex(index);
const durationSeconds = 1_500 + (index % 1_800);
const season: RawSeason = {
air_date: releaseDate,
episode_count: 1,
id: series.series_id * 10 + 1,
name: 'Season 1',
overview: `Synthetic season for ${series.name}.`,
season_number: 1,
cover: '',
cover_big: '',
};
const episode: RawEpisode = {
id: String(EPISODE_ID_BASE + index),
episode_num: 1,
title: `${series.name} S01E01`,
container_extension: 'mkv',
info: {
tmdb_id: series.series_id,
releasedate: releaseDate,
plot: `Synthetic episode for ${series.name}.`,
duration_secs: durationSeconds,
duration: `${Math.floor(durationSeconds / 60)}min`,
movie_image: '',
bitrate: 1_500 + (index % 5_000),
rating: ratingFromIndex(index),
},
custom_sid: '',
added: addedFromIndex(index),
season: 1,
direct_source: '',
};
return {
seasons: [season],
info: {
name: series.name,
cover: '',
plot: series.plot,
cast: series.cast,
director: series.director,
genre: series.genre,
releaseDate: series.releaseDate,
last_modified: series.last_modified,
rating: series.rating,
rating_5based: series.rating_5based,
backdrop_path: [],
youtube_trailer: '',
episode_run_time: series.episode_run_time,
category_id: String(series.category_id),
},
episodes: { '1': [episode] },
};
}
function buildCategories(
label: string,
idBase: number,
count: number
): RawCategory[] {
return Array.from({ length: count }, (_, index) => ({
category_id: String(idBase + index + 1),
category_name: `Performance ${label} ${String(index + 1).padStart(
2,
'0'
)}`,
parent_id: 0,
}));
}
function buildLiveStreams(categories: RawCategory[]): RawLiveStream[] {
const streams: RawLiveStream[] = [];
for (const category of categories) {
for (let itemIndex = 0; itemIndex < ITEMS_PER_CATEGORY; itemIndex++) {
const index = streams.length;
const streamId = LIVE_STREAM_ID_BASE + index;
streams.push({
num: index + 1,
name: `Performance Live ${formatIndex(index)}`,
stream_type: 'live',
stream_id: streamId,
stream_icon: '',
epg_channel_id: `performance-channel-${streamId}`,
added: addedFromIndex(index),
category_id: category.category_id,
custom_sid: '',
direct_source: '',
tv_archive: index % 2,
tv_archive_duration: index % 2 === 0 ? 0 : 3,
rating_imdb: ratingFromIndex(index).toFixed(1),
});
}
}
return streams;
}
function buildVodStreams(categories: RawCategory[]): RawVodStream[] {
const streams: RawVodStream[] = [];
for (const category of categories) {
for (let itemIndex = 0; itemIndex < ITEMS_PER_CATEGORY; itemIndex++) {
const index = streams.length;
const rating = ratingFromIndex(index);
streams.push({
num: index + 1,
name: `Performance Movie ${formatIndex(index)}`,
stream_type: 'movie',
stream_id: VOD_STREAM_ID_BASE + index,
stream_icon: '',
added: addedFromIndex(index),
category_id: category.category_id,
custom_sid: '',
direct_source: '',
rating,
rating_5based: Number((rating / 2).toFixed(1)),
rating_imdb: rating.toFixed(1),
container_extension:
CONTAINER_EXTENSIONS[index % CONTAINER_EXTENSIONS.length],
type: 'movie',
});
}
}
return streams;
}
function buildSeriesItems(categories: RawCategory[]): RawSeriesItem[] {
const items: RawSeriesItem[] = [];
for (const category of categories) {
for (let itemIndex = 0; itemIndex < ITEMS_PER_CATEGORY; itemIndex++) {
const index = items.length;
const rating = ratingFromIndex(index);
items.push({
num: index + 1,
name: `Performance Series ${formatIndex(index)}`,
series_id: SERIES_ID_BASE + index,
cover: '',
plot: `Synthetic series ${formatIndex(index)}.`,
cast: `Performance Actor ${formatIndex(index)}`,
director: `Performance Director ${formatIndex(index)}`,
genre: `Performance Genre ${index % 10}`,
releaseDate: dateFromIndex(index),
last_modified: addedFromIndex(index),
rating: rating.toFixed(1),
rating_5based: Number((rating / 2).toFixed(1)),
backdrop_path: [],
youtube_trailer: '',
episode_run_time: String(24 + (index % 37)),
category_id: Number(category.category_id),
});
}
}
return items;
}
function addedFromIndex(index: number): string {
return String(FIXED_EPOCH_SECONDS - index);
}
function dateFromIndex(index: number, dayRange = 3_650): string {
const timestampSeconds =
FIXED_EPOCH_SECONDS - (index % dayRange) * DAY_SECONDS;
return new Date(timestampSeconds * 1_000).toISOString().slice(0, 10);
}
function formatIndex(index: number): string {
return String(index + 1).padStart(6, '0');
}
function ratingFromIndex(index: number): number {
return Number((6 + (index % 30) / 10).toFixed(1));
}
@@ -0,0 +1,45 @@
import {
isAllowedPerformanceCategory,
parseRuleBody,
PerformanceControlError,
} from './performance-control-validation.js';
describe('Xtream performance category validation', () => {
it.each([
'091101',
'+91101',
' 91101',
'91101 ',
'9.1101e4',
'9007199254740993',
'99999999999999999999999999999999999999999999999999',
])('rejects the non-canonical category ID %p', (categoryId) => {
expect(
isAllowedPerformanceCategory('get_live_streams', categoryId)
).toBe(false);
expect(() =>
parseRuleBody('barrier', {
id: 'non-canonical-category',
match: {
scenario: 'performance-100k',
transport: 'direct',
action: 'get_live_streams',
occurrence: 1,
categoryId,
},
})
).toThrow(PerformanceControlError);
});
it.each([
['get_live_streams', '91101'],
['get_live_streams', '91160'],
['get_vod_streams', '92101'],
['get_vod_streams', '92120'],
['get_series', '93101'],
['get_series', '93120'],
] as const)('accepts canonical %s category %s', (action, categoryId) => {
expect(isAllowedPerformanceCategory(action, categoryId)).toBe(true);
});
});
@@ -0,0 +1,149 @@
import { resetAll } from './data-store.js';
import type { PerformanceControlState } from './performance-control.js';
import { createXtreamMockApp } from './server.js';
import {
openAbortableGet,
postJson,
requestJson,
startLoopbackServer,
waitFor,
} from './testing/http-server.fixture.js';
jest.mock('@faker-js/faker', () => ({
faker: { seed: jest.fn() },
}));
const TOKEN = 'delay-abort-token';
const AUTH_HEADERS = {
'x-iptvnator-performance-token': TOKEN,
};
describe('Xtream performance delay cancellation', () => {
it('aborts an active delay without later dispatch or response lifecycle', async () => {
resetAll();
const running = await startLoopbackServer(
createXtreamMockApp({
control: { enabled: true, token: TOKEN },
host: '127.0.0.1',
port: 0,
})
);
try {
const rule = await postJson(
running.origin,
'/__control/delays',
TOKEN,
{
id: 'abort-delay',
match: {
scenario: 'performance-100k',
transport: 'direct',
action: 'get_live_streams',
occurrence: 1,
categoryId: '91101',
},
milliseconds: 100,
}
);
expect(rule.response.status).toBe(201);
const request = openAbortableGet(
running.origin,
'/player_api.php?username=performance&password=performance&action=get_live_streams&category_id=91101'
);
await waitFor(async () => {
expect(
(await state(running.origin)).ledger.map(
(entry) => entry.status
)
).toEqual(expect.arrayContaining(['arrived', 'delayed']));
});
request.abort();
await request.completed;
await waitFor(async () => {
const snapshot = await state(running.origin);
expect(snapshot.heldCount).toBe(0);
expect(snapshot.rules).toEqual([]);
expect(snapshot.ledger.map((entry) => entry.status)).toEqual(
expect.arrayContaining(['aborted'])
);
});
await new Promise((resolve) => setTimeout(resolve, 125));
const terminal = await state(running.origin);
expect(
terminal.ledger.some((entry) => entry.status === 'responded')
).toBe(false);
expect(terminal.occurrences).toEqual([
expect.objectContaining({
action: 'get_live_streams',
categoryId: '91101',
count: 1,
}),
]);
} finally {
await running.close();
resetAll();
}
});
it('settles an active delay when observations reset', async () => {
resetAll();
const running = await startLoopbackServer(
createXtreamMockApp({
control: { enabled: true, token: TOKEN },
host: '127.0.0.1',
port: 0,
})
);
try {
await postJson(running.origin, '/__control/delays', TOKEN, {
id: 'reset-delay',
match: {
scenario: 'performance-100k',
transport: 'direct',
action: 'get_live_streams',
occurrence: 1,
categoryId: '91101',
},
milliseconds: 100,
});
const delayedResponse = fetch(
`${running.origin}/player_api.php?username=performance&password=performance&action=get_live_streams&category_id=91101`
);
await waitFor(async () => {
expect(
(await state(running.origin)).ledger.map(
(entry) => entry.status
)
).toContain('delayed');
});
const reset = await postJson(
running.origin,
'/__control/reset',
TOKEN,
{ mode: 'observations' }
);
expect(reset.response.status).toBe(200);
expect((await delayedResponse).status).toBe(409);
await new Promise((resolve) => setTimeout(resolve, 125));
expect(await state(running.origin)).toMatchObject({
heldCount: 0,
ledger: [],
occurrences: [],
rules: [],
});
} finally {
await running.close();
resetAll();
}
});
});
async function state(origin: string): Promise<PerformanceControlState> {
const result = await requestJson(origin, '/__control/state', {
headers: AUTH_HEADERS,
});
expect(result.response.status).toBe(200);
return result.body as PerformanceControlState;
}
@@ -0,0 +1,83 @@
import { EventEmitter } from 'node:events';
import type { Request, Response } from 'express';
import { XtreamPerformanceController } from './performance-control.js';
jest.mock('@faker-js/faker', () => ({
faker: { seed: jest.fn() },
}));
describe('Xtream performance controller lifecycle', () => {
it.each([
['observations', 1],
['all', 2],
] as const)(
'invalidates an unmatched pre-reset response on %s reset',
async (mode, expectedEpoch) => {
const controller = new XtreamPerformanceController();
const request = createRequest();
const response = createResponse();
await controller.intercept(
request,
response as unknown as Response,
'direct',
() => undefined
);
expect(controller.snapshot().ledger).toHaveLength(1);
controller.reset(mode);
expect(response.listenerCount('finish')).toBe(0);
expect(response.listenerCount('close')).toBe(0);
expect(response.destroyed).toBe(true);
response.writableFinished = true;
response.emit('finish');
response.emit('close');
expect(controller.snapshot()).toMatchObject({
epoch: expectedEpoch,
heldIds: [],
heldCount: 0,
occurrences: [],
ledger: [],
});
}
);
});
function createRequest(): Request {
return Object.assign(new EventEmitter(), {
aborted: false,
query: {
username: 'performance',
password: 'performance',
},
}) as unknown as Request;
}
function createResponse(): EventEmitter & {
destroyed: boolean;
headersSent: boolean;
writableFinished: boolean;
destroy(): void;
json(): unknown;
status(): unknown;
} {
const response = Object.assign(new EventEmitter(), {
destroyed: false,
headersSent: false,
writableFinished: false,
destroy() {
response.destroyed = true;
response.emit('close');
},
json() {
return response;
},
status() {
return response;
},
});
return response;
}
@@ -0,0 +1,131 @@
import { timingSafeEqual } from 'node:crypto';
import cors from 'cors';
import express, {
type ErrorRequestHandler,
type NextFunction,
type Request,
type Response,
} from 'express';
import {
PERFORMANCE_CONTROL_LIMITS,
XtreamPerformanceController,
} from './performance-control.js';
import {
assertSafeRuleId,
parsePrepareBody,
parseReleaseBody,
parseResetBody,
parseRuleBody,
PerformanceControlError,
} from './performance-control-validation.js';
const TOKEN_HEADER = 'x-iptvnator-performance-token';
export function installPerformanceControlRoutes(
app: express.Express,
controller: XtreamPerformanceController,
token: string
): void {
const router = express.Router();
router.use(authorize(token));
router.use(cors());
router.use(
express.json({
limit: PERFORMANCE_CONTROL_LIMITS.jsonBytes,
strict: true,
})
);
router.post('/prepare', (request, response) => {
executeControl(response, () => {
parsePrepareBody(request.body);
response.json(controller.prepare());
});
});
router.post('/reset', (request, response) => {
executeControl(response, () => {
const { mode } = parseResetBody(request.body);
response.json(controller.reset(mode));
});
});
router.post('/barriers', (request, response) => {
executeControl(response, () => {
const rule = controller.addRule(
parseRuleBody('barrier', request.body)
);
response.status(201).json(rule);
});
});
router.post('/barriers/:id/release', (request, response) => {
executeControl(response, () => {
const id = request.params['id'] ?? '';
assertSafeRuleId(id);
parseReleaseBody(request.body);
if (!controller.releaseBarrier(id)) {
throw new PerformanceControlError(
404,
'held-barrier-not-found'
);
}
response.json({ id, released: true });
});
});
router.post('/delays', (request, response) => {
executeControl(response, () => {
const rule = controller.addRule(
parseRuleBody('delay', request.body)
);
response.status(201).json(rule);
});
});
router.get('/state', (_request, response) => {
response.json(controller.snapshot());
});
router.use(jsonErrorHandler);
app.use('/__control', router);
}
function authorize(token: string) {
const expected = Buffer.from(token, 'utf8');
return (request: Request, response: Response, next: NextFunction) => {
const provided = request.get(TOKEN_HEADER);
if (!provided || !constantTimeEqual(expected, provided)) {
response.status(401).json({ error: 'unauthorized' });
return;
}
next();
};
}
function constantTimeEqual(expected: Buffer, provided: string): boolean {
const candidate = Buffer.from(provided, 'utf8');
return (
candidate.length === expected.length &&
timingSafeEqual(candidate, expected)
);
}
function executeControl(response: Response, execute: () => void): void {
try {
execute();
} catch (error) {
if (error instanceof PerformanceControlError) {
response.status(error.status).json({ error: error.code });
return;
}
response.status(500).json({ error: 'control-internal-error' });
}
}
const jsonErrorHandler: ErrorRequestHandler = (
error: Error & { status?: number; type?: string },
_request,
response,
_next
) => {
void _next;
const tooLarge = error.status === 413 || error.type === 'entity.too.large';
response
.status(tooLarge ? 413 : 400)
.json({ error: tooLarge ? 'control-json-too-large' : 'invalid-json' });
};
@@ -0,0 +1,129 @@
import { getPortalData, resetAll } from './data-store.js';
import { type PerformanceControlState } from './performance-control.js';
import { createXtreamMockApp } from './server.js';
import {
postJson,
requestJson,
startLoopbackServer,
type RunningTestServer,
} from './testing/http-server.fixture.js';
jest.mock('@faker-js/faker', () => ({
faker: { seed: jest.fn() },
}));
const TOKEN = 'security-token';
const ORIGIN = 'http://127.0.0.1:4200';
describe('Xtream performance control request boundary', () => {
let running: RunningTestServer;
beforeEach(async () => {
resetAll();
running = await startLoopbackServer(
createXtreamMockApp({
control: { enabled: true, token: TOKEN },
host: '127.0.0.1',
port: 0,
})
);
});
afterEach(async () => {
await running.close();
resetAll();
});
it('requires the exact token before handling control preflight', async () => {
const missing = await controlOptions(running);
const wrong = await controlOptions(running, `${TOKEN}-wrong`);
const valid = await controlOptions(running, TOKEN);
expect(missing.status).toBe(401);
expect(wrong.status).toBe(401);
expect(valid.status).toBe(204);
expect(valid.headers.get('access-control-allow-origin')).toBe('*');
});
it('preserves unauthenticated CORS preflight for ordinary routes', async () => {
const response = await fetch(`${running.origin}/player_api.php`, {
method: 'OPTIONS',
headers: {
origin: ORIGIN,
'access-control-request-method': 'GET',
},
});
expect(response.status).toBe(204);
expect(response.headers.get('access-control-allow-origin')).toBe('*');
});
it('rejects legacy reset without invalidating prepared control state', async () => {
const prepared = await postJson(
running.origin,
'/__control/prepare',
TOKEN,
{ scenario: 'performance-100k' }
);
const cachedBefore = getPortalData('performance', 'performance');
const reset = await requestJson(running.origin, '/reset', {
method: 'POST',
});
const state = await requestJson(running.origin, '/__control/state', {
headers: {
'x-iptvnator-performance-token': TOKEN,
},
});
expect(reset.response.status).toBe(410);
expect(reset.body).toEqual({
error: 'use-performance-control-reset',
});
expect((state.body as PerformanceControlState).prepared).toEqual(
prepared.body
);
expect(getPortalData('performance', 'performance')).toBe(cachedBefore);
});
});
describe('Xtream normal reset compatibility', () => {
afterEach(() => resetAll());
it('keeps the legacy reset available outside control mode', async () => {
const running = await startLoopbackServer(
createXtreamMockApp({ host: '127.0.0.1', port: 0 })
);
const cachedBefore = getPortalData('performance', 'performance');
try {
const reset = await requestJson(running.origin, '/reset', {
method: 'POST',
});
expect(reset.response.status).toBe(200);
expect(reset.body).toEqual({ status: 'reset' });
expect(getPortalData('performance', 'performance')).not.toBe(
cachedBefore
);
} finally {
await running.close();
}
});
});
function controlOptions(
running: RunningTestServer,
token?: string
): Promise<Response> {
return fetch(`${running.origin}/__control/state`, {
method: 'OPTIONS',
headers: {
origin: ORIGIN,
'access-control-request-method': 'GET',
'access-control-request-headers': 'x-iptvnator-performance-token',
...(token === undefined
? {}
: { 'x-iptvnator-performance-token': token }),
},
});
}
@@ -0,0 +1,341 @@
import {
PERFORMANCE_CONTROL_LIMITS,
type PerformanceControlState,
} from './performance-control.js';
import { resetAll } from './data-store.js';
import { createXtreamMockApp } from './server.js';
import {
postJson,
requestJson,
startLoopbackServer,
type RunningTestServer,
} from './testing/http-server.fixture.js';
jest.mock('@faker-js/faker', () => ({
faker: { seed: jest.fn() },
}));
const TOKEN = 'validation-token';
const AUTH_HEADERS = {
'x-iptvnator-performance-token': TOKEN,
};
jest.setTimeout(60_000);
describe('Xtream performance control validation and bounds', () => {
let running: RunningTestServer;
beforeEach(async () => {
resetAll();
running = await startLoopbackServer(
createXtreamMockApp({
control: { enabled: true, token: TOKEN },
host: '127.0.0.1',
port: 0,
})
);
});
afterEach(async () => {
await running.close();
resetAll();
});
it.each([
['/__control/prepare', {}, 400],
[
'/__control/prepare',
{ scenario: 'performance-100k', extra: true },
400,
],
['/__control/prepare', { scenario: 'large' }, 400],
['/__control/reset', { mode: 'observations', extra: true }, 400],
['/__control/reset', { mode: 'wrong' }, 400],
['/__control/barriers', { id: '../unsafe', match: validMatch(1) }, 400],
[
'/__control/barriers',
{
id: 'bad-scenario',
match: { ...validMatch(1), scenario: 'large' },
},
400,
],
[
'/__control/barriers',
{
id: 'bad-transport',
match: { ...validMatch(1), transport: 'remote' },
},
400,
],
[
'/__control/barriers',
{
id: 'bad-action',
match: { ...validMatch(1), action: 'delete_everything' },
},
400,
],
[
'/__control/barriers',
{ id: 'zero-occurrence', match: validMatch(0) },
400,
],
[
'/__control/barriers',
{ id: 'fraction-occurrence', match: validMatch(1.5) },
400,
],
[
'/__control/barriers',
{
id: 'bad-category',
match: { ...validMatch(1), categoryId: '99999' },
},
400,
],
[
'/__control/barriers',
{
id: 'unknown-field',
match: { ...validMatch(1), url: 'http://forbidden.invalid' },
},
400,
],
[
'/__control/delays',
{ id: 'negative-delay', match: validMatch(1), milliseconds: -1 },
400,
],
[
'/__control/delays',
{ id: 'long-delay', match: validMatch(1), milliseconds: 5001 },
400,
],
[
'/__control/delays',
{
id: 'fraction-delay',
match: validMatch(1),
milliseconds: 1.5,
},
400,
],
])('rejects invalid body %#', async (path, body, status) => {
const result = await postJson(running.origin, path, TOKEN, body);
expect(result.response.status).toBe(status);
});
it('rejects malformed and oversized JSON safely', async () => {
const malformed = await fetch(`${running.origin}/__control/prepare`, {
method: 'POST',
headers: {
...AUTH_HEADERS,
'content-type': 'application/json',
},
body: '{"scenario":',
});
const oversized = await fetch(`${running.origin}/__control/prepare`, {
method: 'POST',
headers: {
...AUTH_HEADERS,
'content-type': 'application/json',
},
body: JSON.stringify({
scenario: 'performance-100k',
padding: 'x'.repeat(PERFORMANCE_CONTROL_LIMITS.jsonBytes + 1),
}),
});
expect(malformed.status).toBe(400);
expect(oversized.status).toBe(413);
expect(await malformed.text()).not.toContain('scenario');
expect(await oversized.text()).not.toContain('padding');
});
it('rejects arbitrary release bodies before looking up a barrier', async () => {
const response = await fetch(
`${running.origin}/__control/barriers/not-held/release`,
{
method: 'POST',
headers: {
...AUTH_HEADERS,
'content-type': 'application/json',
},
body: JSON.stringify({ url: 'http://forbidden.invalid' }),
}
);
expect(response.status).toBe(400);
expect(await response.text()).not.toContain('forbidden');
});
it('rejects duplicate IDs, duplicate matches, and past occurrences', async () => {
const first = await postJson(
running.origin,
'/__control/barriers',
TOKEN,
{ id: 'first', match: validMatch(1) }
);
const duplicateId = await postJson(
running.origin,
'/__control/barriers',
TOKEN,
{
id: 'first',
match: {
...validMatch(1),
categoryId: '91102',
},
}
);
const duplicateMatch = await postJson(
running.origin,
'/__control/delays',
TOKEN,
{ id: 'same-match', match: validMatch(1), milliseconds: 0 }
);
expect(first.response.status).toBe(201);
expect(duplicateId.response.status).toBe(409);
expect(duplicateMatch.response.status).toBe(409);
await postJson(running.origin, '/__control/reset', TOKEN, {
mode: 'observations',
});
await fetch(
`${running.origin}/player_api.php?username=performance&password=performance&action=get_live_streams&category_id=91101`
);
const past = await postJson(
running.origin,
'/__control/barriers',
TOKEN,
{ id: 'past', match: validMatch(1) }
);
expect(past.response.status).toBe(409);
});
it('caps accepted rules and the serialized ledger', async () => {
for (
let occurrence = 1;
occurrence <= PERFORMANCE_CONTROL_LIMITS.rules;
occurrence++
) {
const result = await postJson(
running.origin,
'/__control/barriers',
TOKEN,
{
id: `rule-${occurrence}`,
match: validMatch(occurrence),
}
);
expect(result.response.status).toBe(201);
}
const overflow = await postJson(
running.origin,
'/__control/barriers',
TOKEN,
{
id: 'rule-overflow',
match: validMatch(PERFORMANCE_CONTROL_LIMITS.rules + 1),
}
);
expect(overflow.response.status).toBe(409);
await postJson(running.origin, '/__control/reset', TOKEN, {
mode: 'observations',
});
for (
let index = 0;
index < PERFORMANCE_CONTROL_LIMITS.ledger;
index++
) {
await fetch(
`${running.origin}/player_api.php?username=u${index}&password=p${index}&action=unknown-${index}`
);
}
const snapshot = await state(running);
expect(snapshot.ledger).toHaveLength(PERFORMANCE_CONTROL_LIMITS.ledger);
expect(snapshot.occurrences.length).toBeLessThanOrEqual(
PERFORMANCE_CONTROL_LIMITS.occurrences
);
});
it('never serializes credentials, raw queries, URLs, or response payloads', async () => {
const username = 'SECRET_USER_MARKER';
const password = 'SECRET_PASSWORD_MARKER';
const rawUrl = 'http://RAW_URL_MARKER.invalid/private';
const responseMarker = 'RESPONSE_PAYLOAD_MARKER';
await postJson(running.origin, '/__control/prepare', TOKEN, {
scenario: 'performance-100k',
});
const response = await fetch(
`${running.origin}/xtream?url=${encodeURIComponent(
rawUrl
)}&username=${username}&password=${password}&action=${responseMarker}&query=RAW_QUERY_MARKER`
);
expect(response.status).toBe(400);
expect(await response.text()).toContain(responseMarker);
const serialized = JSON.stringify(await state(running));
for (const marker of [
username,
password,
rawUrl,
'RAW_URL_MARKER',
'RAW_QUERY_MARKER',
responseMarker,
TOKEN,
'Performance Live 000001',
]) {
expect(serialized).not.toContain(marker);
}
});
it('accepts the legacy action alias but rejects arbitrary cardinalities', async () => {
const alias = await postJson(
running.origin,
'/__control/barriers',
TOKEN,
{
id: 'legacy-alias',
match: {
scenario: 'performance-100k',
transport: 'direct',
action: 'get_simple_date_table',
occurrence: 1,
categoryId: 'all',
},
}
);
const arbitraryCategory = await postJson(
running.origin,
'/__control/barriers',
TOKEN,
{
id: 'arbitrary-category',
match: { ...validMatch(2), categoryId: '91161' },
}
);
expect(alias.response.status).toBe(201);
expect(arbitraryCategory.response.status).toBe(400);
});
});
function validMatch(occurrence: number) {
return {
scenario: 'performance-100k',
transport: 'direct',
action: 'get_live_streams',
occurrence,
categoryId: '91101',
};
}
async function state(
running: RunningTestServer
): Promise<PerformanceControlState> {
const result = await requestJson(running.origin, '/__control/state', {
headers: AUTH_HEADERS,
});
expect(result.response.status).toBe(200);
return result.body as PerformanceControlState;
}
@@ -0,0 +1,185 @@
import type {
PerformanceAction,
PerformanceRuleInput,
PerformanceRuleMatch,
PerformanceTransport,
} from './performance-control.types.js';
const RULE_ID_PATTERN = /^[A-Za-z0-9][A-Za-z0-9_-]{0,63}$/;
const PERFORMANCE_SCENARIO = 'performance-100k';
const TRANSPORTS: readonly PerformanceTransport[] = ['direct', 'proxy'];
const ACTIONS: readonly PerformanceAction[] = [
'get_account_info',
'get_live_categories',
'get_vod_categories',
'get_series_categories',
'get_live_streams',
'get_vod_streams',
'get_series',
'get_vod_info',
'get_series_info',
'get_short_epg',
'get_simple_data_table',
'get_simple_date_table',
];
export class PerformanceControlError extends Error {
constructor(
readonly status: number,
readonly code: string
) {
super(code);
}
}
export function parsePrepareBody(value: unknown): {
scenario: typeof PERFORMANCE_SCENARIO;
} {
requireExactObject(value, ['scenario']);
if (value['scenario'] !== PERFORMANCE_SCENARIO) {
invalid();
}
return { scenario: PERFORMANCE_SCENARIO };
}
export function parseResetBody(value: unknown): {
mode: 'all' | 'observations';
} {
requireExactObject(value, ['mode']);
const mode = value['mode'];
if (mode !== 'all' && mode !== 'observations') invalid();
return { mode };
}
export function parseRuleBody(
kind: 'barrier' | 'delay',
value: unknown
): PerformanceRuleInput {
const keys =
kind === 'barrier' ? ['id', 'match'] : ['id', 'match', 'milliseconds'];
requireExactObject(value, keys);
const id = value['id'];
if (typeof id !== 'string' || !RULE_ID_PATTERN.test(id)) invalid();
const match = parseMatch(value['match']);
if (kind === 'barrier') {
return { id, kind, match };
}
const milliseconds = value['milliseconds'];
if (
!Number.isInteger(milliseconds) ||
(milliseconds as number) < 0 ||
(milliseconds as number) > 5_000
) {
invalid();
}
return { id, kind, match, milliseconds: milliseconds as number };
}
export function parseReleaseBody(value: unknown): void {
if (value === undefined) return;
requireExactObject(value, []);
}
export function assertSafeRuleId(value: string): void {
if (!RULE_ID_PATTERN.test(value)) invalid();
}
export function isPerformanceAction(value: string): value is PerformanceAction {
return ACTIONS.includes(value as PerformanceAction);
}
export function canonicalObservedAction(
value: unknown
): PerformanceAction | 'unknown' {
const raw = typeof value === 'string' ? value : '';
if (raw === '') return 'get_account_info';
return isPerformanceAction(raw) ? raw : 'unknown';
}
function parseMatch(value: unknown): PerformanceRuleMatch {
requireExactObject(value, [
'scenario',
'transport',
'action',
'occurrence',
'categoryId',
]);
const scenario = value['scenario'];
const transport = value['transport'];
const action = value['action'];
const occurrence = value['occurrence'];
const categoryId = value['categoryId'];
if (scenario !== PERFORMANCE_SCENARIO) invalid();
if (
typeof transport !== 'string' ||
!TRANSPORTS.includes(transport as PerformanceTransport)
) {
invalid();
}
if (typeof action !== 'string' || !isPerformanceAction(action)) invalid();
if (!Number.isInteger(occurrence) || (occurrence as number) <= 0) invalid();
if (
typeof categoryId !== 'string' ||
!isAllowedPerformanceCategory(action, categoryId)
) {
invalid();
}
return {
scenario,
transport: transport as PerformanceTransport,
action,
occurrence: occurrence as number,
categoryId,
};
}
export function isAllowedPerformanceCategory(
action: PerformanceAction,
categoryId: string
): boolean {
if (categoryId === 'all') return true;
if (!/^\d+$/.test(categoryId)) return false;
const numericId = Number(categoryId);
if (!Number.isSafeInteger(numericId) || categoryId !== String(numericId)) {
return false;
}
if (action === 'get_live_streams') {
return numericId >= 91_101 && numericId <= 91_160;
}
if (action === 'get_vod_streams') {
return numericId >= 92_101 && numericId <= 92_120;
}
if (action === 'get_series') {
return numericId >= 93_101 && numericId <= 93_120;
}
return false;
}
function requireExactObject(
value: unknown,
expectedKeys: readonly string[]
): asserts value is Record<string, unknown> {
if (
value === null ||
typeof value !== 'object' ||
Array.isArray(value) ||
Object.getPrototypeOf(value) !== Object.prototype
) {
invalid();
}
const actualKeys = Object.keys(value).sort();
const sortedExpected = [...expectedKeys].sort();
if (
actualKeys.length !== sortedExpected.length ||
actualKeys.some((key, index) => key !== sortedExpected[index])
) {
invalid();
}
}
function invalid(): never {
throw new PerformanceControlError(400, 'invalid-control-request');
}
@@ -0,0 +1,380 @@
import { resetAll } from './data-store.js';
import { type PerformanceControlState } from './performance-control.js';
import { createXtreamMockApp } from './server.js';
import {
openAbortableGet,
postJson,
requestJson,
startLoopbackServer,
waitFor,
type RunningTestServer,
} from './testing/http-server.fixture.js';
jest.mock('@faker-js/faker', () => ({
faker: { seed: jest.fn() },
}));
const TOKEN = 'benchmark-token';
const AUTH_HEADERS = {
'x-iptvnator-performance-token': TOKEN,
};
jest.setTimeout(120_000);
describe('Xtream performance control API', () => {
let running: RunningTestServer;
beforeEach(async () => {
resetAll();
running = await startLoopbackServer(
createXtreamMockApp({
control: { enabled: true, token: TOKEN },
host: '127.0.0.1',
port: 0,
})
);
});
afterEach(async () => {
await running.close();
resetAll();
});
it('requires the exact token for every control route', async () => {
const missing = await fetch(`${running.origin}/__control/state`);
const incorrect = await fetch(`${running.origin}/__control/state`, {
headers: {
'x-iptvnator-performance-token': `${TOKEN}-wrong`,
},
});
const valid = await fetch(`${running.origin}/__control/state`, {
headers: AUTH_HEADERS,
});
expect([missing.status, incorrect.status, valid.status]).toEqual([
401, 401, 200,
]);
});
it('prepares the exact 100k fixture with a stable redacted manifest', async () => {
const first = await postJson(
running.origin,
'/__control/prepare',
TOKEN,
{ scenario: 'performance-100k' }
);
expect(first.response.status).toBe(200);
expect(Object.keys(first.body as object)).toEqual([
'epoch',
'scenario',
'seed',
'counts',
'bytes',
'catalogSha256',
]);
expect(first.body).toMatchObject({
epoch: 1,
scenario: 'performance-100k',
seed: 91_001,
counts: {
categories: { live: 60, vod: 20, series: 20, total: 100 },
items: {
live: 60_000,
vod: 20_000,
series: 20_000,
total: 100_000,
},
},
});
expect((first.body as { bytes: number }).bytes).toBeGreaterThan(0);
expect((first.body as { catalogSha256: string }).catalogSha256).toMatch(
/^[a-f0-9]{64}$/
);
const reset = await postJson(
running.origin,
'/__control/reset',
TOKEN,
{ mode: 'all' }
);
expect(reset.response.status).toBe(200);
const second = await postJson(
running.origin,
'/__control/prepare',
TOKEN,
{ scenario: 'performance-100k' }
);
expect(second.body).toMatchObject({
epoch: 2,
bytes: (first.body as { bytes: number }).bytes,
catalogSha256: (first.body as { catalogSha256: string })
.catalogSha256,
});
});
it('distinguishes observation reset from all-state reset', async () => {
const prepared = await prepare(running);
await postJson(running.origin, '/__control/barriers', TOKEN, {
id: 'future-live',
match: liveMatch('direct', '91101', 1),
});
await fetch(
`${running.origin}/player_api.php?username=performance&password=performance&action=get_account_info`
);
const before = await state(running);
expect(before.prepared).toEqual(prepared);
expect(before.rules).toHaveLength(1);
expect(before.ledger.length).toBeGreaterThan(0);
await postJson(running.origin, '/__control/reset', TOKEN, {
mode: 'observations',
});
const observationsReset = await state(running);
expect(observationsReset).toMatchObject({
epoch: 1,
prepared,
rules: [],
heldIds: [],
heldCount: 0,
occurrences: [],
ledger: [],
});
await postJson(running.origin, '/__control/reset', TOKEN, {
mode: 'all',
});
const allReset = await state(running);
expect(allReset).toMatchObject({
epoch: 2,
prepared: null,
rules: [],
heldIds: [],
heldCount: 0,
occurrences: [],
ledger: [],
});
});
it('matches category barriers independently of parallel arrival order', async () => {
await prepare(running);
await createBarrier(
running,
'category-a',
liveMatch('direct', '91101', 1)
);
await createBarrier(
running,
'category-b',
liveMatch('direct', '91102', 1)
);
const categoryB = fetch(catalogUrl(running, '91102', 'direct'));
await expectHeld(running, ['category-b']);
const categoryA = fetch(catalogUrl(running, '91101', 'direct'));
await expectHeld(running, ['category-a', 'category-b']);
await release(running, 'category-a');
await release(running, 'category-b');
const [bodyA, bodyB] = await Promise.all([
categoryA.then((response) => response.json()),
categoryB.then((response) => response.json()),
]);
expect(bodyA).toHaveLength(1_000);
expect(bodyB).toHaveLength(1_000);
const snapshot = await state(running);
expect(
snapshot.occurrences.filter(
(entry) =>
entry.action === 'get_live_streams' &&
['91101', '91102'].includes(entry.categoryId)
)
).toEqual(
expect.arrayContaining([
expect.objectContaining({ categoryId: '91101', count: 1 }),
expect.objectContaining({ categoryId: '91102', count: 1 }),
])
);
});
it('releases a barrier and records the terminal response', async () => {
await prepare(running);
await createBarrier(running, 'account-gate', {
scenario: 'performance-100k',
transport: 'direct',
action: 'get_account_info',
occurrence: 1,
categoryId: 'all',
});
const request = fetch(
`${running.origin}/player_api.php?username=performance&password=performance`
);
await expectHeld(running, ['account-gate']);
const releaseResponse = await release(running, 'account-gate');
expect(releaseResponse.status).toBe(200);
expect((await request).status).toBe(200);
await waitFor(async () => {
const snapshot = await state(running);
expect(snapshot.heldCount).toBe(0);
expect(snapshot.ledger.map((entry) => entry.status)).toEqual(
expect.arrayContaining(['arrived', 'blocked', 'responded'])
);
});
});
it('settles an aborted held request without leaking it', async () => {
await prepare(running);
await createBarrier(
running,
'abort-gate',
liveMatch('direct', '91101', 1)
);
const request = openAbortableGet(
running.origin,
'/player_api.php?username=performance&password=performance&action=get_live_streams&category_id=91101'
);
await expectHeld(running, ['abort-gate']);
request.abort();
await request.completed;
await waitFor(async () => {
const snapshot = await state(running);
expect(snapshot.heldCount).toBe(0);
expect(
snapshot.ledger.some(
(entry) =>
entry.action === 'get_live_streams' &&
entry.status === 'aborted'
)
).toBe(true);
});
});
it('settles held requests when observations are reset', async () => {
await prepare(running);
await createBarrier(
running,
'reset-gate',
liveMatch('direct', '91101', 1)
);
const heldResponse = fetch(
`${running.origin}/player_api.php?username=performance&password=performance&action=get_live_streams&category_id=91101`
);
await expectHeld(running, ['reset-gate']);
const reset = await postJson(
running.origin,
'/__control/reset',
TOKEN,
{ mode: 'observations' }
);
expect(reset.response.status).toBe(200);
expect((await heldResponse).status).toBe(409);
expect(await state(running)).toMatchObject({
epoch: 1,
heldIds: [],
heldCount: 0,
occurrences: [],
ledger: [],
rules: [],
});
});
it('applies a configured short delay and records its lifecycle', async () => {
await prepare(running);
const created = await postJson(
running.origin,
'/__control/delays',
TOKEN,
{
id: 'short-delay',
match: liveMatch('proxy', '91101', 1),
milliseconds: 1,
}
);
expect(created.response.status).toBe(201);
const response = await fetch(catalogUrl(running, '91101', 'proxy'));
expect(response.status).toBe(200);
const snapshot = await state(running);
expect(snapshot.heldCount).toBe(0);
expect(snapshot.ledger.map((entry) => entry.status)).toEqual(
expect.arrayContaining(['arrived', 'delayed', 'responded'])
);
});
});
async function prepare(
running: RunningTestServer
): Promise<PerformanceControlState['prepared']> {
const result = await postJson(running.origin, '/__control/prepare', TOKEN, {
scenario: 'performance-100k',
});
expect(result.response.status).toBe(200);
return result.body as PerformanceControlState['prepared'];
}
async function state(
running: RunningTestServer
): Promise<PerformanceControlState> {
const result = await requestJson(running.origin, '/__control/state', {
headers: AUTH_HEADERS,
});
expect(result.response.status).toBe(200);
return result.body as PerformanceControlState;
}
function liveMatch(
transport: 'direct' | 'proxy',
categoryId: string,
occurrence: number
) {
return {
scenario: 'performance-100k',
transport,
action: 'get_live_streams',
occurrence,
categoryId,
};
}
async function createBarrier(
running: RunningTestServer,
id: string,
match: ReturnType<typeof liveMatch> | Record<string, unknown>
) {
const result = await postJson(
running.origin,
'/__control/barriers',
TOKEN,
{ id, match }
);
expect(result.response.status).toBe(201);
}
async function release(
running: RunningTestServer,
id: string
): Promise<Response> {
return fetch(`${running.origin}/__control/barriers/${id}/release`, {
method: 'POST',
headers: AUTH_HEADERS,
});
}
async function expectHeld(
running: RunningTestServer,
ids: string[]
): Promise<void> {
await waitFor(async () => {
const snapshot = await state(running);
expect(snapshot.heldIds.sort()).toEqual([...ids].sort());
expect(snapshot.heldCount).toBe(ids.length);
});
}
function catalogUrl(
running: RunningTestServer,
categoryId: string,
transport: 'direct' | 'proxy'
): string {
const route = transport === 'direct' ? '/player_api.php' : '/xtream';
return `${running.origin}${route}?url=ignored&username=performance&password=performance&action=get_live_streams&category_id=${categoryId}`;
}
@@ -0,0 +1,299 @@
import { performance } from 'node:perf_hooks';
import type { Request, Response } from 'express';
import { resetAll } from './data-store.js';
import { PerformanceInterceptionLifecycle } from './performance-interception-lifecycle.js';
import { buildPerformanceManifest } from './performance-manifest.js';
import { PerformanceControlError } from './performance-control-validation.js';
import {
PERFORMANCE_CONTROL_LIMITS,
type PerformanceControlState,
type PerformanceIdentity,
type PerformanceLedgerEntry,
type PerformanceLifecycleStatus,
type PerformanceManifest,
type PerformanceRuleInput,
type PerformanceTransport,
} from './performance-control.types.js';
import {
identityKey,
PerformanceOccurrenceTracker,
ruleMatchKey,
} from './performance-occurrences.js';
export {
PERFORMANCE_CONTROL_LIMITS,
type PerformanceAction,
type PerformanceControlState,
type PerformanceLedgerEntry,
type PerformanceLifecycleStatus,
type PerformanceManifest,
type PerformanceOccurrenceState,
type PerformanceRuleInput,
type PerformanceRuleMatch,
type PerformanceTransport,
} from './performance-control.types.js';
type HeldOutcome = 'released' | 'aborted' | 'reset' | 'shutdown';
interface PendingRequest {
readonly kind: 'barrier' | 'delay';
settle(outcome: HeldOutcome): void;
}
export class XtreamPerformanceController {
private epoch = 1;
private prepared: PerformanceManifest | null = null;
private readonly rules = new Map<string, PerformanceRuleInput>();
private readonly matchKeys = new Set<string>();
private readonly seenRuleIds = new Set<string>();
private acceptedRuleCount = 0;
private readonly pending = new Map<string, PendingRequest>();
private readonly interceptionLifecycle =
new PerformanceInterceptionLifecycle();
private readonly occurrences = new PerformanceOccurrenceTracker();
private readonly ledger: PerformanceLedgerEntry[] = [];
prepare(): PerformanceManifest {
this.prepared ??= buildPerformanceManifest(this.epoch);
return this.prepared;
}
reset(mode: 'all' | 'observations'): { epoch: number; mode: string } {
this.interceptionLifecycle.reset();
this.settlePending('reset');
this.clearObservations();
if (mode === 'all') {
resetAll();
this.prepared = null;
this.epoch += 1;
}
return { epoch: this.epoch, mode };
}
shutdown(): PerformanceControlState {
if (this.interceptionLifecycle.shutdown()) {
this.settlePending('shutdown');
this.clearObservations();
}
return this.snapshot();
}
addRule(input: PerformanceRuleInput): PerformanceRuleInput {
if (this.acceptedRuleCount >= PERFORMANCE_CONTROL_LIMITS.rules) {
conflict('control-rule-limit');
}
if (this.seenRuleIds.has(input.id)) conflict('duplicate-rule-id');
const matchKey = ruleMatchKey(input.match);
if (this.matchKeys.has(matchKey)) conflict('duplicate-rule-match');
if (input.match.occurrence <= this.occurrences.count(input.match)) {
conflict('past-rule-occurrence');
}
this.acceptedRuleCount += 1;
this.seenRuleIds.add(input.id);
this.matchKeys.add(matchKey);
this.rules.set(input.id, input);
return input;
}
releaseBarrier(id: string): boolean {
const pending = this.pending.get(id);
if (!pending || pending.kind !== 'barrier') return false;
this.settleOnePending(id, 'released');
return true;
}
snapshot(): PerformanceControlState {
const heldIds = [...this.pending.entries()]
.filter(([, request]) => request.kind === 'barrier')
.map(([id]) => id);
return {
epoch: this.epoch,
prepared: this.prepared,
rules: [...this.rules.values()],
heldIds,
heldCount: heldIds.length,
occurrences: this.occurrences.snapshot(),
ledger: [...this.ledger],
};
}
async intercept(
request: Request,
response: Response,
transport: PerformanceTransport,
dispatch: () => void | Promise<void>
): Promise<void> {
if (this.interceptionLifecycle.shuttingDown) {
if (!response.destroyed) response.destroy();
return;
}
const generation = this.interceptionLifecycle.generation;
const identity = this.occurrences.next(request, transport, this.epoch);
this.record(identity, 'arrived');
const activeId = Symbol('performance-interception');
let terminal = false;
let suppressTerminal = false;
let pendingId: string | undefined;
const cleanup = () => {
request.off('aborted', onAborted);
response.off('finish', onFinished);
response.off('close', onClosed);
this.interceptionLifecycle.unregister(activeId);
};
const finish = (status: 'responded' | 'aborted') => {
if (terminal) return;
terminal = true;
cleanup();
if (pendingId) {
this.settleOnePending(pendingId, 'aborted');
pendingId = undefined;
}
if (
!suppressTerminal &&
this.interceptionLifecycle.isCurrent(generation)
) {
this.record(identity, status);
}
};
const onAborted = () => finish('aborted');
const onFinished = () => finish('responded');
const onClosed = () => {
if (!response.writableFinished) finish('aborted');
};
request.once('aborted', onAborted);
response.once('finish', onFinished);
response.once('close', onClosed);
this.interceptionLifecycle.register(activeId, (abortHeldResponse) => {
if (terminal) return;
terminal = true;
suppressTerminal = true;
cleanup();
if (
(abortHeldResponse || pendingId === undefined) &&
!response.writableFinished &&
!response.destroyed
) {
response.destroy();
}
});
const rule = this.consumeRule(identity);
if (rule?.kind === 'barrier') {
this.record(identity, 'blocked');
pendingId = rule.id;
const outcome = await new Promise<HeldOutcome>((resolve) => {
this.pending.set(rule.id, {
kind: 'barrier',
settle: (value) => {
if (value === 'reset') {
suppressTerminal = true;
cleanup();
}
resolve(value);
},
});
});
pendingId = undefined;
if (this.finishControlledWait(outcome, response)) return;
} else if (rule?.kind === 'delay') {
this.record(identity, 'delayed');
pendingId = rule.id;
const outcome = await new Promise<HeldOutcome>((resolve) => {
const timer = setTimeout(
() => this.settleOnePending(rule.id, 'released'),
rule.milliseconds
);
this.pending.set(rule.id, {
kind: 'delay',
settle: (value) => {
clearTimeout(timer);
if (value === 'reset') {
suppressTerminal = true;
cleanup();
}
resolve(value);
},
});
});
pendingId = undefined;
if (this.finishControlledWait(outcome, response)) return;
}
if (request.aborted || response.destroyed) {
finish('aborted');
return;
}
await dispatch();
}
private consumeRule(
identity: PerformanceIdentity
): PerformanceRuleInput | undefined {
for (const rule of this.rules.values()) {
if (ruleMatchKey(rule.match) !== identityKey(identity)) continue;
this.rules.delete(rule.id);
this.matchKeys.delete(ruleMatchKey(rule.match));
return rule;
}
return undefined;
}
private record(
identity: PerformanceIdentity,
status: PerformanceLifecycleStatus
): void {
this.ledger.push({
...identity,
status,
atMs: Number(performance.now().toFixed(3)),
});
if (this.ledger.length > PERFORMANCE_CONTROL_LIMITS.ledger) {
this.ledger.splice(
0,
this.ledger.length - PERFORMANCE_CONTROL_LIMITS.ledger
);
}
}
private finishControlledWait(
outcome: HeldOutcome,
response: Response
): boolean {
if (outcome === 'released') return false;
if (
outcome === 'reset' &&
!response.headersSent &&
!response.destroyed
) {
response.status(409).json({ error: 'control-reset' });
}
return true;
}
private settlePending(outcome: HeldOutcome): void {
for (const id of [...this.pending.keys()]) {
this.settleOnePending(id, outcome);
}
}
private settleOnePending(id: string, outcome: HeldOutcome): void {
const pending = this.pending.get(id);
if (!pending) return;
this.pending.delete(id);
pending.settle(outcome);
}
private clearObservations(): void {
this.rules.clear();
this.matchKeys.clear();
this.seenRuleIds.clear();
this.acceptedRuleCount = 0;
this.occurrences.clear();
this.ledger.length = 0;
}
}
function conflict(code: string): never {
throw new PerformanceControlError(409, code);
}
@@ -0,0 +1,98 @@
export const PERFORMANCE_CONTROL_LIMITS = {
jsonBytes: 16 * 1024,
ledger: 128,
occurrences: 512,
rules: 32,
} as const;
export type PerformanceTransport = 'direct' | 'proxy';
export type PerformanceAction =
| 'get_account_info'
| 'get_live_categories'
| 'get_vod_categories'
| 'get_series_categories'
| 'get_live_streams'
| 'get_vod_streams'
| 'get_series'
| 'get_vod_info'
| 'get_series_info'
| 'get_short_epg'
| 'get_simple_data_table'
| 'get_simple_date_table';
export type PerformanceLifecycleStatus =
'arrived' | 'blocked' | 'delayed' | 'responded' | 'aborted';
export interface PerformanceRuleMatch {
readonly scenario: 'performance-100k';
readonly transport: PerformanceTransport;
readonly action: PerformanceAction;
readonly occurrence: number;
readonly categoryId: string;
}
export type PerformanceRuleInput =
| {
readonly id: string;
readonly kind: 'barrier';
readonly match: PerformanceRuleMatch;
}
| {
readonly id: string;
readonly kind: 'delay';
readonly match: PerformanceRuleMatch;
readonly milliseconds: number;
};
export interface PerformanceManifest {
readonly epoch: number;
readonly scenario: 'performance-100k';
readonly seed: number;
readonly counts: {
readonly categories: {
readonly live: number;
readonly vod: number;
readonly series: number;
readonly total: number;
};
readonly items: {
readonly live: number;
readonly vod: number;
readonly series: number;
readonly total: number;
};
};
readonly bytes: number;
readonly catalogSha256: string;
}
export interface PerformanceIdentity {
readonly epoch: number;
readonly scenario: string;
readonly transport: PerformanceTransport;
readonly action: PerformanceAction | 'unknown';
readonly categoryId: string;
readonly occurrence: number;
}
export interface PerformanceOccurrenceState {
readonly scenario: string;
readonly transport: PerformanceTransport;
readonly action: PerformanceAction | 'unknown';
readonly categoryId: string;
readonly count: number;
}
export interface PerformanceLedgerEntry extends PerformanceIdentity {
readonly status: PerformanceLifecycleStatus;
readonly atMs: number;
}
export interface PerformanceControlState {
readonly epoch: number;
readonly prepared: PerformanceManifest | null;
readonly rules: ReadonlyArray<PerformanceRuleInput>;
readonly heldIds: string[];
readonly heldCount: number;
readonly occurrences: ReadonlyArray<PerformanceOccurrenceState>;
readonly ledger: ReadonlyArray<PerformanceLedgerEntry>;
}
@@ -0,0 +1,45 @@
type InvalidateInterception = (abortHeldResponse: boolean) => void;
export class PerformanceInterceptionLifecycle {
private currentGeneration = 1;
private stopped = false;
private readonly active = new Map<symbol, InvalidateInterception>();
get generation(): number {
return this.currentGeneration;
}
get shuttingDown(): boolean {
return this.stopped;
}
isCurrent(generation: number): boolean {
return generation === this.currentGeneration;
}
register(id: symbol, invalidate: InvalidateInterception): void {
this.active.set(id, invalidate);
}
unregister(id: symbol): void {
this.active.delete(id);
}
reset(): void {
this.invalidate(false);
}
shutdown(): boolean {
if (this.stopped) return false;
this.stopped = true;
this.invalidate(true);
return true;
}
private invalidate(abortHeldResponses: boolean): void {
this.currentGeneration += 1;
for (const invalidate of [...this.active.values()]) {
invalidate(abortHeldResponses);
}
}
}
@@ -0,0 +1,56 @@
import { createHash } from 'node:crypto';
import { getPortalData } from './data-store.js';
import type { PerformanceManifest } from './performance-control.types.js';
const PERFORMANCE_USERNAME = 'performance';
const PERFORMANCE_PASSWORD = 'performance';
/**
* The property insertion order is part of the fixture identity contract.
* Never add credentials, server origins, or response envelopes to this input.
*/
export function buildPerformanceManifest(epoch: number): PerformanceManifest {
const data = getPortalData(PERFORMANCE_USERNAME, PERFORMANCE_PASSWORD);
const catalog = {
liveCategories: data.liveCategories,
vodCategories: data.vodCategories,
seriesCategories: data.seriesCategories,
liveCatalog: data.liveStreams,
vodCatalog: data.vodStreams,
seriesCatalog: data.seriesItems,
};
const serialized = JSON.stringify(catalog);
const categoryCounts = {
live: data.liveCategories.length,
vod: data.vodCategories.length,
series: data.seriesCategories.length,
};
const itemCounts = {
live: data.liveStreams.length,
vod: data.vodStreams.length,
series: data.seriesItems.length,
};
return {
epoch,
scenario: 'performance-100k',
seed: data.scenario.seed,
counts: {
categories: {
...categoryCounts,
total:
categoryCounts.live +
categoryCounts.vod +
categoryCounts.series,
},
items: {
...itemCounts,
total: itemCounts.live + itemCounts.vod + itemCounts.series,
},
},
bytes: Buffer.byteLength(serialized, 'utf8'),
catalogSha256: createHash('sha256')
.update(serialized, 'utf8')
.digest('hex'),
};
}
@@ -0,0 +1,108 @@
import type { Request } from 'express';
import { PerformanceOccurrenceTracker } from './performance-occurrences.js';
import { PERFORMANCE_CONTROL_LIMITS } from './performance-control.types.js';
describe('PerformanceOccurrenceTracker', () => {
it('does not restart an old valid identity after more than 128 peers', () => {
const tracker = new PerformanceOccurrenceTracker();
const firstRequest = requestFor('get_live_streams', '91101');
expect(tracker.next(firstRequest, 'direct', 1).occurrence).toBe(1);
for (let category = 91_102; category <= 91_160; category++) {
tracker.next(
requestFor('get_live_streams', String(category)),
'direct',
1
);
}
for (let category = 91_101; category <= 91_160; category++) {
tracker.next(
requestFor('get_live_streams', String(category)),
'proxy',
1
);
}
for (let category = 92_101; category <= 92_120; category++) {
tracker.next(
requestFor('get_vod_streams', String(category)),
'direct',
1
);
}
expect(tracker.snapshot()).toHaveLength(140);
expect(tracker.next(firstRequest, 'direct', 1).occurrence).toBe(2);
});
it('collapses non-canonical category aliases into one bounded identity', () => {
const tracker = new PerformanceOccurrenceTracker();
const canonical = tracker.next(
requestFor('get_live_streams', '91101'),
'direct',
1
);
const aliases = [
'+91101',
' 91101',
'91101 ',
'9.1101e4',
'9007199254740993',
];
expect(canonical).toMatchObject({
categoryId: '91101',
occurrence: 1,
});
aliases.forEach((categoryId, index) => {
expect(
tracker.next(
requestFor('get_live_streams', categoryId),
'direct',
1
)
).toMatchObject({
categoryId: 'all',
occurrence: index + 1,
});
});
for (
let zeroCount = 1;
zeroCount <= PERFORMANCE_CONTROL_LIMITS.occurrences + 1;
zeroCount++
) {
tracker.next(
requestFor('get_live_streams', `${'0'.repeat(zeroCount)}91101`),
'direct',
1
);
}
expect(tracker.snapshot()).toHaveLength(2);
expect(tracker.snapshot()).toEqual(
expect.arrayContaining([
expect.objectContaining({
categoryId: '91101',
count: 1,
}),
expect.objectContaining({
categoryId: 'all',
count:
aliases.length +
PERFORMANCE_CONTROL_LIMITS.occurrences +
1,
}),
])
);
});
});
function requestFor(action: string, categoryId: string): Request {
return {
query: {
username: 'performance',
password: 'performance',
action,
category_id: categoryId,
},
} as unknown as Request;
}
@@ -0,0 +1,104 @@
import type { Request } from 'express';
import {
canonicalObservedAction,
isAllowedPerformanceCategory,
PerformanceControlError,
} from './performance-control-validation.js';
import {
PERFORMANCE_CONTROL_LIMITS,
type PerformanceAction,
type PerformanceIdentity,
type PerformanceOccurrenceState,
type PerformanceRuleMatch,
type PerformanceTransport,
} from './performance-control.types.js';
import { getScenario } from './scenarios.js';
export class PerformanceOccurrenceTracker {
private readonly entries = new Map<string, PerformanceOccurrenceState>();
next(
request: Request,
transport: PerformanceTransport,
epoch: number
): PerformanceIdentity {
const username = stringQuery(request.query['username']);
const password = stringQuery(request.query['password']);
const scenario = getScenario(username, password).name;
const action = canonicalObservedAction(request.query['action']);
const rawCategoryId = stringQuery(request.query['category_id']);
const categoryId =
scenario === 'performance-100k' &&
action !== 'unknown' &&
rawCategoryId !== 'all' &&
isAllowedPerformanceCategory(action, rawCategoryId)
? rawCategoryId
: 'all';
const base = { scenario, transport, action, categoryId };
const key = baseIdentityKey(base);
const previous = this.entries.get(key);
if (
!previous &&
this.entries.size >= PERFORMANCE_CONTROL_LIMITS.occurrences
) {
throw new PerformanceControlError(503, 'control-occurrence-limit');
}
const count = (previous?.count ?? 0) + 1;
this.entries.set(key, { ...base, count });
return { epoch, ...base, occurrence: count };
}
count(match: PerformanceRuleMatch): number {
return this.entries.get(baseIdentityKey(match))?.count ?? 0;
}
snapshot(): PerformanceOccurrenceState[] {
return [...this.entries.values()];
}
clear(): void {
this.entries.clear();
}
}
export function ruleMatchKey(match: PerformanceRuleMatch): string {
return JSON.stringify([
match.scenario,
match.transport,
match.action,
match.categoryId,
match.occurrence,
]);
}
export function identityKey(identity: PerformanceIdentity): string {
return JSON.stringify([
identity.scenario,
identity.transport,
identity.action,
identity.categoryId,
identity.occurrence,
]);
}
function baseIdentityKey(
identity:
| Omit<PerformanceRuleMatch, 'occurrence'>
| {
scenario: string;
transport: PerformanceTransport;
action: PerformanceAction | 'unknown';
categoryId: string;
}
): string {
return JSON.stringify([
identity.scenario,
identity.transport,
identity.action,
identity.categoryId,
]);
}
function stringQuery(value: unknown): string {
return typeof value === 'string' ? value : '';
}
@@ -0,0 +1,28 @@
import { readFileSync } from 'node:fs';
import { join } from 'node:path';
describe('Xtream mock Nx serve environment', () => {
it.each(['serve', 'serve-with-watch'])(
'lets the caller select PORT for the %s target',
(targetName) => {
const project = JSON.parse(
readFileSync(
join(process.cwd(), 'apps/xtream-mock-server/project.json'),
'utf8'
)
) as {
targets: Record<
string,
{ options: { env?: Record<string, string> } }
>;
};
expect(project.targets[targetName]?.options.env).not.toHaveProperty(
'PORT'
);
expect(project.targets[targetName]?.options.env).toMatchObject({
NODE_ENV: 'development',
});
}
);
});
+28 -3
View File
@@ -17,6 +17,10 @@ export interface ScenarioConfig {
vodDetailsFixture?: 'empty-metadata';
/** Optional fictional release-marketing dataset with local demo artwork. */
marketingFixture?: true;
/** Optional deterministic catalog profile reserved for performance tests. */
performanceFixture?: 'catalog-100k';
/** Build series details on demand instead of during portal initialization. */
deferSeriesDetails?: true;
}
/**
@@ -59,9 +63,23 @@ export const SCENARIOS: Record<string, ScenarioConfig> = {
accountStatus: 'Active',
expiryDate: '2099-12-31',
},
'performance:performance': {
name: 'performance-100k',
description: 'Deterministic local-only 100k performance catalog',
seed: 91001,
categoryCount: { live: 60, vod: 20, series: 20 },
itemsPerCategory: 1000,
seasonsPerSeries: 1,
episodesPerSeason: 1,
accountStatus: 'Active',
expiryDate: '2099-12-31',
performanceFixture: 'catalog-100k',
deferSeriesDetails: true,
},
'series:series': {
name: 'series-heavy',
description: 'Series-heavy — 15 series categories, 6 seasons × 10 episodes',
description:
'Series-heavy — 15 series categories, 6 seasons × 10 episodes',
seed: 2002,
categoryCount: { live: 3, vod: 4, series: 15 },
itemsPerCategory: 30,
@@ -147,11 +165,18 @@ export const SCENARIOS: Record<string, ScenarioConfig> = {
/** Convert credential pair to a numeric seed for unknown credentials. */
export function credentialsToSeed(username: string, password: string): number {
const str = `${username}:${password}`;
return str.split('').reduce((acc, ch) => (acc * 31 + ch.charCodeAt(0)) | 0, 0) >>> 0;
return (
str
.split('')
.reduce((acc, ch) => (acc * 31 + ch.charCodeAt(0)) | 0, 0) >>> 0
);
}
/** Return the scenario config for a given username+password pair. */
export function getScenario(username: string, password: string): ScenarioConfig {
export function getScenario(
username: string,
password: string
): ScenarioConfig {
const key = `${username}:${password}`;
if (SCENARIOS[key]) return SCENARIOS[key];
return {
@@ -0,0 +1,55 @@
import type { Server } from 'node:http';
const SHUTDOWN_DEADLINE_MS = 1_000;
export interface XtreamMockShutdownApplication {
shutdown(): unknown;
}
export function createXtreamMockServerShutdown(
server: Server,
application: XtreamMockShutdownApplication
): () => Promise<void> {
let shutdownPromise: Promise<void> | undefined;
return () => {
shutdownPromise ??= beginShutdown(server, application);
return shutdownPromise;
};
}
function beginShutdown(
server: Server,
application: XtreamMockShutdownApplication
): Promise<void> {
try {
application.shutdown();
} catch {
// HTTP shutdown must remain bounded even if an application hook fails.
}
return closeHttpServer(server);
}
function closeHttpServer(server: Server): Promise<void> {
return new Promise((resolve) => {
let settled = false;
const finish = () => {
if (settled) return;
settled = true;
clearTimeout(deadline);
resolve();
};
const deadline = setTimeout(() => {
server.closeAllConnections();
finish();
}, SHUTDOWN_DEADLINE_MS);
deadline.unref();
try {
server.close(() => finish());
server.closeIdleConnections();
server.closeAllConnections();
} catch {
finish();
}
});
}
@@ -0,0 +1,160 @@
import { type AddressInfo } from 'node:net';
import { createServer, type Server, type ServerResponse } from 'node:http';
import type { PerformanceControlState } from './performance-control.js';
import {
createXtreamMockApp,
createXtreamMockServerShutdown,
} from './server.js';
import {
openAbortableGet,
postJson,
requestJson,
waitFor,
} from './testing/http-server.fixture.js';
jest.mock('@faker-js/faker', () => ({
faker: { seed: jest.fn() },
}));
const TOKEN = 'shutdown-token';
const AUTH_HEADERS = {
'x-iptvnator-performance-token': TOKEN,
};
describe('Xtream mock HTTP shutdown', () => {
it.each([
['barrier', 'blocked'],
['delay', 'delayed'],
] as const)(
'settles an active %s and releases its listener',
async (kind, lifecycleStatus) => {
const app = createXtreamMockApp({
control: { enabled: true, token: TOKEN },
host: '127.0.0.1',
port: 0,
});
const server = createServer(app);
await listen(server);
const address = server.address() as AddressInfo;
const origin = `http://127.0.0.1:${address.port}`;
let heldClient: ReturnType<typeof openAbortableGet> | undefined;
try {
await createRule(origin, kind);
heldClient = openAbortableGet(origin, catalogPath());
await waitFor(async () => {
const snapshot = await state(origin);
expect(
snapshot.ledger.map((entry) => entry.status)
).toContain(lifecycleStatus);
});
let shutdownState: PerformanceControlState | null | undefined;
const applicationShutdown = app.shutdown.bind(app);
app.shutdown = () => {
shutdownState = applicationShutdown();
return shutdownState;
};
const shutdown = createXtreamMockServerShutdown(server, app);
const first = shutdown();
const second = shutdown();
expect(second).toBe(first);
await within(first, 2_000);
await within(heldClient.completed, 2_000);
expect(shutdownState).toMatchObject({
heldIds: [],
heldCount: 0,
rules: [],
occurrences: [],
ledger: [],
});
const rebound = createServer(
(_request, response: ServerResponse) => response.end('ok')
);
try {
await listen(rebound, address.port);
} finally {
await close(rebound);
}
} finally {
heldClient?.abort();
server.closeAllConnections();
if (server.listening) await close(server);
}
}
);
});
async function createRule(
origin: string,
kind: 'barrier' | 'delay'
): Promise<void> {
const result = await postJson(
origin,
kind === 'barrier' ? '/__control/barriers' : '/__control/delays',
TOKEN,
{
id: `shutdown-${kind}`,
match: {
scenario: 'performance-100k',
transport: 'direct',
action: 'get_live_streams',
occurrence: 1,
categoryId: '91101',
},
...(kind === 'delay' ? { milliseconds: 5_000 } : {}),
}
);
expect(result.response.status).toBe(201);
}
async function state(origin: string): Promise<PerformanceControlState> {
const result = await requestJson(origin, '/__control/state', {
headers: AUTH_HEADERS,
});
expect(result.response.status).toBe(200);
return result.body as PerformanceControlState;
}
function catalogPath(): string {
return '/player_api.php?username=performance&password=performance&action=get_live_streams&category_id=91101';
}
function listen(server: Server, port = 0): Promise<void> {
return new Promise((resolve, reject) => {
server.once('error', reject);
server.listen(port, '127.0.0.1', () => {
server.off('error', reject);
resolve();
});
});
}
function close(server: Server): Promise<void> {
return new Promise((resolve) => {
server.close(() => resolve());
server.closeAllConnections();
});
}
async function within<T>(
promise: Promise<T>,
milliseconds: number
): Promise<T> {
let timer: NodeJS.Timeout | undefined;
try {
return await Promise.race([
promise,
new Promise<never>((_resolve, reject) => {
timer = setTimeout(
() => reject(new Error('operation exceeded deadline')),
milliseconds
);
}),
]);
} finally {
if (timer) clearTimeout(timer);
}
}
@@ -0,0 +1,225 @@
import { connect } from 'node:net';
import { resetAll } from './data-store.js';
import {
createXtreamMockApp,
parseXtreamMockServerEnvironment,
} from './server.js';
import { startLoopbackServer } from './testing/http-server.fixture.js';
jest.mock('@faker-js/faker', () => ({
faker: { seed: jest.fn() },
}));
jest.setTimeout(60_000);
describe('Xtream mock server factory', () => {
afterEach(() => {
resetAll();
jest.restoreAllMocks();
});
it('keeps performance controls absent unless explicitly enabled', async () => {
const running = await startLoopbackServer(
createXtreamMockApp({ host: '127.0.0.1', port: 0 })
);
try {
const response = await fetch(`${running.origin}/__control/state`);
expect(response.status).toBe(404);
} finally {
await running.close();
}
});
it('reports the configured port and preserves direct/PWA responses', async () => {
const running = await startLoopbackServer(
createXtreamMockApp({ host: '127.0.0.1', port: 0 })
);
try {
const health = await fetch(`${running.origin}/health`).then((r) =>
r.json()
);
const direct = await fetch(
`${running.origin}/player_api.php?username=performance&password=performance&action=get_live_categories`
);
const directBody = await direct.json();
const proxy = await fetch(
`${running.origin}/xtream?url=http%3A%2F%2Fignored.invalid&username=performance&password=performance&action=get_live_categories`
);
const proxyBody = await proxy.json();
expect(health).toEqual({
status: 'ok',
server: 'xtream-mock-server',
port: 0,
});
expect(direct.status).toBe(200);
expect(directBody).toHaveLength(60);
expect(proxy.status).toBe(200);
expect(proxyBody).toEqual({
payload: directBody,
action: 'get_live_categories',
});
} finally {
await running.close();
}
});
it('preserves the basic M3U fixture and ordinary stream redirects', async () => {
const running = await startLoopbackServer(
createXtreamMockApp({ host: '127.0.0.1', port: 0 })
);
try {
const playlist = await fetch(`${running.origin}/playlist.m3u`);
const stream = await fetch(
`${running.origin}/live/user1/pass1/10000.m3u8`,
{ redirect: 'manual' }
);
expect(playlist.status).toBe(200);
expect(playlist.headers.get('content-type')).toContain(
'audio/x-mpegurl'
);
expect(await playlist.text()).toContain('#EXTM3U');
expect(stream.status).toBe(302);
expect(stream.headers.get('location')).toBe(
'https://test-streams.mux.dev/x36xhzz/x36xhzz.m3u8'
);
} finally {
await running.close();
}
});
it('rejects non-loopback performance binds before opening a listener', () => {
expect(() =>
createXtreamMockApp({
control: { enabled: true, token: 'test-token' },
host: '0.0.0.0',
port: 0,
})
).toThrow(/loopback/i);
});
it('makes local-only performance media routes terminal without outbound requests', async () => {
const running = await startLoopbackServer(
createXtreamMockApp({
control: { enabled: true, token: 'test-token' },
host: '127.0.0.1',
port: 0,
})
);
try {
const httpModule =
jest.requireActual<typeof import('node:http')>('node:http');
const httpsModule =
jest.requireActual<typeof import('node:https')>('node:https');
const httpRequest = jest.spyOn(httpModule, 'request');
const httpsRequest = jest.spyOn(httpsModule, 'request');
const globalFetch = jest.spyOn(globalThis, 'fetch');
const mediaResponses = await Promise.all(
[
'/playlist.m3u',
'/live/performance/performance/1000000.m3u8',
'/live/performance/performance/1000000.ts',
'/movie/performance/performance/2000000.mkv',
'/series/performance/performance/3000000.mkv',
'/timeshift/performance/performance/60/start/1000000.ts',
'/streaming/timeshift.php?username=performance&password=performance',
].map((path) => rawLoopbackGet(running.origin, path))
);
expect(mediaResponses.map(({ status }) => status)).toEqual(
Array(mediaResponses.length).fill(410)
);
expect(
mediaResponses.every(
({ headers }) => !headers.includes('location:')
)
).toBe(true);
expect(httpRequest).not.toHaveBeenCalled();
expect(httpsRequest).not.toHaveBeenCalled();
expect(globalFetch).not.toHaveBeenCalled();
} finally {
await running.close();
}
});
});
describe('Xtream mock environment parsing', () => {
it('uses safe defaults and enables control only for the exact flag', () => {
expect(parseXtreamMockServerEnvironment({})).toEqual({
host: '0.0.0.0',
port: 3211,
});
expect(
parseXtreamMockServerEnvironment({
IPTVNATOR_XTREAM_MOCK_CONTROL: '1',
IPTVNATOR_XTREAM_MOCK_CONTROL_TOKEN: 'secret-token',
})
).toEqual({
control: { enabled: true, token: 'secret-token' },
host: '127.0.0.1',
port: 3211,
});
expect(
parseXtreamMockServerEnvironment({
HOST: '::1',
PORT: '4321',
IPTVNATOR_XTREAM_MOCK_CONTROL: '1',
IPTVNATOR_XTREAM_MOCK_CONTROL_TOKEN: 'secret-token',
})
).toEqual({
control: { enabled: true, token: 'secret-token' },
host: '::1',
port: 4321,
});
expect(
parseXtreamMockServerEnvironment({
IPTVNATOR_XTREAM_MOCK_CONTROL: 'true',
IPTVNATOR_XTREAM_MOCK_CONTROL_TOKEN: 'ignored',
})
).toEqual({ host: '0.0.0.0', port: 3211 });
});
it.each([
[{ PORT: '-1' }, /port/i],
[{ PORT: '12x' }, /port/i],
[{ PORT: '65536' }, /port/i],
[{ IPTVNATOR_XTREAM_MOCK_CONTROL: '1' }, /token/i],
[
{
HOST: '192.0.2.1',
IPTVNATOR_XTREAM_MOCK_CONTROL: '1',
IPTVNATOR_XTREAM_MOCK_CONTROL_TOKEN: 'secret-token',
},
/loopback/i,
],
])('rejects invalid environment %#', (environment, message) => {
expect(() => parseXtreamMockServerEnvironment(environment)).toThrow(
message
);
});
});
async function rawLoopbackGet(
origin: string,
path: string
): Promise<{ headers: string; status: number }> {
const url = new URL(origin);
return new Promise((resolve, reject) => {
const socket = connect(Number(url.port), url.hostname);
let raw = '';
socket.setEncoding('utf8');
socket.once('error', reject);
socket.on('data', (chunk) => (raw += chunk));
socket.once('end', () => {
const [headers = ''] = raw.split('\r\n\r\n');
const status = Number(headers.split(' ')[1]);
resolve({ headers: headers.toLowerCase(), status });
});
socket.once('connect', () => {
socket.write(
`GET ${path} HTTP/1.1\r\nHost: ${url.host}\r\nConnection: close\r\n\r\n`
);
});
});
}
+274
View File
@@ -0,0 +1,274 @@
import { existsSync, readFileSync } from 'node:fs';
import { join } from 'node:path';
import cors from 'cors';
import express, {
type NextFunction,
type Request,
type Response,
} from 'express';
import { resetAll } from './data-store.js';
import { renderMarketingAssetSvg } from './generators/marketing.generator.js';
import { installPerformanceControlRoutes } from './performance-control-routes.js';
import {
type PerformanceControlState,
XtreamPerformanceController,
} from './performance-control.js';
import { dispatchAction } from './routes/dispatch.js';
export { createXtreamMockServerShutdown } from './server-lifecycle.js';
export interface XtreamMockServerOptions {
readonly control?: {
readonly enabled: boolean;
readonly token: string;
};
readonly host: string;
readonly port: number;
}
export interface XtreamMockApplication extends express.Express {
shutdown(): PerformanceControlState | null;
}
const M3U_FIXTURE = `#EXTM3U
#EXTINF:0 tvg-id="1" tvg-logo="http://channel.icons.url/img/1.png" group-title="News", Channel 1
https://example.channels/path-to-file/1.m3u8
#EXTINF:0 tvg-id="2" tvg-logo="http://channel.icons.url/img/2.png" group-title="News", Positive News TV
https://example.channels/path-to-file/2.m3u8
#EXTINF:0 tvg-id="3" tvg-logo="http://channel.icons.url/img/3.png" group-title="Sport", Sport TVX
https://example.channels/path-to-file/3.m3u8
#EXTINF:0 tvg-id="4" tvg-logo="http://channel.icons.url/img/4.png" group-title="Kids", HappyKids TV
https://example.channels/path-to-file/4.m3u8
`;
const HLS_STUB = 'https://test-streams.mux.dev/x36xhzz/x36xhzz.m3u8';
const DEFAULT_PORT = 3211;
const DEFAULT_NORMAL_HOST = '0.0.0.0';
const DEFAULT_CONTROL_HOST = '127.0.0.1';
const PERFORMANCE_USERNAME = 'performance';
const PERFORMANCE_PASSWORD = 'performance';
const marketingRasterAssetRoot = join(
process.cwd(),
'apps/xtream-mock-server/public/marketing'
);
export function createXtreamMockApp(
options: XtreamMockServerOptions
): XtreamMockApplication {
validateServerOptions(options);
const app = express() as XtreamMockApplication;
const controlEnabled = options.control?.enabled === true;
const controller = controlEnabled
? new XtreamPerformanceController()
: undefined;
app.shutdown = () => controller?.shutdown() ?? null;
if (controller && options.control) {
installPerformanceControlRoutes(app, controller, options.control.token);
}
app.use(cors());
app.use(express.json());
app.get('/health', (_request, response) => {
response.json({
status: 'ok',
server: 'xtream-mock-server',
port: options.port,
});
});
app.post('/reset', (_request, response) => {
if (controlEnabled) {
response
.status(410)
.json({ error: 'use-performance-control-reset' });
return;
}
resetAll();
response.json({ status: 'reset' });
});
app.get('/playlist.m3u', (_request, response) => {
if (controlEnabled) {
response.status(410).json({ error: 'performance-media-disabled' });
return;
}
response
.type('audio/x-mpegurl')
.set('Cache-Control', 'no-store')
.send(M3U_FIXTURE);
});
installMarketingAssetRoute(app);
app.get('/player_api.php', (request, response, next) => {
dispatchWithControl(controller, request, response, next, 'direct', () =>
dispatchAction(request, response)
);
});
installStreamRoutes(app, controlEnabled);
app.get('/xtream', (request, response, next) => {
dispatchWithControl(controller, request, response, next, 'proxy', () =>
dispatchProxyAction(request, response)
);
});
return app;
}
export function parseXtreamMockServerEnvironment(
environment: NodeJS.ProcessEnv
): XtreamMockServerOptions {
const rawPort = environment['PORT'] ?? String(DEFAULT_PORT);
if (!/^\d+$/.test(rawPort)) {
throw new Error('Xtream mock port must be an integer');
}
const port = Number(rawPort);
const controlEnabled = environment['IPTVNATOR_XTREAM_MOCK_CONTROL'] === '1';
const host =
environment['HOST'] ??
(controlEnabled ? DEFAULT_CONTROL_HOST : DEFAULT_NORMAL_HOST);
const token = environment['IPTVNATOR_XTREAM_MOCK_CONTROL_TOKEN'] ?? '';
const options: XtreamMockServerOptions = {
...(controlEnabled
? { control: { enabled: true, token } as const }
: {}),
host,
port,
};
validateServerOptions(options);
return options;
}
function validateServerOptions(options: XtreamMockServerOptions): void {
if (
!Number.isInteger(options.port) ||
options.port < 0 ||
options.port > 65_535
) {
throw new Error('Xtream mock port must be in 0..65535');
}
if (!options.host) throw new Error('Xtream mock host is required');
if (!options.control?.enabled) return;
if (options.control.token.trim().length === 0) {
throw new Error('Xtream performance control token is required');
}
if (options.host !== '127.0.0.1' && options.host !== '::1') {
throw new Error(
'Xtream performance control requires a loopback bind host'
);
}
}
function dispatchWithControl(
controller: XtreamPerformanceController | undefined,
request: Request,
response: Response,
next: NextFunction,
transport: 'direct' | 'proxy',
dispatch: () => void
): void {
if (!controller) {
dispatch();
return;
}
void controller
.intercept(request, response, transport, dispatch)
.catch(next);
}
function dispatchProxyAction(request: Request, response: Response): void {
const { url, action, ...rest } = request.query as Record<string, string>;
void url;
const syntheticRequest = {
query: { action, ...rest },
headers: request.headers,
params: {},
} as unknown as Request;
let payload: unknown;
let statusCode = 200;
const syntheticResponse = {
json(data: unknown) {
payload = data;
return this;
},
status(code: number) {
statusCode = code;
return this;
},
} as unknown as Response;
dispatchAction(syntheticRequest, syntheticResponse);
if (statusCode !== 200) {
response.status(statusCode).json({ error: payload });
return;
}
response.json({ payload, action });
}
function installMarketingAssetRoute(app: express.Express): void {
app.get('/assets/marketing/:kind/:slug', (request, response) => {
const kind = request.params['kind'] as
'backdrop' | 'episode' | 'logo' | 'poster' | 'season';
const slug = request.params['slug'] ?? '';
const size =
typeof request.query['size'] === 'string'
? request.query['size']
: undefined;
if (
!['backdrop', 'episode', 'logo', 'poster', 'season'].includes(kind)
) {
response.status(404).send('Unknown marketing asset kind');
return;
}
const rasterSlug = slug.replace(/\.(svg|png)$/i, '');
const rasterPath = join(
marketingRasterAssetRoot,
kind,
`${rasterSlug}.png`
);
if (existsSync(rasterPath)) {
response
.type('image/png')
.set('Cache-Control', 'public, max-age=3600')
.send(readFileSync(rasterPath));
return;
}
response
.type('image/svg+xml')
.set('Cache-Control', 'public, max-age=3600')
.send(renderMarketingAssetSvg(kind, slug, size));
});
}
function installStreamRoutes(
app: express.Express,
controlEnabled: boolean
): void {
const streamResponse = (request: Request, response: Response) => {
if (isPerformanceMediaRequest(request, controlEnabled)) {
response.status(410).json({ error: 'performance-media-disabled' });
return;
}
response.redirect(HLS_STUB);
};
app.get('/live/:username/:password/:streamId.m3u8', streamResponse);
app.get('/live/:username/:password/:streamId.ts', streamResponse);
app.get('/movie/:username/:password/:streamId.:ext', streamResponse);
app.get('/series/:username/:password/:streamId.:ext', streamResponse);
app.all(
'/timeshift/:username/:password/:duration/:start/:streamId.ts',
streamResponse
);
app.all('/streaming/timeshift.php', streamResponse);
}
function isPerformanceMediaRequest(
request: Request,
controlEnabled: boolean
): boolean {
if (!controlEnabled) return false;
const username =
request.params['username'] ?? request.query['username'] ?? '';
const password =
request.params['password'] ?? request.query['password'] ?? '';
return (
username === PERFORMANCE_USERNAME && password === PERFORMANCE_PASSWORD
);
}
@@ -0,0 +1,126 @@
import { type AddressInfo } from 'node:net';
import {
createServer,
get as httpGet,
type Server,
type ServerResponse,
} from 'node:http';
import type express from 'express';
export interface RunningTestServer {
readonly origin: string;
close(): Promise<void>;
}
export interface AbortableRequest {
readonly completed: Promise<void>;
abort(): void;
}
export async function startLoopbackServer(
app: express.Express
): Promise<RunningTestServer> {
const server = createServer(app);
await new Promise<void>((resolve, reject) => {
server.once('error', reject);
server.listen(0, '127.0.0.1', () => {
server.off('error', reject);
resolve();
});
});
const address = server.address() as AddressInfo;
return {
origin: `http://127.0.0.1:${address.port}`,
close: () => closeServer(server),
};
}
export async function requestJson(
origin: string,
path: string,
init: RequestInit = {}
): Promise<{ body: unknown; response: Response }> {
const response = await fetch(`${origin}${path}`, init);
const body = await response.json();
return { body, response };
}
export async function postJson(
origin: string,
path: string,
token: string,
body?: unknown
): Promise<{ body: unknown; response: Response }> {
return requestJson(origin, path, {
method: 'POST',
headers: {
'content-type': 'application/json',
'x-iptvnator-performance-token': token,
},
...(body === undefined ? {} : { body: JSON.stringify(body) }),
});
}
export function openAbortableGet(
origin: string,
path: string,
headers: Record<string, string> = {}
): AbortableRequest {
let settled = false;
let settle!: () => void;
const completed = new Promise<void>((resolve) => {
settle = resolve;
});
const request = httpGet(`${origin}${path}`, { headers }, (response) => {
drain(response, () => {
settled = true;
settle();
});
});
request.once('error', () => {
settled = true;
settle();
});
return {
completed,
abort() {
if (!settled) request.destroy();
},
};
}
export async function waitFor(
assertion: () => Promise<void>,
timeoutMs = 2_000
): Promise<void> {
const deadline = Date.now() + timeoutMs;
let lastError: unknown;
while (Date.now() < deadline) {
try {
await assertion();
return;
} catch (error) {
lastError = error;
await new Promise((resolve) => setTimeout(resolve, 5));
}
}
throw lastError;
}
function closeServer(server: Server): Promise<void> {
return new Promise((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
server.closeAllConnections();
});
}
function drain(
response: ServerResponse | NodeJS.ReadableStream,
done: () => void
) {
response.on('data', () => undefined);
response.once('end', done);
response.once('error', done);
}
+65 -193
View File
@@ -1,201 +1,73 @@
import http from 'http';
import { existsSync, readFileSync } from 'fs';
import { join } from 'path';
import express from 'express';
import cors from 'cors';
import { dispatchAction } from './app/routes/dispatch.js';
import { resetAll } from './app/data-store.js';
import { renderMarketingAssetSvg } from './app/generators/marketing.generator.js';
import { createServer } from 'node:http';
import {
createXtreamMockServerShutdown,
createXtreamMockApp,
parseXtreamMockServerEnvironment,
} from './app/server.js';
const app = express();
const PORT = parseInt(process.env['PORT'] ?? '3211', 10);
const M3U_FIXTURE = `#EXTM3U
#EXTINF:0 tvg-id="1" tvg-logo="http://channel.icons.url/img/1.png" group-title="News", Channel 1
https://example.channels/path-to-file/1.m3u8
#EXTINF:0 tvg-id="2" tvg-logo="http://channel.icons.url/img/2.png" group-title="News", Positive News TV
https://example.channels/path-to-file/2.m3u8
#EXTINF:0 tvg-id="3" tvg-logo="http://channel.icons.url/img/3.png" group-title="Sport", Sport TVX
https://example.channels/path-to-file/3.m3u8
#EXTINF:0 tvg-id="4" tvg-logo="http://channel.icons.url/img/4.png" group-title="Kids", HappyKids TV
https://example.channels/path-to-file/4.m3u8
`;
const marketingRasterAssetRoot = join(
process.cwd(),
'apps/xtream-mock-server/public/marketing'
);
app.use(cors());
app.use(express.json());
// ─── Health check ──────────────────────────────────────────────────────────────
app.get('/health', (_req, res) => {
res.json({ status: 'ok', server: 'xtream-mock-server', port: PORT });
});
// ─── Reset all cached data ─────────────────────────────────────────────────────
app.post('/reset', (_req, res) => {
resetAll();
res.json({ status: 'reset' });
});
// ─── M3U fixture endpoint for self-hosted PWA URL import tests ───────────────
app.get('/playlist.m3u', (_req, res) => {
res.type('audio/x-mpegurl')
.set('Cache-Control', 'no-store')
.send(M3U_FIXTURE);
});
// ─── Local fictional artwork for release screenshots ──────────────────────────
app.get('/assets/marketing/:kind/:slug', (req, res) => {
const kind = req.params['kind'] as
| 'backdrop'
| 'episode'
| 'logo'
| 'poster'
| 'season';
const slug = req.params['slug'] ?? '';
const size =
typeof req.query['size'] === 'string' ? req.query['size'] : undefined;
if (!['backdrop', 'episode', 'logo', 'poster', 'season'].includes(kind)) {
res.status(404).send('Unknown marketing asset kind');
return;
function start(): void {
let options;
try {
options = parseXtreamMockServerEnvironment(process.env);
} catch (error) {
const message =
error instanceof Error ? error.message : 'Invalid configuration';
console.error(`[xtream-mock] ${message}`);
process.exit(1);
}
const rasterSlug = slug.replace(/\.(svg|png)$/i, '');
const rasterPath = join(
marketingRasterAssetRoot,
kind,
`${rasterSlug}.png`
);
const application = createXtreamMockApp(options);
const server = createServer(application);
const shutdownServer = createXtreamMockServerShutdown(server, application);
const displayHost = options.host === '::1' ? '[::1]' : options.host;
let shuttingDown = false;
if (existsSync(rasterPath)) {
res.type('image/png')
.set('Cache-Control', 'public, max-age=3600')
.send(readFileSync(rasterPath));
return;
}
server.on('error', (error: NodeJS.ErrnoException) => {
if (error.code === 'EADDRINUSE') {
console.error(
`[xtream-mock] Port ${options.port} is already in use.`
);
} else {
console.error('[xtream-mock] Server error:', error.message);
}
process.exit(1);
});
res.type('image/svg+xml')
.set('Cache-Control', 'public, max-age=3600')
.send(renderMarketingAssetSvg(kind, slug, size));
});
server.listen(options.port, options.host, () => {
const origin = `http://${displayHost}:${options.port}`;
console.log(`[xtream-mock] Listening on ${origin}`);
console.log(`[xtream-mock] Direct API: ${origin}/player_api.php`);
console.log(`[xtream-mock] PWA proxy: ${origin}/xtream`);
console.log(`[xtream-mock] Health: ${origin}/health`);
if (options.control?.enabled) {
console.log('[xtream-mock] Performance control enabled');
}
});
// ─── Direct Xtream player_api.php endpoint ─────────────────────────────────────
app.get('/player_api.php', dispatchAction);
const shutdown = () => {
if (!shuttingDown) {
shuttingDown = true;
console.log('\n[xtream-mock] Shutting down...');
}
void shutdownServer().then(() => process.exit(0));
};
// ─── Stream stub endpoints (Xtream stream URLs reference these) ────────────────
// Returns a valid HLS playlist redirect for any stream type.
// In production these would be actual video streams.
const HLS_STUB = 'https://test-streams.mux.dev/x36xhzz/x36xhzz.m3u8';
process.on('SIGINT', shutdown);
process.on('SIGTERM', shutdown);
process.on('uncaughtException', (error) => {
console.error('[xtream-mock] Uncaught exception:', error.message);
process.exit(1);
});
process.on('unhandledRejection', (reason) => {
const message =
reason instanceof Error ? reason.message : 'Unknown rejection';
console.error('[xtream-mock] Unhandled rejection:', message);
process.exit(1);
});
process.stdin.resume();
process.stdin.on('end', () => {
/* The HTTP server handle intentionally keeps this process alive. */
});
}
app.get('/live/:username/:password/:streamId.m3u8', (_req, res) => {
res.redirect(HLS_STUB);
});
app.get('/live/:username/:password/:streamId.ts', (_req, res) => {
res.redirect(HLS_STUB);
});
app.get('/movie/:username/:password/:streamId.:ext', (_req, res) => {
res.redirect(HLS_STUB);
});
app.get('/series/:username/:password/:streamId.:ext', (_req, res) => {
res.redirect(HLS_STUB);
});
app.all(
'/timeshift/:username/:password/:duration/:start/:streamId.ts',
(_req, res) => {
res.redirect(HLS_STUB);
}
);
app.all('/streaming/timeshift.php', (_req, res) => {
res.redirect(HLS_STUB);
});
// ─── PWA CORS proxy endpoint ────────────────────────────────────────────────────
// IPTVnator PWA routes Xtream calls through:
// GET /xtream?url=<serverUrl>&action=<action>&username=X&password=Y
// and expects: { payload: <data>, action: <action> }
app.get('/xtream', async (req, res) => {
const { url: _url, action, ...rest } = req.query as Record<string, string>;
const syntheticReq = {
query: { action, ...rest },
headers: req.headers,
params: {},
} as unknown as express.Request;
// Capture response via monkey-patched json method
let payload: unknown;
let statusCode = 200;
const syntheticRes = {
json(data: unknown) {
payload = data;
return this;
},
status(code: number) {
statusCode = code;
return this;
},
} as unknown as express.Response;
dispatchAction(syntheticReq, syntheticRes);
if (statusCode !== 200) {
res.status(statusCode).json({ error: payload });
return;
}
res.json({ payload, action });
});
// ─── Server lifecycle ──────────────────────────────────────────────────────────
const server = http.createServer(app);
server.on('error', (err: NodeJS.ErrnoException) => {
if (err.code === 'EADDRINUSE') {
console.error(`[xtream-mock] Port ${PORT} is already in use.`);
} else {
console.error('[xtream-mock] Server error:', err.message);
}
process.exit(1);
});
server.listen(PORT, () => {
console.log(`[xtream-mock] Listening on http://localhost:${PORT}`);
console.log(
`[xtream-mock] Direct API: http://localhost:${PORT}/player_api.php?username=user1&password=pass1&action=get_account_info`
);
console.log(
`[xtream-mock] PWA proxy: http://localhost:${PORT}/xtream?url=http://localhost:${PORT}&username=user1&password=pass1&action=get_account_info`
);
console.log(`[xtream-mock] Health: http://localhost:${PORT}/health`);
});
const shutdown = () => {
console.log('\n[xtream-mock] Shutting down...');
server.close(() => process.exit(0));
};
process.on('SIGINT', shutdown);
process.on('SIGTERM', shutdown);
process.on('uncaughtException', (err) => {
console.error('[xtream-mock] Uncaught exception:', err);
process.exit(1);
});
process.on('unhandledRejection', (reason) => {
console.error('[xtream-mock] Unhandled rejection:', reason);
process.exit(1);
});
// When Nx (or any process manager) closes stdin, prevent auto-exit.
// The HTTP server handle is what keeps the process alive.
process.stdin.resume();
process.stdin.on('end', () => {
/* ignore stdin close */
});
start();
@@ -0,0 +1,17 @@
{
"extends": "../../tsconfig.base.json",
"compilerOptions": {
"outDir": "../../dist/out-tsc",
"target": "ES2022",
"module": "commonjs",
"moduleResolution": "node10",
"types": ["jest", "node"],
"strict": true
},
"include": [
"jest.config.ts",
"src/**/*.test.ts",
"src/**/*.spec.ts",
"src/**/*.d.ts"
]
}