fix(web-backend): validate and pin provider redirect hops (#1553)

* fix(web-backend): validate and pin every provider redirect hop

* fix(web-backend): separate provider metadata from connection authority
This commit is contained in:
4gray authored and GitHub committed 2026-09-06 08:59:28 +02:00
1 parent 9de480826c
commit 5febe28eba
24 files changed
+1694 -331

No files matched your search

+18 -6
View File
@@ -13,9 +13,9 @@
* see `@iptvnator/shared/host-health`) fast-fails the rest once a host has
* refused to answer twice in a row.
*
* On the axios timeout semantics: with the default (follow-redirects)
* transport, `timeout` is NOT a wall-clock deadline for the whole response. It
* bounds the time to response headers, and then continues as the socket's
* On the axios timeout semantics: with axios 1.20's native HTTP
* transport (`maxRedirects: 0`), each hop's `timeout` is NOT a wall-clock
* deadline for the whole response. It bounds the time to response headers, and then continues as the socket's
* inactivity timeout for the body. A large XMLTV or M3U download that keeps
* delivering bytes is therefore never cut off mid-transfer — only a stalled one
* is, which is exactly the case these numbers are meant to catch.
@@ -29,6 +29,7 @@ import {
portalEndpointKeyOf,
} from '@iptvnator/shared/host-health';
import { buildHostConnectivityFastFailMessage } from '@iptvnator/shared/interfaces';
import { ProviderRequestError } from './provider-request-error';
import type { NormalizedProviderError } from './provider-error';
/**
@@ -176,8 +177,8 @@ export function resetProviderHost(
* `countFailures: false` is for requests exempt from the guard (endpoint
* discovery). Their failures are expected and must not count — but the report
* must still happen, because an error that carries an HTTP response proves the
* endpoint answered. This route sets no `validateStatus`, so axios rejects every
* non-2xx WITH `error.response`; dropping those instead of reporting them is
* endpoint answered. The redirect wrapper accepts 3xx, but axios still rejects
* other non-2xx WITH `error.response`; dropping those instead of reporting them is
* what would let the breaker open in the middle of discovery.
*/
export function reportProviderRequestFailure(
@@ -190,6 +191,14 @@ export function reportProviderRequestFailure(
return;
}
const manualChain = error instanceof ProviderRequestError;
if (error instanceof ProviderRequestError) {
if (error.initialResponded) {
guard.reportSuccess(token);
return;
}
error = error.cause;
}
const countFailures = options.countFailures ?? true;
switch (classifyHostRequestFailure(error)) {
case 'host-level':
@@ -199,7 +208,10 @@ export function reportProviderRequestFailure(
// other response — otherwise an ordinary timeout, a probe that was
// redirected to a dead destination, and another ordinary timeout
// still read as two consecutive failures.
if (failedAfterRedirect(error, token, options.requestUrl)) {
if (
!manualChain &&
failedAfterRedirect(error, token, options.requestUrl)
) {
guard.reportSuccess(token);
break;
}
@@ -0,0 +1,108 @@
import axios from 'axios';
import {
request as httpRequest,
ClientRequest,
IncomingMessage,
RequestOptions,
} from 'node:http';
import {
request as httpsRequest,
Agent as HttpsAgent,
AgentOptions,
} from 'node:https';
import { isIP } from 'node:net';
import { checkServerIdentity } from 'node:tls';
import type {
ProviderTransportOptions,
WebBackendHttpClient,
WebBackendHttpResponse,
} from './validated-http-client';
/**
* Keep HTTP authority metadata separate from socket destination selection.
* Only the hop's pinned agent can select an address; axios receives a fixed
* logical hostname, never a user-selected connection authority. Host and TLS
* identity still describe the validated provider, including IP certificates.
*/
export class ProviderAxiosTransport implements WebBackendHttpClient {
async get<T>(
url: string,
options: ProviderTransportOptions = {}
): Promise<WebBackendHttpResponse<T>> {
const target = new URL(axios.getUri({ url, params: options.params }));
const secure = target.protocol === 'https:';
const hostname = target.hostname.replace(/^\[|\]$/g, '');
const port = Number(target.port || (secure ? 443 : 80));
const root = secure
? 'https://provider.invalid/'
: 'http://provider.invalid/';
if (!options.httpAgent || !options.httpsAgent) {
throw new Error('Provider transport requires pinned agents');
}
const agent = options.httpsAgent as HttpsAgent & {
options: AgentOptions;
};
agent.options.servername = isIP(hostname) ? '' : hostname;
agent.options.checkServerIdentity = (_name, certificate) =>
checkServerIdentity(hostname, certificate);
const response = await axios.get<T>(root, {
...options,
params: undefined,
transport: {
request: (
requestOptions: RequestOptions,
callback: (response: IncomingMessage) => void
) => {
// Path is HTTP metadata, never parsed as an authority.
// Even //host/path remains a path on the pinned connection.
const request = (secure ? httpsRequest : httpRequest)(
{
...requestOptions,
hostname: 'provider.invalid',
port,
path: target.pathname + target.search,
},
callback
);
retainHeaderTimeout(request, options.timeout);
return request;
},
},
headers: { ...options.headers, Host: target.host },
});
return {
...response,
headers: {
location:
typeof response.headers['location'] === 'string'
? response.headers['location']
: undefined,
},
};
}
}
/** Axios identifies native transports by identity; a custom one needs this timer. */
function retainHeaderTimeout(request: ClientRequest, timeout?: number): void {
if (!timeout) return;
const timer = setTimeout(
() =>
request.destroy(
Object.assign(
new Error('Provider response headers timed out'),
{ code: 'ECONNABORTED' }
)
),
timeout
);
timer.unref();
const clear = () => {
clearTimeout(timer);
request.removeListener('response', clear);
request.removeListener('error', clear);
request.removeListener('close', clear);
};
request.once('response', clear);
request.once('error', clear);
request.once('close', clear);
}
@@ -9,6 +9,9 @@
* routinely holds Xtream credentials.
*/
import { ProviderRequestError } from './provider-request-error';
import { providerUrlErrorBody } from './provider-url-policy';
export interface ProviderError extends Error {
readonly code?: unknown;
readonly cause?: unknown;
@@ -78,6 +81,10 @@ function visitErrorNode(
export function normalizeProviderError(
error: unknown
): NormalizedProviderError {
if (error instanceof ProviderRequestError) {
if (error.policyError) return providerUrlErrorBody(error.policyError);
return normalizeProviderError(error.cause);
}
const providerError = error as ProviderError | null | undefined;
const response = providerError?.response;
if (
@@ -111,6 +118,10 @@ export function logProviderRequestFailure(options: {
}
function describeProviderFailure(error: unknown): string {
if (error instanceof ProviderRequestError)
return error.policyError
? `URL policy ${error.policyError.status}`
: describeProviderFailure(error.cause);
const providerError = error as ProviderError | null | undefined;
const response = providerError?.response;
if (
@@ -0,0 +1,12 @@
import type { ProviderUrlError } from './provider-url-policy';
/** Internal chain evidence, separate from the deliberately small public body. */
export class ProviderRequestError extends Error {
constructor(
readonly initialResponded: boolean,
readonly cause: unknown,
readonly policyError?: ProviderUrlError
) {
super('Provider request failed');
}
}
@@ -0,0 +1,51 @@
import { Socket } from 'node:net';
import { Readable } from 'node:stream';
/** Redirect bodies are irrelevant, including invalid gzip and endless bodies. */
export function discardProviderBody(data: unknown): void {
if (data instanceof Readable) data.destroy();
}
export function discardProviderErrorBody(error: unknown): void {
const response = (error as { response?: { data?: unknown } } | null)
?.response;
discardProviderBody(response?.data);
}
/** Match axios's buffered default and arraybuffer modes for the final hop. */
export async function readProviderBody<T>(
data: T,
arraybuffer: boolean,
timeout?: number,
socket?: Socket
): Promise<T> {
if (!(data instanceof Readable)) return data;
const chunks: Buffer[] = [];
// Axios resolves stream responses at headers and then ignores its own
// request timeout callback. Keep the native socket's inactivity timer
// alive until the body finishes; arriving wire bytes reset it, even if
// decompression has not emitted another decoded chunk yet.
const onTimeout = () =>
data.destroy(
Object.assign(new Error('Provider response body timed out'), {
code: 'ECONNABORTED',
})
);
if (timeout && socket) socket.setTimeout(timeout, onTimeout);
try {
for await (const chunk of data) chunks.push(Buffer.from(chunk));
} finally {
if (timeout && socket) {
socket.removeListener('timeout', onTimeout);
socket.setTimeout(0);
}
}
const buffer = Buffer.concat(chunks);
if (arraybuffer) return buffer as T;
const text = buffer.toString('utf8').replace(/^\uFEFF/, '');
try {
return JSON.parse(text) as T;
} catch {
return text as T;
}
}
@@ -0,0 +1,117 @@
import { validateProviderUrl } from './provider-url-policy';
const publicV4 = '93.184.216.34';
const publicV6 = '2606:4700:4700::1111';
const policy = {
allowPrivateNetworkTargets: false,
resolveHostname: async () => [publicV4],
};
const blocked = [
'0.0.0.0',
'10.1.2.3',
'100.64.0.1',
'127.0.0.1',
'169.254.169.254',
'172.16.0.1',
'192.168.0.1',
'192.0.2.1',
'198.18.0.1',
'198.51.100.1',
'203.0.113.1',
'224.0.0.1',
'255.255.255.255',
'::',
'::1',
'fc00::1',
'fd12::1',
'fe80::1',
'febf::1',
'ff02::1',
'::ffff:127.0.0.1',
'::ffff:7f00:1',
'::ffff:a00:1',
'0:0:0:0:0:ffff:a9fe:a9fe',
'64:ff9b::a00:1',
'2001:db8::1',
'2001::1',
'2002:7f00:1::',
'3fff::1',
];
describe('provider URL policy', () => {
it.each(blocked)('rejects literal and DNS answer %s', async (address) => {
const url = `http://${address.includes(':') ? `[${address}]` : address}/`;
await expect(validateProviderUrl(url, policy)).resolves.toMatchObject({
status: 400,
});
await expect(
validateProviderUrl('https://provider.example', {
...policy,
resolveHostname: async () => [publicV4, address],
})
).resolves.toMatchObject({ status: 400 });
});
it.each([
'ftp://host/',
'file:///etc/passwd',
'data:text/plain,hello',
'http://user:secret@host',
'invalid',
'http://localhost.',
'http://sub.localhost',
])('rejects %s', async (url) => {
await expect(validateProviderUrl(url, policy)).resolves.toMatchObject({
status: 400,
});
});
it.each([
[],
['bad-address'],
[publicV4, 'host.example'],
['127.0.0.1%zone'],
])(
'rejects malformed DNS records %j even with LAN opt-in',
async (...addresses) => {
// Jest spreads array table rows; collect the records back into a list.
for (const allowPrivateNetworkTargets of [false, true]) {
await expect(
validateProviderUrl('https://provider.example', {
allowPrivateNetworkTargets,
resolveHostname: async () => addresses as string[],
})
).resolves.toMatchObject({ status: 400 });
}
}
);
it.each([publicV4, publicV6, '::ffff:5db8:d822'])(
'accepts public address %s and retains the hostname',
async (address) => {
const result = await validateProviderUrl(
'https://provider.example/path',
{
...policy,
resolveHostname: async () => [address],
}
);
expect(result).toEqual({
url: new URL('https://provider.example/path'),
addresses: [address],
});
}
);
it('resolves and pins trusted LAN hosts while still rejecting schemes and credentials', async () => {
const trusted = {
allowPrivateNetworkTargets: true,
resolveHostname: async () => ['127.0.0.1', '::1'],
};
await expect(
validateProviderUrl('https://lan.example', trusted)
).resolves.toMatchObject({ addresses: ['127.0.0.1', '::1'] });
await expect(
validateProviderUrl('https://user:secret@lan.example', trusted)
).resolves.toMatchObject({ status: 400 });
await expect(
validateProviderUrl('file:///tmp/test', trusted)
).resolves.toMatchObject({ status: 400 });
});
});
@@ -0,0 +1,140 @@
import { lookup } from 'node:dns/promises';
import { BlockList, isIP } from 'node:net';
export interface ProviderUrlPolicy {
readonly allowPrivateNetworkTargets: boolean;
readonly resolveHostname: (hostname: string) => Promise<readonly string[]>;
}
export interface ProviderUrlError {
readonly message: string;
readonly status: number;
readonly lookupError?: unknown;
}
export interface ValidatedProviderTarget {
readonly url: URL;
readonly addresses: readonly string[];
}
export function providerUrlErrorBody(error: ProviderUrlError): {
message: string;
status: number;
} {
return { message: error.message, status: error.status };
}
export async function resolveHostname(
hostname: string
): Promise<readonly string[]> {
const records = await lookup(hostname, { all: true, verbatim: true });
return records.map((record) => record.address);
}
export async function validateProviderUrl(
rawUrl: string,
policy: ProviderUrlPolicy
): Promise<ValidatedProviderTarget | ProviderUrlError> {
let url: URL;
try {
url = new URL(rawUrl);
} catch {
return { message: 'Provider URL is not a valid URL', status: 400 };
}
if (url.protocol !== 'http:' && url.protocol !== 'https:') {
return {
message: 'Only http and https provider URLs are supported',
status: 400,
};
}
if (url.username || url.password) {
return {
message: 'Provider URL credentials are not supported',
status: 400,
};
}
const hostname = url.hostname.replace(/^\[|\]$/g, '').toLowerCase();
const local = hostname.replace(/\.$/, '');
if (
!policy.allowPrivateNetworkTargets &&
(local === 'localhost' ||
local.endsWith('.localhost') ||
(isIP(hostname) !== 0 && !isPublicAddress(hostname)))
) {
return privateAddressError();
}
let addresses: readonly string[];
try {
addresses = isIP(hostname)
? [hostname]
: await policy.resolveHostname(hostname);
} catch (lookupError) {
return {
message: 'Provider URL host could not be resolved',
status: 400,
lookupError,
};
}
// Validate every answer; never let a malformed record trigger a second DNS
// lookup in the transport. Even trusted LAN mode requires concrete IPs.
if (!addresses.length || addresses.some((address) => !isIP(address))) {
return {
message: 'Provider URL host could not be resolved',
status: 400,
};
}
if (
!policy.allowPrivateNetworkTargets &&
addresses.some((address) => !isPublicAddress(address))
) {
return privateAddressError();
}
return { url, addresses: [...addresses] };
}
function privateAddressError(): ProviderUrlError {
return {
message: 'Provider URL points to a private or local network address',
status: 400,
};
}
const blockedV4 = new BlockList();
for (const [network, prefix] of [
['0.0.0.0', 8],
['10.0.0.0', 8],
['100.64.0.0', 10],
['127.0.0.0', 8],
['169.254.0.0', 16],
['172.16.0.0', 12],
['192.0.0.0', 24],
['192.0.2.0', 24],
['192.88.99.0', 24],
['192.168.0.0', 16],
['198.18.0.0', 15],
['198.51.100.0', 24],
['203.0.113.0', 24],
['224.0.0.0', 3],
] as const)
blockedV4.addSubnet(network, prefix, 'ipv4');
const globalV6 = new BlockList();
globalV6.addSubnet('2000::', 3, 'ipv6');
const blockedV6 = new BlockList();
for (const [network, prefix] of [
['2001::', 23], // Protocol assignments (including Teredo / benchmarking).
['2001:db8::', 32],
['2002::', 16], // Documentation / 6to4.
['3fff::', 20], // Documentation.
] as const)
blockedV6.addSubnet(network, prefix, 'ipv6');
const mappedV4 = new BlockList();
mappedV4.addSubnet('::ffff:0:0', 96, 'ipv6');
function isPublicAddress(address: string): boolean {
if (isIP(address) === 4) return !blockedV4.check(address, 'ipv4');
// Node BlockList handles IPv4-mapped IPv6 in both dotted and hex forms.
if (mappedV4.check(address, 'ipv6'))
return !blockedV4.check(address, 'ipv6');
return globalV6.check(address, 'ipv6') && !blockedV6.check(address, 'ipv6');
}
@@ -0,0 +1,28 @@
-----BEGIN PRIVATE KEY-----
MIIEvQIBADANBgkqhkiG9w0BAQEFAASCBKcwggSjAgEAAoIBAQC+U4KN4bcXoh+W
TMyRIBLi2jMWmqv8ZB5rcyiVvcqp4Z8aJOOuiSZjMJJxxsVwu6qKLZwm8KThYo4M
m5x5tXMOvUKhqxDoXhQcql5dAiydoFXCJD56cQnRpS/BT5PtqtSBG9VVo3v1YF13
rKBF6ETcIiHFT7+jlSzPUP3lnyB7QL+yDFfaqX68TGXRkFG3thl3YMF+JhFILe7V
s8H2yZAiI3HuIeBUdTHOqnqVFW2vH1Hl4soXOsEQlQF5VyAFbCJ44duDfCe5uSbX
8lA+1VNJSRLKaQUdqvqgjF4UKmnk6a7xMAW0WTiC7RFoIUkH1WosprwEygYohhiF
XoG6QbKhAgMBAAECggEAA/H6NlOz9mbzbauo3+dAzPgF8BWDtCclJEgOUtBM16mo
ISQbnh4UsCCtIHOk2xngxp18a6g4Wr2uwR8mprU2rdsJew1vO8nbc96qNxZY82mD
7ZLPwrz+HZzleQXbxKTyY7y+dth9NNBrD5SB/AD9EG0asxrcl5j7hU6h/LUIONXN
GmOwEs7N7IsZgQtfItX8Zgo8+PUfQMp2QUawM4BAXwOSW3rpyioNmMQCa0Y/uKAf
8CKvNcP2XcbKucZNc9KHpNtaAjq7PDJ+7L/mSP3MGGNlYRjaasPiIemo2C6Pg1Pb
pr+yhwt28PFTqf2LLZcuivBtcPRlMmawcO9KFLjRewKBgQDm28ikYl8bqrROEgNQ
QL5rY/DzaaFjyEwvj+qWXWzy6cykQpw4vdLIfe3Kn9Evc4dU1o//RYnoqIWi7EJ0
Ma6OBp4lMocHajklZdCG+ArQENcPcdRgVWGf6LZtEbPGaTk88mDOEMm50/FfnKGW
7FWsmpmcKqZduq2NX2BvzBDmtwKBgQDTDbTGIGelbe1Cl4ex/zBbcFA9CUm/4583
1NMpfZNJaDbi9E3rOZK4ryCHAy4B9GQnmVW+A7ChzqMz2y9jKKM0RHlBXzvkroas
IQpRip11Vb/6ZXH3rM+7zgRYHi61+L7de2YbtCUfEOxQ+DcE5yNg1WzgRlYh0yfQ
SJ9f2XsZZwKBgEsVYn1sbShvbbMSkrdQR15gI+bXDSGJ7JVvhkmfWybqOZ+W9n5R
5rNEmclUD1ISjgpeuni44jCkVsp1cuudmPsiVd8dPuN/fdSW96peFA412+xvBjbK
rjS3GFYC8uhuIqqa3jdHKITi1NdW9wtCFF9N7PXovTEw3O9k/NV/lmOjAoGAWsDZ
DB0hFHS5gloQYozeKWOZTTWyPc5OR76/cmbqL7WdbGgrHUvreHjt3sCSRwrlClYY
FZYWnO1zJjhJHzV5QF91WJPv+DzH8jpe6oNVg//0hmKa6CqqRRKosY+A/ITS5gBK
/vyuvbYUOBkT54rQnrIHmEUGgpL+2sRvq9Kj6V8CgYEAxJ46AcgAW0SriYSn+oMX
KNMuAeFFS87fMEc+L3e/kH8ctWOGF8WAJew1v7TV+hHJvycMWqBbuvq7N1ArFBH4
Hxq76pRoo6DMz5PxIVajSxj9+LuNWn9onuX9Ni0w2p9clH5nD1+Arnb2G91kZby8
CqB+d4aWPHBq/0GQDp7xDR8=
-----END PRIVATE KEY-----
@@ -0,0 +1,20 @@
-----BEGIN CERTIFICATE-----
MIIDNDCCAhygAwIBAgIUL/DZdKkHvo0op+ZSNbwEJ8/q98UwDQYJKoZIhvcNAQEL
BQAwGzEZMBcGA1UEAwwQcHJvdmlkZXIuZXhhbXBsZTAeFw0yNjA5MDYwNTU4Mjda
Fw0zNjA5MDMwNTU4MjdaMBsxGTAXBgNVBAMMEHByb3ZpZGVyLmV4YW1wbGUwggEi
MA0GCSqGSIb3DQEBAQUAA4IBDwAwggEKAoIBAQC+U4KN4bcXoh+WTMyRIBLi2jMW
mqv8ZB5rcyiVvcqp4Z8aJOOuiSZjMJJxxsVwu6qKLZwm8KThYo4Mm5x5tXMOvUKh
qxDoXhQcql5dAiydoFXCJD56cQnRpS/BT5PtqtSBG9VVo3v1YF13rKBF6ETcIiHF
T7+jlSzPUP3lnyB7QL+yDFfaqX68TGXRkFG3thl3YMF+JhFILe7Vs8H2yZAiI3Hu
IeBUdTHOqnqVFW2vH1Hl4soXOsEQlQF5VyAFbCJ44duDfCe5uSbX8lA+1VNJSRLK
aQUdqvqgjF4UKmnk6a7xMAW0WTiC7RFoIUkH1WosprwEygYohhiFXoG6QbKhAgMB
AAGjcDBuMB0GA1UdDgQWBBSiXrRLEeZImTh1I0IjfqRQoHIUfzAfBgNVHSMEGDAW
gBSiXrRLEeZImTh1I0IjfqRQoHIUfzAPBgNVHRMBAf8EBTADAQH/MBsGA1UdEQQU
MBKCEHByb3ZpZGVyLmV4YW1wbGUwDQYJKoZIhvcNAQELBQADggEBAFtT545lqFbw
3pgHsnFI2kpXOso1WhjJa44cliFt0L2rg9JgZCih27ba1OZmr55Y6X2lG2xmQtqF
Gzm00pyC8EW/Yjwe2xcf0Tk2pGOTQ8MJlXhKOvwcxL6c6d/mc7Uvy+B/cwWbn9F9
IeOgDiaPGTRWE/UgiVqq1ZARI3ExRN5G0QWejdw0Q3VYWh2L5lM9kKeRomudaX0T
YKQFuuuZji9/RQTN9ZRfjVG0imJOMfqi0Nrl4eva4kHP1FWzH7yR5ybjsuG9WiXx
6o7Ktfva+sRr//utJjluL8RvrxJ0d7Dd+BIKQQxBP3yL3YokweL1xzS2WdQM66QW
N+E360kOO9k=
-----END CERTIFICATE-----
@@ -0,0 +1,323 @@
import { classifyHostRequestFailure } from '@iptvnator/shared/host-health';
import { ProviderRequestError } from './provider-request-error';
import { ProviderAxiosTransport } from './provider-axios-transport';
import express from 'express';
import { readFileSync } from 'node:fs';
import { createServer, Agent as HttpsAgent } from 'node:https';
import { Server } from 'node:http';
import { AddressInfo } from 'node:net';
import { TLSSocket } from 'node:tls';
import { gzipSync } from 'node:zlib';
import {
ValidatedHttpClient,
WebBackendHttpClient,
} from './validated-http-client';
import { withServer } from './web-backend-app.spec-helpers';
const lan = {
allowPrivateNetworkTargets: true,
resolveHostname: async () => ['127.0.0.1'],
};
async function withListeningServer<T>(
server: Server,
hostname: string,
run: (port: number) => Promise<T>
): Promise<T> {
await new Promise<void>((resolve, reject) => {
server.once('error', reject);
server.listen(0, hostname, resolve);
});
try {
return await run((server.address() as AddressInfo).port);
} finally {
server.closeAllConnections();
await new Promise<void>((resolve, reject) =>
server.close((error) => (error ? reject(error) : resolve()))
);
}
}
describe('real axios validated transport', () => {
it('uses exactly the resolved address at connect time, preserves Host, and bypasses ambient proxy settings', async () => {
const proxy = express().use((_req, res) =>
res.status(500).send('unexpected proxy')
);
await withServer(proxy, async (proxyUrl) => {
const provider = express().use((req, res) =>
res.json({ host: req.headers.host })
);
await withServer(provider, async (providerUrl) => {
const { port } = new URL(providerUrl);
const resolveHostname = jest
.fn()
.mockResolvedValueOnce(['127.0.0.1'])
.mockResolvedValue(['127.0.0.2']);
const previous = process.env['http_proxy'];
process.env['http_proxy'] = proxyUrl;
try {
const response = await new ValidatedHttpClient({
...lan,
resolveHostname,
}).get(`http://provider.example:${port}/`);
expect(response.data).toEqual({
host: `provider.example:${port}`,
});
expect(resolveHostname).toHaveBeenCalledTimes(1);
} finally {
if (previous === undefined)
delete process.env['http_proxy'];
else process.env['http_proxy'] = previous;
}
});
});
});
it('uses the new validated DNS answer on each same-host redirect instead of a pooled socket', async () => {
const provider = express().use((_req, res) => res.redirect('/next'));
await withServer(provider, async (providerUrl) => {
const { port } = new URL(providerUrl);
const resolveHostname = jest
.fn()
.mockResolvedValueOnce(['127.0.0.1'])
.mockResolvedValueOnce(['127.0.0.2']);
await expect(
new ValidatedHttpClient({ ...lan, resolveHostname }).get(
`http://provider.example:${port}/`,
{ timeout: 500 }
)
).rejects.toMatchObject({
initialResponded: true,
cause: {
code: expect.stringMatching(
/^(ECONNREFUSED|ECONNABORTED)$/
),
},
});
expect(resolveHostname).toHaveBeenCalledTimes(2);
});
});
it.each(['invalid-gzip', 'unfinished'])(
'discards the %s redirect body immediately',
async (kind) => {
const provider = express();
provider.get('/start', (_req, res) => {
res.status(302).set('Location', '/final');
if (kind === 'invalid-gzip')
res.set('Content-Encoding', 'gzip').end('not gzip');
else res.write('a body that never ends');
});
provider.get('/final', (_req, res) => res.json({ ok: true }));
await withServer(provider, async (url) => {
await expect(
new ValidatedHttpClient(lan).get(`${url}/start`, {
timeout: 500,
})
).resolves.toMatchObject({ data: { ok: true } });
});
}
);
it('preserves double-slash paths and escaped segments without changing connection authority', async () => {
const provider = express().use((req, res) =>
res.json({ path: req.url, host: req.headers.host })
);
await withServer(provider, async (url) => {
const target = new URL(url);
await expect(
new ValidatedHttpClient(lan).get(
`${url}//other.example/a%2Fb?token=a%2Bb`
)
).resolves.toMatchObject({
data: {
path: '//other.example/a%2Fb?token=a%2Bb',
host: target.host,
},
});
});
});
it('keeps query serialization and sends only Location query on a real redirect', async () => {
const requests: string[] = [];
const provider = express().use((req, res) => {
requests.push(req.url);
if (requests.length === 1) res.redirect('/final?ticket=issued');
else res.send('done');
});
await withServer(provider, async (url) => {
await new ValidatedHttpClient(lan).get(`${url}/start?fixed=yes`, {
params: { password: 'a b+c' },
});
expect(requests).toEqual([
'/start?fixed=yes&password=a+b%2Bc',
'/final?ticket=issued',
]);
});
});
it('retains final arraybuffer, decompression, BOM text and JSON behavior', async () => {
const provider = express();
provider.get('/binary', (_req, res) =>
res.end(Buffer.from([0, 255, 128]))
);
provider.get('/compressed', (_req, res) =>
res.set('Content-Encoding', 'gzip').end(gzipSync('{"ok":true}'))
);
provider.get('/text', (_req, res) => res.end('\uFEFFhello'));
await withServer(provider, async (url) => {
const client = new ValidatedHttpClient(lan);
expect(
(
await client.get(`${url}/binary`, {
responseType: 'arraybuffer',
})
).data
).toEqual(Buffer.from([0, 255, 128]));
expect((await client.get(`${url}/compressed`)).data).toEqual({
ok: true,
});
expect((await client.get(`${url}/text`)).data).toEqual('hello');
});
});
it('cancels an unfinished final body and closes its connection', async () => {
const controller = new AbortController();
let started!: () => void;
const bodyStarted = new Promise<void>((resolve) => (started = resolve));
const provider = express().use((_req, res) => {
res.write('unfinished');
started();
});
await withServer(provider, async (url) => {
const pending = new ValidatedHttpClient(lan).get(url, {
signal: controller.signal,
});
await bodyStarted;
controller.abort();
await expect(pending).rejects.toMatchObject({
cause: { code: 'ERR_CANCELED' },
});
});
});
it.each(['truncated', 'invalid-gzip', 'timeout'])(
'retains response evidence for a %s final body',
async (kind) => {
const provider = express().use((_req, res) => {
if (kind === 'invalid-gzip') {
res.set('Content-Encoding', 'gzip').end('not gzip');
return;
}
res.set('Content-Length', '100').write('short');
if (kind === 'truncated') setTimeout(() => res.destroy(), 10);
});
await withServer(provider, async (url) => {
const error = await new ValidatedHttpClient(lan)
.get(url, { timeout: 100 })
.catch((error: unknown) => error);
expect(error).toBeInstanceOf(ProviderRequestError);
const failure = error as ProviderRequestError;
expect(failure.cause).toMatchObject({
response: { status: 200, statusText: 'OK' },
});
expect(classifyHostRequestFailure(failure.cause)).toBe(
'responded'
);
});
}
);
it('times out before headers on the custom native transport', async () => {
const provider = express().use(() => {
/* Deliberately silent synthetic provider. */
});
await withServer(provider, async (url) => {
await expect(
new ValidatedHttpClient(lan).get(url, { timeout: 100 })
).rejects.toMatchObject({
initialResponded: false,
cause: { code: 'ECONNABORTED' },
});
});
});
it('does not cap a healthy trickling body at the inactivity timeout', async () => {
const provider = express().use((_req, res) => {
let count = 0;
res.write('start');
const interval = setInterval(() => {
res.write('x');
if (++count === 8) {
clearInterval(interval);
res.end();
}
}, 25);
res.on('close', () => clearInterval(interval));
});
await withServer(provider, async (url) => {
await expect(
new ValidatedHttpClient(lan).get(url, { timeout: 100 })
).resolves.toMatchObject({ data: 'startxxxxxxxx' });
});
});
it('preserves TLS SNI and hostname verification with a pinned address', async () => {
// This key/certificate is exclusively a synthetic test fixture.
const cert = readFileSync(`${__dirname}/testing/provider-test.pem`);
const key = readFileSync(`${__dirname}/testing/provider-test.key`);
let requests = 0;
const server = createServer({ key, cert }, (req, res) => {
requests++;
res.setHeader('Content-Type', 'application/json');
res.end(
JSON.stringify({
host: req.headers.host,
sni: (req.socket as TLSSocket & { servername: string })
.servername,
})
);
});
const trustedTransport: WebBackendHttpClient = {
get: (url, options) => {
const agent = options?.httpsAgent as HttpsAgent & {
options: { ca?: Buffer };
};
// Supply the local test CA, retaining the production lookup,
// SNI and rejectUnauthorized defaults on the actual agent.
agent.options.ca = cert;
return new ProviderAxiosTransport().get(url, options);
},
};
await withListeningServer(server, '127.0.0.1', async (port) => {
await expect(
new ValidatedHttpClient(lan).get(
`https://provider.example:${port}`
)
).rejects.toMatchObject({
cause: { code: 'DEPTH_ZERO_SELF_SIGNED_CERT' },
});
expect(requests).toBe(0);
const client = new ValidatedHttpClient(lan, trustedTransport);
await expect(
client.get(`https://provider.example:${port}`)
).resolves.toMatchObject({
data: {
host: `provider.example:${port}`,
sni: 'provider.example',
},
});
await expect(
client.get(`https://wrong.example:${port}`)
).rejects.toMatchObject({
cause: { code: 'ERR_TLS_CERT_ALTNAME_INVALID' },
});
expect(requests).toBe(1);
});
});
it('connects to a pinned IPv6 address and an IPv6 literal in trusted LAN mode', async () => {
const server = new Server((_req, res) => res.end('ipv6'));
await withListeningServer(server, '::1', async (port) => {
const client = new ValidatedHttpClient({
...lan,
resolveHostname: async () => ['::1'],
});
expect(
(await client.get(`http://provider.example:${port}/`)).data
).toBe('ipv6');
expect((await client.get(`http://[::1]:${port}/`)).data).toBe(
'ipv6'
);
});
});
});
@@ -0,0 +1,198 @@
import { Agent } from 'node:http';
import { LookupFunction } from 'node:net';
import { ValidatedHttpClient } from './validated-http-client';
import {
resolvePublicHost,
StubHttpClient,
} from './web-backend-app.spec-helpers';
const policy = {
allowPrivateNetworkTargets: false,
resolveHostname: resolvePublicHost,
};
describe('validated HTTP redirect chain', () => {
it.each([301, 302, 303, 307, 308])(
'follows relative Location for %s, without replaying original params',
async (status) => {
const transport = new StubHttpClient();
transport.queueRedirect('../next?ticket=provider', status);
transport.queueResponse('done');
const result = await new ValidatedHttpClient(policy, transport).get(
'https://provider.example/base/start',
{
params: { username: 'demo', password: 'secret' },
timeout: 1234,
}
);
expect(result.data).toBe('done');
expect(transport.requests[0].params).toEqual({
username: 'demo',
password: 'secret',
});
expect(transport.requests[1]).toMatchObject({
url: 'https://provider.example/next?ticket=provider',
timeout: 1234,
});
expect(transport.requests[1].params).toBeUndefined();
}
);
it('resolves fragment-only Location against the sent query and detects the cycle', async () => {
const transport = new StubHttpClient();
transport.queueRedirect('#fragment');
await expect(
new ValidatedHttpClient(policy, transport).get(
'https://provider.example/start',
{
params: { token: 'secret' },
}
)
).rejects.toMatchObject({
policyError: { status: 502, message: 'Redirect cycle detected' },
});
expect(transport.requests).toHaveLength(1);
});
it.each([
['https://provider.example/next', true],
['https://provider.example:8443/next', false],
['http://provider.example/next', false],
['https://other.example/next', false],
])(
'scopes sensitive headers on redirect to %s',
async (location, retained) => {
const transport = new StubHttpClient();
transport.queueRedirect(location as string);
transport.queueResponse('done');
const headers = {
Authorization: 'Bearer secret',
cOoKiE: 'mac=secret',
'Proxy-Authorization': 'secret',
sN: 'serial',
'User-Agent': 'player',
};
await new ValidatedHttpClient(policy, transport).get(
'https://provider.example/start',
{ headers }
);
expect(transport.requests[1].headers).toEqual(
retained ? headers : { 'User-Agent': 'player' }
);
expect(headers.Authorization).toBe('Bearer secret');
}
);
it.each([
'http://127.0.0.1/',
'https://user:secret@provider.example/',
'file:///tmp/test',
'http://[invalid',
])(
'refuses unsafe Location %s without another transport call',
async (location) => {
const transport = new StubHttpClient();
transport.queueRedirect(location);
await expect(
new ValidatedHttpClient(policy, transport).get(
'https://provider.example/start'
)
).rejects.toMatchObject({ initialResponded: true });
expect(transport.requests).toHaveLength(1);
}
);
it('accepts five redirects and rejects a sixth before its destination', async () => {
for (const count of [5, 6]) {
const transport = new StubHttpClient();
for (let i = 0; i < count; i++) transport.queueRedirect(`/hop${i}`);
transport.queueResponse('done');
const request = new ValidatedHttpClient(policy, transport).get(
'https://provider.example/start'
);
if (count === 5)
await expect(request).resolves.toMatchObject({ data: 'done' });
else
await expect(request).rejects.toMatchObject({
policyError: { status: 502, message: 'Too many redirects' },
});
expect(transport.requests).toHaveLength(6);
}
});
it('rejects a redirect without Location', async () => {
const transport = new StubHttpClient();
transport.queueRedirect('');
await expect(
new ValidatedHttpClient(policy, transport).get(
'https://provider.example'
)
).rejects.toMatchObject({ policyError: { status: 502 } });
});
it('revalidates DNS on same-host hops and blocks a changed answer before connect', async () => {
const transport = new StubHttpClient();
transport.queueRedirect('/next');
const resolveHostname = jest
.fn()
.mockResolvedValueOnce(['93.184.216.34'])
.mockResolvedValueOnce(['127.0.0.1']);
await expect(
new ValidatedHttpClient(
{ ...policy, resolveHostname },
transport
).get('https://provider.example/start')
).rejects.toMatchObject({
policyError: { status: 400 },
initialResponded: true,
});
expect(resolveHostname).toHaveBeenCalledTimes(2);
expect(transport.requests).toHaveLength(1);
});
it('pins all validated IPv4/IPv6 records without a second DNS lookup', async () => {
const transport = new StubHttpClient();
transport.queueResponse('done');
const get = jest.spyOn(transport, 'get');
const resolveHostname = jest
.fn()
.mockResolvedValue(['93.184.216.34', '2606:4700:4700::1111']);
await new ValidatedHttpClient(
{ ...policy, resolveHostname },
transport
).get('https://provider.example/start');
const options = get.mock.calls[0][1];
if (!options) throw new Error('Missing transport options');
expect(options).toMatchObject({
maxRedirects: 0,
proxy: false,
adapter: 'http',
});
const agent = options.httpsAgent as Agent & {
options: { lookup: LookupFunction; proxyEnv: object };
};
const lookup = agent.options.lookup as LookupFunction;
const callback = jest.fn();
lookup('provider.example', { all: true }, callback);
expect(callback).toHaveBeenLastCalledWith(null, [
{ address: '93.184.216.34', family: 4 },
{ address: '2606:4700:4700::1111', family: 6 },
]);
lookup('provider.example', { family: 6 }, callback);
expect(callback).toHaveBeenLastCalledWith(
null,
'2606:4700:4700::1111',
6
);
expect(resolveHostname).toHaveBeenCalledTimes(1);
expect(agent.options).toMatchObject({ proxyEnv: {} });
});
it('does not send a hop when canceled during its DNS validation', async () => {
const controller = new AbortController();
const transport = new StubHttpClient();
const resolveHostname = async () => {
controller.abort();
return ['93.184.216.34'];
};
await expect(
new ValidatedHttpClient(
{ ...policy, resolveHostname },
transport
).get('https://provider.example', { signal: controller.signal })
).rejects.toMatchObject({ initialResponded: false });
expect(transport.requests).toHaveLength(0);
});
});
@@ -0,0 +1,216 @@
import { ProviderAxiosTransport } from './provider-axios-transport';
import {
discardProviderBody,
discardProviderErrorBody,
readProviderBody,
} from './provider-response';
import axios, { AxiosRequestConfig } from 'axios';
import { Agent as HttpAgent, ClientRequest } from 'node:http';
import { Agent as HttpsAgent } from 'node:https';
import { isIP, LookupFunction } from 'node:net';
import { ProviderRequestError } from './provider-request-error';
import { ProviderUrlPolicy, validateProviderUrl } from './provider-url-policy';
export interface WebBackendHttpGetOptions {
readonly headers?: Record<string, string>;
readonly params?: Record<string, string>;
readonly responseType?: 'arraybuffer';
readonly timeout?: number;
readonly signal?: AbortSignal;
}
export type ProviderTransportOptions = Omit<
WebBackendHttpGetOptions,
'responseType'
> &
Pick<
AxiosRequestConfig,
| 'responseType'
| 'maxRedirects'
| 'validateStatus'
| 'httpAgent'
| 'httpsAgent'
| 'proxy'
| 'adapter'
>;
export interface WebBackendHttpResponse<T> {
readonly data: T;
readonly status: number;
readonly statusText?: string;
readonly request?: ClientRequest;
readonly headers: { readonly location?: string };
}
/** Injected transports must honor the same options and status contract as axios. */
export interface WebBackendHttpClient {
get<T>(
url: string,
options?: ProviderTransportOptions
): Promise<WebBackendHttpResponse<T>>;
}
const REDIRECT_STATUSES = new Set([301, 302, 303, 307, 308]);
const SENSITIVE_HEADERS = new Set([
'authorization',
'cookie',
'proxy-authorization',
'sn',
]);
export class ValidatedHttpClient implements WebBackendHttpClient {
constructor(
private readonly policy: ProviderUrlPolicy,
private readonly transport: WebBackendHttpClient = new ProviderAxiosTransport()
) {}
async get<T>(
rawUrl: string,
options: WebBackendHttpGetOptions = {}
): Promise<WebBackendHttpResponse<T>> {
let currentUrl = rawUrl;
let params = options.params;
let headers = options.headers;
let initialResponded = false;
const visited = new Set<string>();
try {
for (let redirects = 0; ; redirects++) {
options.signal?.throwIfAborted();
const target = await validateProviderUrl(
currentUrl,
this.policy
);
options.signal?.throwIfAborted();
if ('message' in target) {
throw new ProviderRequestError(
initialResponded,
target.lookupError,
target
);
}
// Match axios serialization once. Location is resolved against
// the URL actually sent; original params are never replayed.
const sentUrl = new URL(
axios.getUri({ url: target.url.href, params })
);
sentUrl.hash = '';
if (visited.has(sentUrl.href))
throw redirectError('Redirect cycle detected');
visited.add(sentUrl.href);
const lookup = pinnedLookup(target.addresses);
// Fresh agents prevent socket-pool reuse across validations.
// Empty proxyEnv also disables Node's native env proxy support.
const agentOptions = { lookup, proxyEnv: {} };
const httpAgent = new HttpAgent(agentOptions);
const httpsAgent = new HttpsAgent(agentOptions);
let response: WebBackendHttpResponse<T>;
try {
response = await this.transport.get<T>(target.url.href, {
...options,
headers,
params,
adapter: 'http',
proxy: false,
maxRedirects: 0,
responseType: 'stream',
httpAgent,
httpsAgent,
validateStatus: (status) =>
REDIRECT_STATUSES.has(status) ||
(status >= 200 && status < 300),
});
if (!REDIRECT_STATUSES.has(response.status)) {
try {
return {
...response,
data: await readProviderBody(
response.data,
options.responseType === 'arraybuffer',
options.timeout,
response.request?.socket ?? undefined
),
};
} catch (error) {
if (axios.isCancel(error)) throw error;
// A broken/stalled body still proves the endpoint
// answered. Preserve buffered axios error semantics.
throw Object.assign(
new Error('Provider response body failed'),
{
cause: error,
response: {
status: response.status,
statusText: response.statusText,
},
}
);
}
}
initialResponded = true;
discardProviderBody(response.data);
} catch (error) {
discardProviderErrorBody(error);
throw error;
} finally {
httpAgent.destroy();
httpsAgent.destroy();
}
if (redirects >= 5) throw redirectError('Too many redirects');
const location = response.headers.location;
if (!location)
throw redirectError(
'Redirect response did not include a location'
);
let nextUrl: URL;
try {
nextUrl = new URL(location, sentUrl);
} catch {
throw redirectError('Redirect location is not a valid URL');
}
if (nextUrl.origin !== sentUrl.origin) {
headers = Object.fromEntries(
Object.entries(headers ?? {}).filter(
([name]) =>
!SENSITIVE_HEADERS.has(name.toLowerCase()) &&
name.toLowerCase() !== 'host'
)
);
}
params = undefined;
currentUrl = nextUrl.href;
}
} catch (error) {
if (error instanceof ProviderRequestError) throw error;
throw new ProviderRequestError(initialResponded, error);
}
}
}
function redirectError(message: string): ProviderRequestError {
return new ProviderRequestError(true, undefined, { message, status: 502 });
}
function pinnedLookup(addresses: readonly string[]): LookupFunction {
const records = addresses.map((address) => ({
address,
family: isIP(address),
}));
return (_hostname, options, callback) => {
const eligible = options.family
? records.filter((record) => record.family === options.family)
: records;
if (!eligible.length) {
callback(
Object.assign(
new Error('No validated address for the requested family'),
{ code: 'ENOTFOUND' }
),
[]
);
} else if (options.all) {
callback(null, eligible);
} else {
callback(null, eligible[0].address, eligible[0].family);
}
};
}
@@ -418,7 +418,7 @@ describe('web backend host connectivity guard', () => {
// not consume the single trial the breaker allows.
resolvable = false;
const refusedByPolicy = await call();
expect(refusedByPolicy.status).toBe(400);
expect(refusedByPolicy.status).toBe(200);
resolvable = true;
const trial = await call();
@@ -469,7 +469,7 @@ describe('web backend host connectivity guard', () => {
);
const first = await call();
expect(first.status).toBe(400);
expect(first.status).toBe(200);
// The DNS error is internal: the client sees what it always saw.
await expect(first.json()).resolves.toEqual({
message: 'Provider URL host could not be resolved',
@@ -492,21 +492,11 @@ describe('web backend host connectivity guard', () => {
});
it('does not fast-fail a provider whose redirect destination is dead', async () => {
// The shape this route actually produces, verified against axios
// 1.19.0: follow-redirects walks the chain inside one `get()`, so
// `config` still holds the URL we asked for and only
// `request._currentUrl` names the hop that failed. Reading `config.url`
// alone would compare the original URL with itself, find no redirect,
// and charge the dead destination to the provider that answered.
const redirectedFailure = () =>
Object.assign(new Error('connect ECONNREFUSED'), {
code: 'ECONNREFUSED',
config: { url: 'http://xtream.example/player_api.php' },
request: { _currentUrl: 'http://cdn.dead.example/stream' },
});
const httpClient = new StubHttpClient();
httpClient.queueNetworkError(redirectedFailure());
httpClient.queueNetworkError(redirectedFailure());
for (let i = 0; i < 2; i++) {
httpClient.queueRedirect('http://cdn.dead.example/stream');
httpClient.queueNetworkError(hostLevelFailure());
}
httpClient.queueResponse({ user_info: { username: 'demo' } });
const { guard } = createTestGuard();
@@ -534,7 +524,7 @@ describe('web backend host connectivity guard', () => {
action: 'get_account_info',
payload: { user_info: { username: 'demo' } },
});
expect(httpClient.requests).toHaveLength(3);
expect(httpClient.requests).toHaveLength(5);
}
);
});
@@ -689,31 +679,32 @@ describe('web backend host connectivity guard', () => {
const started = new Promise<void>((resolve) => {
arrived = resolve;
});
const transport = jest
.spyOn(httpClient, 'get')
.mockImplementationOnce(async () => {
arrived();
await pending;
if (outcome !== 'success') {
throw Object.assign(
hostLevelFailure(
outcome === 'cancel'
? 'ERR_CANCELED'
: 'ETIMEDOUT'
),
{
request: {
_currentUrl:
outcome ===
'redirect-failure'
? 'http://cdn.example/slow'
: `http://portal.example/${route === 'xtream' ? 'player_api.php' : ''}`,
},
}
);
}
return { data: [] as never };
});
const transport = jest.spyOn(httpClient, 'get');
if (outcome === 'redirect-failure') {
transport.mockImplementationOnce(async () => ({
data: '' as never,
status: 302,
headers: {
location: 'http://cdn.example/slow',
},
}));
}
transport.mockImplementationOnce(async () => {
arrived();
await pending;
if (outcome !== 'success') {
throw hostLevelFailure(
outcome === 'cancel'
? 'ERR_CANCELED'
: 'ETIMEDOUT'
);
}
return {
data: [] as never,
status: 200,
headers: {},
};
});
const trial = call();
try {
await started;
@@ -729,7 +720,9 @@ describe('web backend host connectivity guard', () => {
).message
)
).toBe(true);
expect(transport).toHaveBeenCalledTimes(1);
expect(transport).toHaveBeenCalledTimes(
outcome === 'redirect-failure' ? 2 : 1
);
} finally {
settle();
await trial;
@@ -0,0 +1,238 @@
import {
HostConnectivityGuard,
OPEN_DURATION_MS,
} from '@iptvnator/shared/host-health';
import { isHostConnectivityFastFailMessage } from '@iptvnator/shared/interfaces';
import axios from 'axios';
import express from 'express';
import { createWebBackendApp, WebBackendHttpClient } from './web-backend-app';
import {
registerProviderTarget,
resolvePublicHost,
StubHttpClient,
withServer,
} from './web-backend-app.spec-helpers';
describe('provider proxy redirect boundary', () => {
it.each(['/xtream', '/stalker', '/parse', '/parse-xml'])(
'%s refuses a private redirect before contacting its destination',
async (route) => {
let destinationRequests = 0;
const destination = express().get('/private', (_req, res) => {
destinationRequests++;
res.send('synthetic private response');
});
await withServer(destination, async (destinationUrl) => {
const provider = express().use((_req, res) => {
res.redirect(`${destinationUrl}/private`);
});
await withServer(provider, async (providerUrl) => {
// Only the first public endpoint is mapped to a synthetic
// local provider. Axios and its redirect transport are real.
const httpClient: WebBackendHttpClient = {
get: (url, options) =>
axios.get(
new URL(url).hostname === 'provider.example'
? providerUrl
: url,
options
),
};
await withServer(
createWebBackendApp({
httpClient,
resolveHostname: resolvePublicHost,
allowPrivateNetworkTargets: false,
}),
async (backend) => {
const id = await registerProviderTarget(
backend,
'http://provider.example'
);
const response = await fetch(
`${backend}${route}?targetId=${id}`
);
const body = await response.json();
expect(destinationRequests).toBe(0);
expect(response.status).toBe(
route.startsWith('/parse') ? 400 : 200
);
expect(body).toEqual({
status: 400,
message:
'Provider URL points to a private or local network address',
});
}
);
});
});
}
);
});
describe('redirect routes and admission lifecycle', () => {
it.each(['/xtream', '/stalker', '/parse', '/parse-xml'])(
'%s follows an allowed local redirect and parses the final body',
async (route) => {
const provider = express().use((req, res) => {
if (req.path !== '/final') {
res.redirect('/final');
return;
}
if (route === '/parse')
res.send(
'#EXTM3U\n#EXTINF:-1,Test\nhttps://example.com/test.ts'
);
else if (route === '/parse-xml')
res.type('xml').send(
'<?xml version="1.0"?><tv><channel id="test"><display-name>Test</display-name></channel></tv>'
);
else res.json({ ok: true });
});
await withServer(provider, async (url) => {
await withServer(
createWebBackendApp({ allowPrivateNetworkTargets: true }),
async (backend) => {
const id = await registerProviderTarget(backend, url);
const response = await fetch(
`${backend}${route}?targetId=${id}&action=test`
);
const body = await response.json();
expect(response.status).toBe(200);
if (route === '/parse')
expect(body).toMatchObject({ count: 1 });
else if (route === '/parse-xml')
expect(body).toHaveProperty('channels');
else
expect(body).toEqual({
action: 'test',
payload: { ok: true },
});
}
);
});
}
);
it.each(['/xtream', '/stalker', '/parse', '/parse-xml'])(
'%s rechecks a registered hostname before its initial connection',
async (route) => {
const transport = new StubHttpClient();
const resolveHostname = jest
.fn()
.mockResolvedValueOnce(['93.184.216.34'])
.mockResolvedValue(['127.0.0.1']);
await withServer(
createWebBackendApp({ httpClient: transport, resolveHostname }),
async (backend) => {
const id = await registerProviderTarget(
backend,
'https://provider.example'
);
const response = await fetch(
`${backend}${route}?targetId=${id}`
);
expect(response.status).toBe(
route.startsWith('/parse') ? 400 : 200
);
expect(await response.json()).toMatchObject({
status: 400,
});
expect(transport.requests).toHaveLength(0);
}
);
}
);
it.each(['/xtream', '/stalker'])(
'%s keeps query-only redirect failures off the initial endpoint record',
async (route) => {
const transport = new StubHttpClient();
const networkFailure = () =>
Object.assign(
new Error('secret http://user:password@provider.example'),
{ code: 'ENOTFOUND' }
);
transport.queueNetworkError(networkFailure());
transport.queueRedirect('?next=1');
transport.queueNetworkError(networkFailure());
transport.queueNetworkError(networkFailure());
transport.queueResponse({ ok: true });
const guard = new HostConnectivityGuard();
await withServer(
createWebBackendApp({
httpClient: transport,
hostGuard: guard,
resolveHostname: resolvePublicHost,
}),
async (backend) => {
const id = await registerProviderTarget(
backend,
'https://provider.example'
);
const call = () =>
fetch(`${backend}${route}?targetId=${id}`);
await call();
const redirected = await call();
expect(await redirected.text()).not.toContain('secret');
await call();
expect(await (await call()).json()).toMatchObject({
payload: { ok: true },
});
expect(transport.requests).toHaveLength(5);
}
);
}
);
it.each(['/xtream', '/stalker'])(
'%s releases a half-open trial after redirect DNS refusal and counts initial DNS failures',
async (route) => {
let now = 1000;
const guard = new HostConnectivityGuard({ now: () => now });
const transport = new StubHttpClient();
let failDns = false;
const resolveHostname = async (hostname: string) => {
if (failDns || hostname === 'blocked.example')
throw Object.assign(new Error('secret transport details'), {
code: 'ENOTFOUND',
});
return ['93.184.216.34'];
};
await withServer(
createWebBackendApp({
httpClient: transport,
hostGuard: guard,
resolveHostname,
}),
async (backend) => {
const id = await registerProviderTarget(
backend,
'https://provider.example'
);
const call = () =>
fetch(`${backend}${route}?targetId=${id}`);
failDns = true;
await call();
await call();
const refused = (await (await call()).json()) as {
message: string;
};
expect(
isHostConnectivityFastFailMessage(refused.message)
).toBe(true);
expect(transport.requests).toHaveLength(0);
now += OPEN_DURATION_MS + 1;
failDns = false;
transport.queueRedirect('https://blocked.example/next');
const trial = await call();
expect(await trial.json()).toEqual({
status: 400,
message: 'Provider URL host could not be resolved',
});
transport.queueResponse({ ok: true });
expect(await (await call()).json()).toMatchObject({
payload: { ok: true },
});
}
);
}
);
});
@@ -9,11 +9,11 @@
import { AddressInfo } from 'node:net';
import { Server } from 'node:http';
import { STALKER_MAG_USER_AGENT } from '@iptvnator/shared/interfaces';
import { createWebBackendApp, WebBackendHttpClient } from './web-backend-app';
import {
createWebBackendApp,
WebBackendHttpClient,
WebBackendHttpGetOptions,
} from './web-backend-app';
ProviderTransportOptions,
WebBackendHttpResponse,
} from './validated-http-client';
/** The transport-identity headers every portal-facing Stalker request carries. */
export const STALKER_IDENTITY_HEADERS = {
@@ -41,12 +41,17 @@ export class StubHttpClient implements WebBackendHttpClient {
readonly error?: Error;
readonly status?: number;
readonly statusText?: string;
readonly headers?: { location?: string };
}> = [];
queueResponse(data: unknown): void {
this.queuedResponses.push({ data });
}
queueRedirect(location: string, status = 302): void {
this.queuedResponses.push({ data: '', status, headers: { location } });
}
queueFailure(status: number, statusText = 'Provider failure'): void {
this.queuedResponses.push({ data: null, status, statusText });
}
@@ -61,8 +66,8 @@ export class StubHttpClient implements WebBackendHttpClient {
async get<T>(
url: string,
options: WebBackendHttpGetOptions = {}
): Promise<{ data: T }> {
options: ProviderTransportOptions = {}
): Promise<WebBackendHttpResponse<T>> {
this.requests.push({
headers: options.headers,
params: options.params,
@@ -90,7 +95,13 @@ export class StubHttpClient implements WebBackendHttpClient {
throw error;
}
if (response.status) {
if (
response.status &&
!(
options.validateStatus ??
((status) => status >= 200 && status < 300)
)(response.status)
) {
const error = new Error(response.statusText) as Error & {
response: { status: number; statusText: string };
};
@@ -101,7 +112,11 @@ export class StubHttpClient implements WebBackendHttpClient {
throw error;
}
return { data: response.data as T };
return {
data: response.data as T,
status: response.status ?? 200,
headers: response.headers ?? {},
};
}
}
@@ -147,6 +162,7 @@ export async function withServer<T>(
const address = server.address() as AddressInfo;
return await callback(`http://127.0.0.1:${address.port}`);
} finally {
server.closeAllConnections();
await new Promise<void>((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
});
@@ -893,7 +893,7 @@ https://stream.example/live.m3u8`);
`${baseUrl}/xtream?targetId=${targetId}&action=get_account_info`
);
expect(response.status).toBe(400);
expect(response.status).toBe(200);
await expect(response.json()).resolves.toEqual({
message:
'Provider URL points to a private or local network address',
+29 -244
View File
@@ -1,10 +1,7 @@
import cors from 'cors';
import express, { Express, Request, Response } from 'express';
import { createHash } from 'node:crypto';
import { lookup } from 'node:dns/promises';
import { isIP } from 'node:net';
import zlib from 'node:zlib';
import axios from 'axios';
import epgParser from 'epg-parser';
import parser from 'iptv-playlist-parser';
import {
@@ -33,19 +30,21 @@ import {
ProviderError,
} from './provider-error';
export interface WebBackendHttpGetOptions {
readonly headers?: Record<string, string>;
readonly params?: Record<string, string>;
readonly responseType?: 'arraybuffer';
readonly timeout?: number;
}
export interface WebBackendHttpClient {
get<T>(
url: string,
options?: WebBackendHttpGetOptions
): Promise<{ data: T }>;
}
import {
ValidatedHttpClient,
WebBackendHttpClient,
} from './validated-http-client';
import { ProviderRequestError } from './provider-request-error';
import {
ProviderUrlPolicy,
providerUrlErrorBody,
resolveHostname,
validateProviderUrl,
} from './provider-url-policy';
export type {
WebBackendHttpClient,
WebBackendHttpGetOptions,
} from './validated-http-client';
interface PlaylistParseError {
readonly message: string;
@@ -68,42 +67,12 @@ export interface WebBackendAppOptions {
readonly runtimeBackendUrl?: string;
}
interface ProviderUrlPolicy {
readonly allowPrivateNetworkTargets: boolean;
readonly resolveHostname: (hostname: string) => Promise<readonly string[]>;
}
interface ProviderUrlError {
readonly message: string;
readonly status: number;
/**
* The DNS failure behind a "could not be resolved" refusal, when that is
* what this is. Internal only — {@link providerUrlErrorBody} strips it,
* because the client is told the same thing it always was.
*
* A name that does not resolve is evidence about reachability, exactly like
* the `ENOTFOUND` the transport would have raised a moment later. Without
* it, a host whose DNS died is refused by the policy on every request and
* the breaker never opens, so each call keeps paying for the lookup.
*/
readonly lookupError?: unknown;
}
/** The client-facing half of a {@link ProviderUrlError}. */
function providerUrlErrorBody(error: ProviderUrlError): {
message: string;
status: number;
} {
return { message: error.message, status: error.status };
}
type ProviderTargetRegistry = Map<string, URL>;
export function createWebBackendApp(
options: WebBackendAppOptions = {}
): Express {
const app = express();
const httpClient = (options.httpClient ?? axios) as WebBackendHttpClient;
const guid = options.guid ?? createGuid;
const now = options.now ?? (() => new Date());
const hostGuard =
@@ -125,6 +94,10 @@ export function createWebBackendApp(
isPrivateNetworkProxyAllowed(),
resolveHostname: options.resolveHostname ?? resolveHostname,
};
const httpClient = new ValidatedHttpClient(
providerUrlPolicy,
options.httpClient
);
const providerTargets: ProviderTargetRegistry = new Map();
const corsMiddleware = cors({
@@ -179,8 +152,8 @@ export function createWebBackendApp(
return;
}
const targetId = createProviderTargetId(result);
providerTargets.set(targetId, result);
const targetId = createProviderTargetId(result.url);
providerTargets.set(targetId, result.url);
res.json({ targetId });
}
);
@@ -303,38 +276,10 @@ export function createWebBackendApp(
}
guardToken = admission.token;
const providerUrlError =
await normalizeAndValidateXtreamProviderUrl(
url,
providerUrlPolicy
);
if (providerUrlError) {
// Admitted, then abandoned before any request went out. A
// policy refusal — private address, bad scheme — says nothing
// about reachability, so it only hands the half-open slot back.
// A name that would not resolve is different: that IS the host
// failing to answer, and counting it is what lets the breaker
// stop paying for the same dead lookup on every request.
if (providerUrlError.lookupError !== undefined) {
reportProviderRequestFailure(
hostGuard,
guardToken,
providerUrlError.lookupError
);
} else {
releaseProviderRequest(hostGuard, guardToken);
}
guardToken = null;
res.status(providerUrlError.status).json(
providerUrlErrorBody(providerUrlError)
);
return;
}
url.href = normalizeXtreamServerUrl(url.href);
requestUrl = appendPathSegment(url, 'player_api.php');
// Provider URLs are validated by /provider-targets before they enter the registry.
// codeql[js/request-forgery]
const response = await httpClient.get(requestUrl, {
params: getProxyParams(req, ['targetId']),
timeout: PROVIDER_REQUEST_TIMEOUT_MS.xtream,
@@ -433,8 +378,6 @@ export function createWebBackendApp(
guardToken = observeProviderRequest(hostGuard, requestUrl);
}
// Provider URLs are validated by /provider-targets before they enter the registry.
// codeql[js/request-forgery]
const response = await httpClient.get(requestUrl, {
headers,
// `create_link` gets the longer budget: the portal mints a
@@ -492,78 +435,6 @@ function getRegisteredProviderUrl(
return targetUrl;
}
async function validateProviderUrl(
rawUrl: string,
policy: ProviderUrlPolicy
): Promise<URL | ProviderUrlError> {
let url: URL;
try {
url = new URL(rawUrl);
} catch {
return { message: 'Provider URL is not a valid URL', status: 400 };
}
if (url.protocol !== 'http:' && url.protocol !== 'https:') {
return {
message: 'Only http and https provider URLs are supported',
status: 400,
};
}
if (url.username || url.password) {
return {
message: 'Provider URL credentials are not supported',
status: 400,
};
}
if (policy.allowPrivateNetworkTargets) {
return url;
}
const hostname = normalizeHostname(url.hostname);
if (isLocalHostname(hostname) || isPrivateOrReservedIp(hostname)) {
return {
message:
'Provider URL points to a private or local network address',
status: 400,
};
}
if (isIP(hostname) === 0) {
let addresses: readonly string[];
try {
addresses = await policy.resolveHostname(hostname);
} catch (lookupError) {
return {
message: 'Provider URL host could not be resolved',
status: 400,
lookupError,
};
}
if (
addresses.length === 0 ||
addresses.some((address) =>
isPrivateOrReservedIp(normalizeHostname(address))
)
) {
return {
message:
'Provider URL points to a private or local network address',
status: 400,
};
}
}
return url;
}
async function resolveHostname(hostname: string): Promise<readonly string[]> {
const records = await lookup(hostname, { all: true, verbatim: true });
return records.map((record) => record.address);
}
function createProviderTargetId(url: URL): string {
return createHash('sha256').update(url.href).digest('hex');
}
@@ -636,29 +507,6 @@ function appendPathSegment(url: URL, segment: string): string {
return nextUrl.href;
}
async function normalizeAndValidateXtreamProviderUrl(
url: URL,
policy: ProviderUrlPolicy
): Promise<ProviderUrlError | null> {
let normalizedUrl: URL;
try {
normalizedUrl = new URL(normalizeXtreamServerUrl(url.href));
} catch {
return { message: 'Provider URL is not a valid URL', status: 400 };
}
const validatedUrl = await validateProviderUrl(
appendPathSegment(normalizedUrl, 'player_api.php'),
policy
);
if ('message' in validatedUrl) {
return validatedUrl;
}
url.href = normalizedUrl.href;
return null;
}
async function handlePlaylistParse(options: {
readonly guid: () => string;
readonly httpClient: WebBackendHttpClient;
@@ -667,8 +515,6 @@ async function handlePlaylistParse(options: {
readonly userAgent?: string;
}): Promise<Record<string, unknown> | PlaylistParseError> {
try {
// Provider URLs are validated by /provider-targets before playlist parsing.
// codeql[js/request-forgery]
const response = await options.httpClient.get<string>(options.url, {
timeout: PROVIDER_REQUEST_TIMEOUT_MS.playlist,
...(options.userAgent
@@ -689,14 +535,19 @@ async function handlePlaylistParse(options: {
};
} catch (error) {
logProviderRequestFailure({ error, route: '/parse', url: options.url });
const providerError = error as ProviderError;
if (error instanceof ProviderRequestError && error.policyError) {
return providerUrlErrorBody(error.policyError);
}
const providerError = (
error instanceof ProviderRequestError ? error.cause : error
) as ProviderError;
if (providerError?.response?.statusText !== undefined) {
return {
status: providerError.response.status ?? 500,
message: providerError.response.statusText,
};
}
const code = collectProviderErrorCodes(error)[0];
const code = collectProviderErrorCodes(providerError)[0];
return {
status: providerError?.response?.status ?? 500,
message: code
@@ -712,8 +563,6 @@ async function fetchEpgDataFromUrl(
url: URL
): Promise<unknown> {
const href = url.href;
// Provider URLs are validated by /provider-targets before XMLTV parsing.
// codeql[js/request-forgery]
const response = await httpClient.get<ArrayBuffer | string>(href, {
timeout: PROVIDER_REQUEST_TIMEOUT_MS.epg,
...(url.pathname.endsWith('.gz')
@@ -791,67 +640,3 @@ function getLastUrlSegment(value: string): string {
function createGuid(): string {
return Math.random().toString(36).slice(2);
}
function normalizeHostname(hostname: string): string {
return hostname.trim().replace(/^\[/, '').replace(/\]$/, '').toLowerCase();
}
function isLocalHostname(hostname: string): boolean {
return hostname === 'localhost' || hostname.endsWith('.localhost');
}
function isPrivateOrReservedIp(address: string): boolean {
const version = isIP(address);
if (version === 4) {
return isPrivateOrReservedIpv4(address);
}
if (version === 6) {
return isPrivateOrReservedIpv6(address);
}
return false;
}
function isPrivateOrReservedIpv4(address: string): boolean {
const parts = address.split('.').map((part) => Number(part));
if (
parts.length !== 4 ||
parts.some((part) => !Number.isInteger(part) || part < 0 || part > 255)
) {
return true;
}
const [first, second, third] = parts;
return (
first === 0 ||
first === 10 ||
first === 127 ||
(first === 100 && second >= 64 && second <= 127) ||
(first === 169 && second === 254) ||
(first === 172 && second >= 16 && second <= 31) ||
(first === 192 && second === 168) ||
(first === 192 && second === 0) ||
(first === 192 && second === 0 && third === 2) ||
(first === 198 && (second === 18 || second === 19)) ||
(first === 198 && second === 51 && third === 100) ||
(first === 203 && second === 0 && third === 113) ||
first >= 224
);
}
function isPrivateOrReservedIpv6(address: string): boolean {
const normalized = address.toLowerCase();
if (
normalized === '::' ||
normalized === '::1' ||
normalized.startsWith('fc') ||
normalized.startsWith('fd') ||
normalized.startsWith('fe80:')
) {
return true;
}
const mappedIpv4 = normalized.match(/::ffff:(\d+\.\d+\.\d+\.\d+)$/)?.[1];
return mappedIpv4 ? isPrivateOrReservedIpv4(mappedIpv4) : false;
}
+1
View File
@@ -5,6 +5,7 @@
"sourceRoot": "apps/web-e2e/src",
"implicitDependencies": [
"web",
"web-backend",
"stalker-mock-server",
"xtream-mock-server"
],