perf(portals): fast-fail requests to portal hosts that stopped answering (#1421)

This commit is contained in:
4gray authored and GitHub committed 2026-08-13 07:31:23 +02:00
1 parent 2d7811eb5f
commit e3f72f7dce
50 files changed
+3073 -174

No files matched your search

@@ -682,7 +682,10 @@ const electronApi: ElectronBridgeApi = {
token?: string;
serialNumber?: string;
requestId?: string;
skipConnectionGuard?: boolean;
}) => ipcRenderer.invoke('STALKER_REQUEST', payload),
resetHostConnectivityGuard: (url: string) =>
ipcRenderer.invoke('CONNECTIVITY_GUARD_RESET', { url }),
xtreamRequest: (payload: {
url: string;
params: Record<string, string>;
@@ -0,0 +1,32 @@
/**
* Lets the renderer clear the connectivity guard's memory of a portal host.
*
* The guard refuses requests to a host that just stopped answering, which is
* wrong the moment the user says "try again" or hands over a portal address
* that may now point somewhere else. Every such moment sends this, so nothing
* has to wait out the guard's window.
*
* The host key is derived the same way the request path derives it — both
* `normalizeXtreamServerUrl` and `buildStalkerRequestUrl` rebuild their URL
* from `URL.origin`, so the authority a request ends up using is always the
* authority of the URL stored on the playlist. `URL.host` also leaves any
* `user:pass@` userinfo out, so no credential reaches the key or the log.
*/
import { ipcMain } from 'electron';
import { CONNECTIVITY_GUARD_RESET } from '@iptvnator/shared/interfaces';
import { resetGuardedHost } from '../util/host-connectivity-guard';
export function registerConnectivityGuardHandlers(): void {
ipcMain.handle(
CONNECTIVITY_GUARD_RESET,
async (_event, payload: { url?: string } | undefined) => {
const url = payload?.url;
if (!url) {
return { success: false };
}
return { success: resetGuardedHost(url) };
}
);
}
@@ -0,0 +1,234 @@
import {
STALKER_REQUEST,
buildHostConnectivityFastFailMessage,
} from '@iptvnator/shared/interfaces';
const registeredHandlers = new Map<string, (...args: unknown[]) => unknown>();
const axiosMock = Object.assign(jest.fn(), {
isAxiosError: jest.fn(),
});
jest.mock('electron', () => ({
ipcMain: {
handle: jest.fn(
(channel: string, handler: (...args: unknown[]) => unknown) => {
registeredHandlers.set(channel, handler);
}
),
},
}));
jest.mock('axios', () => ({
__esModule: true,
default: axiosMock,
}));
jest.mock('./portal-debug.events', () => ({
emitPortalDebugEvent: jest.fn(),
}));
jest.mock('../services/stalker-playback-context.service', () => ({
rememberStalkerPlaybackContext: jest.fn(),
}));
const PORTAL_URL = 'http://dead-portal.example.com:8080/portal.php';
const PORTAL_ENDPOINT = 'http://dead-portal.example.com:8080';
const MAC_ADDRESS = '00:1A:79:AA:BB:CC';
describe('StalkerEvents host connectivity guard', () => {
let consoleErrorSpy: jest.SpyInstance;
let consoleWarnSpy: jest.SpyInstance;
let requestHandler: (...args: unknown[]) => unknown;
/** A host-level failure: the portal never answered. */
const connectionRefused = () =>
Object.assign(new Error('connect ECONNREFUSED 10.0.0.1:8080'), {
code: 'ECONNREFUSED',
});
const request = (overrides: Record<string, unknown> = {}) =>
requestHandler(
{},
{
url: PORTAL_URL,
macAddress: MAC_ADDRESS,
params: {
type: 'itv',
action: 'get_all_channels',
JsHttpRequest: '1-xml',
},
...overrides,
}
) as Promise<unknown>;
beforeEach(async () => {
jest.resetModules();
registeredHandlers.clear();
axiosMock.mockReset();
axiosMock.isAxiosError.mockReset();
// Nothing here is an axios error unless a test says so; the handler
// only needs `code` to classify a connection failure.
axiosMock.isAxiosError.mockReturnValue(false);
consoleErrorSpy = jest.spyOn(console, 'error').mockImplementation();
consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation();
await import('./stalker.events');
const handler = registeredHandlers.get(STALKER_REQUEST);
expect(handler).toBeDefined();
requestHandler = handler as (...args: unknown[]) => unknown;
});
afterEach(() => {
consoleErrorSpy.mockRestore();
consoleWarnSpy.mockRestore();
});
it('stops contacting a portal host that refused twice in a row', async () => {
axiosMock.mockRejectedValue(connectionRefused());
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toBeDefined();
expect(axiosMock).toHaveBeenCalledTimes(2);
await expect(request()).rejects.toThrow(
buildHostConnectivityFastFailMessage(PORTAL_ENDPOINT)
);
// The whole point: no third 15-second wait.
expect(axiosMock).toHaveBeenCalledTimes(2);
});
it('rejects with a real Error so the renderer keeps its classification', async () => {
// Electron serializes a rejected plain object to '[object Object]',
// which would destroy the renderer's timeout-vs-connection reading.
axiosMock.mockRejectedValue(connectionRefused());
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toBeInstanceOf(Error);
});
it('does not log a line per skipped request', async () => {
axiosMock.mockRejectedValue(connectionRefused());
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toBeDefined();
consoleErrorSpy.mockClear();
await expect(request()).rejects.toBeDefined();
expect(consoleErrorSpy).not.toHaveBeenCalled();
});
it('keeps trusting a host that answers, whatever the status is', async () => {
axiosMock
.mockRejectedValue(connectionRefused())
.mockRejectedValueOnce(connectionRefused())
.mockResolvedValueOnce({
status: 404,
statusText: 'Not Found',
data: {},
headers: {},
});
await expect(request()).rejects.toThrow('ECONNREFUSED');
// A 404 still proves the host is reachable, so the streak restarts.
await expect(request()).rejects.toThrow('HTTP Error 404');
await expect(request()).rejects.toThrow('ECONNREFUSED');
// Two failures happened in total, but not consecutively.
await expect(request()).rejects.toThrow('ECONNREFUSED');
expect(axiosMock).toHaveBeenCalledTimes(4);
});
describe('endpoint-discovery probes', () => {
it('are never fast-failed, because discovery is how a portal gets reclassified', async () => {
axiosMock.mockRejectedValue(connectionRefused());
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toThrow(
buildHostConnectivityFastFailMessage(PORTAL_ENDPOINT)
);
await expect(
request({ skipConnectionGuard: true })
).rejects.toThrow('ECONNREFUSED');
expect(axiosMock).toHaveBeenCalledTimes(3);
});
it('never count towards the guard', async () => {
// Discovery walks several candidate paths on one host and expects
// most of them to fail; counting that would abandon a slow portal.
axiosMock.mockRejectedValue(connectionRefused());
await expect(
request({ skipConnectionGuard: true })
).rejects.toBeDefined();
await expect(
request({ skipConnectionGuard: true })
).rejects.toBeDefined();
await expect(
request({ skipConnectionGuard: true })
).rejects.toBeDefined();
await expect(request()).rejects.toThrow('ECONNREFUSED');
expect(axiosMock).toHaveBeenCalledTimes(4);
});
it('clear the record on a 5xx too, not just a body', async () => {
// 5xx rejects with a response attached. The probe's failure must not
// count, but the response still proves the origin answered —
// dropping it is what lets the breaker open mid-discovery.
const serverError = () =>
Object.assign(
new Error('Request failed with status code 502'),
{
code: 'ERR_BAD_RESPONSE',
response: { status: 502, statusText: 'Bad Gateway' },
}
);
axiosMock
.mockRejectedValue(connectionRefused())
.mockRejectedValueOnce(connectionRefused())
.mockRejectedValueOnce(serverError())
.mockRejectedValueOnce(connectionRefused());
await expect(request()).rejects.toBeDefined();
await expect(
request({ skipConnectionGuard: true })
).rejects.toBeDefined();
await expect(request()).rejects.toThrow('ECONNREFUSED');
// The probe's 5xx reset the streak, so the failure above is the
// first of a new one and this request still goes out.
await expect(request()).rejects.toThrow('ECONNREFUSED');
expect(axiosMock).toHaveBeenCalledTimes(4);
});
it('still clear the record when a candidate answers', async () => {
// Authentication against auth-gated candidates is NOT exempt, so
// without this the breaker could open in the middle of discovery.
axiosMock
.mockRejectedValue(connectionRefused())
.mockRejectedValueOnce(connectionRefused())
.mockResolvedValueOnce({
status: 200,
statusText: 'OK',
data: { js: [] },
headers: {},
})
.mockRejectedValueOnce(connectionRefused());
await expect(request()).rejects.toBeDefined();
await expect(
request({ skipConnectionGuard: true })
).resolves.toEqual({ js: [] });
await expect(request()).rejects.toThrow('ECONNREFUSED');
// The probe's success reset the streak, so the failure above is the
// first of a new one and this request still goes out. Without that
// reset it would be the second, and this would be fast-failed.
await expect(request()).rejects.toThrow('ECONNREFUSED');
expect(axiosMock).toHaveBeenCalledTimes(4);
});
});
});
@@ -19,6 +19,14 @@ import { emitPortalDebugEvent } from './portal-debug.events';
import { formatPortalRequestError } from './portal-request-error.util';
import { assertRemoteUrlAllowed } from './url-safety';
import { requestWithValidatedRedirects } from '../util/validated-axios';
import {
HostConnectivityGuardError,
HostRequestToken,
beginGuardedHostRequest,
observeGuardedHostRequest,
reportGuardedHostFailure,
reportGuardedHostSuccess,
} from '../util/host-connectivity-guard';
export default class StalkerEvents {
static bootstrapStalkerEvents(): Electron.IpcMain {
@@ -47,11 +55,21 @@ ipcMain.handle(
token?: string;
serialNumber?: string;
requestId?: string;
/**
* Set by endpoint discovery only. Its probes walk several candidate
* paths on one host and expect most of them to fail, so their
* failures must not count towards the connectivity guard — and they
* must not be fast-failed by it either, since discovery is how a
* portal gets reclassified in the first place.
*/
skipConnectionGuard?: boolean;
}
) => {
const startedAt = Date.now();
let debugRequest: Record<string, unknown> | undefined;
let requestUrlForLog = payload.url;
let guardToken: HostRequestToken | null = null;
const countsTowardsGuard = !payload.skipConnectionGuard;
try {
const { url, macAddress, params, token, serialNumber, requestId } =
payload;
@@ -70,6 +88,14 @@ ipcMain.handle(
const fullUrl = buildStalkerRequestUrl(url, requestParams);
requestUrlForLog = fullUrl;
// Navigating a dead portal would otherwise queue dozens of
// 15-second timeouts in a row. Exempt probes still report a
// success, so a candidate that answers clears whatever the other
// candidates' authentication attempts recorded.
guardToken = countsTowardsGuard
? beginGuardedHostRequest(fullUrl)
: observeGuardedHostRequest(fullUrl);
// Determine timeout based on action type
// create_link requests can take longer as server generates stream URL
const isCreateLink = requestParams.action === 'create_link';
@@ -98,6 +124,10 @@ ipcMain.handle(
{ allowPrivateNetworks: true }
);
// The host answered — whatever the status says, it is reachable.
reportGuardedHostSuccess(guardToken);
guardToken = null;
// Check if response is successful
if (response.status >= 400) {
// The numeric code must live in the MESSAGE: ipcRenderer
@@ -174,6 +204,21 @@ ipcMain.handle(
emitPortalDebugEvent(debugEvent);
}
if (error instanceof HostConnectivityGuardError) {
// The failures that tripped the guard were logged when they
// happened; a line per skipped request would be the very log
// spam the guard exists to stop.
throw error;
}
// Exempt probes never count failures, but an error carrying an
// HTTP response still proves the origin answered — dropping that is
// what would let the breaker open mid-discovery.
reportGuardedHostFailure(guardToken, error, {
countFailures: countsTowardsGuard,
requestUrl: requestUrlForLog,
});
console.error(
'[STALKER_REQUEST] Failed',
redactSensitiveData(
@@ -1,4 +1,8 @@
import { XTREAM_CANCEL_SESSION } from '@iptvnator/shared/interfaces';
import {
CONNECTIVITY_GUARD_RESET,
XTREAM_CANCEL_SESSION,
buildHostConnectivityFastFailMessage,
} from '@iptvnator/shared/interfaces';
const registeredHandlers = new Map<string, (...args: unknown[]) => unknown>();
const axiosMock = Object.assign(jest.fn(), {
@@ -276,3 +280,188 @@ describe('XtreamEvents session cancellation', () => {
});
});
});
describe('XtreamEvents host connectivity guard', () => {
const GUARD_DISABLED_ENV = 'IPTVNATOR_DISABLE_CONNECTIVITY_GUARD';
const SERVER_URL = 'http://dead-panel.example.com:8080';
const SERVER_ENDPOINT = 'http://dead-panel.example.com:8080';
let consoleErrorSpy: jest.SpyInstance;
let consoleWarnSpy: jest.SpyInstance;
let requestHandler: (...args: unknown[]) => unknown;
let resetHandler: (...args: unknown[]) => unknown;
const timedOut = () =>
Object.assign(new Error('timeout of 30000ms exceeded'), {
code: 'ETIMEDOUT',
});
const request = (url = SERVER_URL) =>
requestHandler(
{},
{
url,
params: {
action: 'get_live_categories',
password: 'secret',
username: 'user1',
},
suppressErrorLog: true,
}
) as Promise<unknown>;
beforeEach(async () => {
jest.resetModules();
delete process.env[PERF_CAPTURE_ENV];
delete process.env[GUARD_DISABLED_ENV];
registeredHandlers.clear();
axiosMock.mockReset();
axiosMock.isAxiosError.mockReset();
axiosMock.isAxiosError.mockReturnValue(true);
consoleErrorSpy = jest.spyOn(console, 'error').mockImplementation();
consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation();
await import('./xtream.events');
// Same fresh module registry as the handler above, so both talk to the
// same guard instance.
const guardEvents = await import('./connectivity-guard.events');
guardEvents.registerConnectivityGuardHandlers();
requestHandler = registeredHandlers.get('XTREAM_REQUEST') as (
...args: unknown[]
) => unknown;
resetHandler = registeredHandlers.get(CONNECTIVITY_GUARD_RESET) as (
...args: unknown[]
) => unknown;
expect(requestHandler).toBeDefined();
expect(resetHandler).toBeDefined();
});
afterEach(() => {
consoleErrorSpy.mockRestore();
consoleWarnSpy.mockRestore();
delete process.env[GUARD_DISABLED_ENV];
});
it('stops contacting a panel that timed out twice in a row', async () => {
axiosMock.mockRejectedValue(timedOut());
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toBeDefined();
expect(axiosMock).toHaveBeenCalledTimes(2);
await expect(request()).rejects.toThrow(
buildHostConnectivityFastFailMessage(SERVER_ENDPOINT)
);
// The whole point: no third 30-second wait.
expect(axiosMock).toHaveBeenCalledTimes(2);
});
it('keeps other panels reachable', async () => {
axiosMock.mockRejectedValue(timedOut());
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toBeDefined();
// The normal axios path rejects with a plain object, not an Error —
// only the guard's own rejection is an Error instance.
await expect(
request('http://other-panel.example.com')
).rejects.toMatchObject({ message: 'timeout of 30000ms exceeded' });
expect(axiosMock).toHaveBeenCalledTimes(3);
});
it('treats a 5xx as proof the panel is alive', async () => {
const serverError = Object.assign(new Error('Request failed'), {
code: 'ERR_BAD_RESPONSE',
response: { status: 502, statusText: 'Bad Gateway', data: {} },
});
axiosMock
.mockRejectedValue(timedOut())
.mockRejectedValueOnce(timedOut())
.mockRejectedValueOnce(serverError);
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toBeDefined();
// Two timeouts happened, but the 502 between them broke the streak.
await expect(request()).rejects.toBeDefined();
expect(axiosMock).toHaveBeenCalledTimes(4);
});
it('does not count a cancelled request against the panel', async () => {
const cancelled = Object.assign(new Error('canceled'), {
code: 'ERR_CANCELED',
});
axiosMock.mockRejectedValue(cancelled);
await expect(request()).rejects.toMatchObject({ status: 499 });
await expect(request()).rejects.toMatchObject({ status: 499 });
await expect(request()).rejects.toMatchObject({ status: 499 });
expect(axiosMock).toHaveBeenCalledTimes(3);
});
it('does not log a line per skipped request', async () => {
// suppressErrorLog is set above, so assert on the unsuppressed path.
axiosMock.mockRejectedValue(timedOut());
const loud = () =>
requestHandler(
{},
{
url: SERVER_URL,
params: { action: 'get_vod_streams' },
}
) as Promise<unknown>;
await expect(loud()).rejects.toBeDefined();
await expect(loud()).rejects.toBeDefined();
expect(consoleErrorSpy).toHaveBeenCalledTimes(2);
consoleErrorSpy.mockClear();
await expect(loud()).rejects.toBeDefined();
expect(consoleErrorSpy).not.toHaveBeenCalled();
});
it('contacts the panel again once the guard is reset', async () => {
axiosMock.mockRejectedValue(timedOut());
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toThrow(
buildHostConnectivityFastFailMessage(SERVER_ENDPOINT)
);
await expect(resetHandler({}, { url: SERVER_URL })).resolves.toEqual({
success: true,
});
await expect(request()).rejects.toMatchObject({
message: 'timeout of 30000ms exceeded',
});
expect(axiosMock).toHaveBeenCalledTimes(3);
});
it('reports a reset it could not apply instead of pretending it worked', async () => {
await expect(resetHandler({}, { url: 'not a url' })).resolves.toEqual({
success: false,
});
await expect(resetHandler({}, {})).resolves.toEqual({
success: false,
});
await expect(resetHandler({}, undefined)).resolves.toEqual({
success: false,
});
});
it('never blocks a request while the kill switch is set', async () => {
process.env[GUARD_DISABLED_ENV] = '1';
axiosMock.mockRejectedValue(timedOut());
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toBeDefined();
await expect(request()).rejects.toMatchObject({
message: 'timeout of 30000ms exceeded',
});
expect(axiosMock).toHaveBeenCalledTimes(3);
});
});
@@ -16,6 +16,13 @@ import { redactSensitiveData } from '@iptvnator/shared/logging';
import { emitPortalDebugEvent } from './portal-debug.events';
import { formatPortalRequestError } from './portal-request-error.util';
import { requestWithValidatedRedirects } from '../util/validated-axios';
import {
HostConnectivityGuardError,
HostRequestToken,
beginGuardedHostRequest,
reportGuardedHostFailure,
reportGuardedHostSuccess,
} from '../util/host-connectivity-guard';
import {
createXtreamMainPerformanceCaptureForRequest,
createXtreamMeasuredTransformResponse,
@@ -62,6 +69,7 @@ ipcMain.handle(
);
let activeRequestKey: string | null = null;
let requestUrlForLog = payload.url;
let guardToken: HostRequestToken | null = null;
try {
const { url, params, requestId, sessionId } = payload;
@@ -70,6 +78,10 @@ ipcMain.handle(
const apiUrl = buildXtreamApiUrl(url, params);
requestUrlForLog = apiUrl.toString();
// Browsing a dead portal would otherwise queue dozens of 30-second
// timeouts in a row. Throws once the host has stopped answering.
guardToken = beginGuardedHostRequest(requestUrlForLog);
const controller = new AbortController();
if (requestId || sessionId) {
activeRequestKey = requestId ?? crypto.randomUUID();
@@ -114,6 +126,11 @@ ipcMain.handle(
config,
{ allowPrivateNetworks: true }
);
// The host answered — whatever the status says, it is reachable.
reportGuardedHostSuccess(guardToken);
guardToken = null;
// Check if response is successful
if (response.status >= 400) {
throw {
@@ -191,6 +208,17 @@ ipcMain.handle(
emitPortalDebugEvent(debugEvent);
}
if (error instanceof HostConnectivityGuardError) {
// The failures that tripped the guard were logged when they
// happened; a line per skipped request would be the very log
// spam the guard exists to stop.
throw error;
}
reportGuardedHostFailure(guardToken, error, {
requestUrl: requestUrlForLog,
});
if (!payload.suppressErrorLog) {
console.error(
'[XTREAM_REQUEST] Failed',
@@ -0,0 +1,598 @@
import {
buildHostConnectivityFastFailMessage,
isHostConnectivityFastFailMessage,
} from '@iptvnator/shared/interfaces';
import {
HostConnectivityGuard,
HostConnectivityGuardError,
HostRequestToken,
beginGuardedHostRequest,
classifyHostRequestFailure,
portalEndpointKeyOf,
reportGuardedHostFailure,
resetHostConnectivityGuardForTests,
} from './host-connectivity-guard';
const HOST = 'http://portal.example.com:8080';
const GUARD_DISABLED_ENV = 'IPTVNATOR_DISABLE_CONNECTIVITY_GUARD';
const OPEN_DURATION_MS = 30_000;
const FAILURE_WINDOW_MS = 120_000;
const TRIAL_TIMEOUT_MS = 45_000;
function timeoutError(code = 'ETIMEDOUT') {
return Object.assign(new Error('timeout of 30000ms exceeded'), { code });
}
describe('classifyHostRequestFailure', () => {
it('treats connection-level error codes as host-level failures', () => {
for (const code of [
'ECONNABORTED',
'ECONNREFUSED',
'EAI_AGAIN',
'EHOSTUNREACH',
'ENETUNREACH',
'ENOTFOUND',
'ETIMEDOUT',
]) {
expect(classifyHostRequestFailure(timeoutError(code))).toBe(
'host-level'
);
}
});
it('treats an error carrying an HTTP response as proof the host answered', () => {
// 5xx reaches the handlers as a rejection: validateStatus only
// tolerates < 500. The host still answered.
const serverError = Object.assign(new Error('Request failed'), {
code: 'ERR_BAD_RESPONSE',
response: { status: 502, statusText: 'Bad Gateway' },
});
expect(classifyHostRequestFailure(serverError)).toBe('responded');
});
it('does not read reachability into cancellations, SSRF refusals or plain errors', () => {
expect(
classifyHostRequestFailure(
Object.assign(new Error('canceled'), { code: 'ERR_CANCELED' })
)
).toBe('inconclusive');
expect(
classifyHostRequestFailure(
new Error('URL host could not be resolved')
)
).toBe('inconclusive');
expect(classifyHostRequestFailure(undefined)).toBe('inconclusive');
expect(classifyHostRequestFailure('ETIMEDOUT')).toBe('inconclusive');
});
it('never counts a mid-transfer connection reset', () => {
// A reset happens on hosts that are very much alive.
expect(classifyHostRequestFailure(timeoutError('ECONNRESET'))).toBe(
'inconclusive'
);
});
});
describe('portalEndpointKeyOf', () => {
it('keys on host and port so two panels on one machine stay separate', () => {
expect(
portalEndpointKeyOf('http://example.com:8080/player_api.php')
).toBe('http://example.com:8080');
expect(
portalEndpointKeyOf('http://example.com:9090/player_api.php')
).toBe('http://example.com:9090');
});
it('separates HTTP from HTTPS on their default ports', () => {
// `URL.host` omits a default port, so both would collapse onto
// `example.com` — and a panel whose TLS listener is broken while plain
// HTTP works is a routine IPTV setup.
expect(portalEndpointKeyOf('http://example.com/player_api.php')).toBe(
'http://example.com'
);
expect(portalEndpointKeyOf('https://example.com/player_api.php')).toBe(
'https://example.com'
);
});
it('leaves URL credentials out of the key', () => {
expect(portalEndpointKeyOf('http://user:pass@example.com/c')).toBe(
'http://example.com'
);
});
it('returns null for an unparseable URL instead of throwing', () => {
expect(portalEndpointKeyOf('not a url')).toBeNull();
});
});
describe('HostConnectivityGuardError', () => {
it('carries the shared fast-fail message and no status field', () => {
const error = new HostConnectivityGuardError(HOST);
expect(error).toBeInstanceOf(Error);
expect(error.message).toBe(buildHostConnectivityFastFailMessage(HOST));
// getStalkerRequestErrorStatus reads `status` first; a number there
// would read as "the endpoint answered".
expect(
(error as unknown as { status?: unknown }).status
).toBeUndefined();
expect(isHostConnectivityFastFailMessage(error.message)).toBe(true);
});
});
describe('HostConnectivityGuard', () => {
let clock = 1_000_000;
let opened: string[];
let guard: HostConnectivityGuard;
const advance = (ms: number) => {
clock += ms;
};
/** Runs a request that fails at the host level after `durationMs`. */
const failRequest = (durationMs = 0): void => {
const check = guard.check(HOST);
if (!check.allowed) {
throw new Error('expected the request to be allowed');
}
advance(durationMs);
guard.reportFailure(check.token);
};
const expectAllowed = (): HostRequestToken => {
const check = guard.check(HOST);
expect(check.allowed).toBe(true);
if (!check.allowed) {
throw new Error('unreachable');
}
return check.token;
};
const expectBlocked = (): void => {
expect(guard.check(HOST).allowed).toBe(false);
};
beforeEach(() => {
clock = 1_000_000;
opened = [];
delete process.env[GUARD_DISABLED_ENV];
guard = new HostConnectivityGuard({
now: () => clock,
onOpen: (host) => opened.push(host),
});
});
it('allows requests to a host it has never seen', () => {
expect(guard.check(HOST).allowed).toBe(true);
expect(opened).toEqual([]);
});
it('still allows the request after a single failure', () => {
failRequest(30_000);
expect(guard.check(HOST).allowed).toBe(true);
expect(opened).toEqual([]);
});
it('opens after two consecutive failures and reports the wait', () => {
failRequest(30_000);
failRequest(30_000);
expect(guard.check(HOST)).toEqual({
allowed: false,
retryAfterMs: OPEN_DURATION_MS,
});
expect(opened).toEqual([HOST]);
});
it('logs the open transition once, not per blocked request', () => {
failRequest(30_000);
failRequest(30_000);
expectBlocked();
expectBlocked();
expect(opened).toEqual([HOST]);
});
it('keeps other hosts untouched', () => {
failRequest(30_000);
failRequest(30_000);
expect(guard.check('http://other.example.com').allowed).toBe(true);
});
it('keeps the same host on another scheme reachable', () => {
// Regression: keying by `URL.host` dropped the default port, so a dead
// HTTPS panel fast-failed the working HTTP one on the same machine.
failRequest(30_000);
failRequest(30_000);
expectBlocked();
expect(guard.check('https://portal.example.com:8080').allowed).toBe(
true
);
});
it('counts a parallel fan-out that fails together as one failure', () => {
// Catalog init loads live/vod/series at once. One wifi hiccup failing
// all three is one piece of evidence, not a trip.
const first = expectAllowed();
const second = expectAllowed();
const third = expectAllowed();
advance(30_000);
guard.reportFailure(first);
guard.reportFailure(second);
guard.reportFailure(third);
expect(guard.check(HOST).allowed).toBe(true);
expect(opened).toEqual([]);
});
it('opens when a second fan-out fails after the first one did', () => {
const first = expectAllowed();
const second = expectAllowed();
advance(30_000);
guard.reportFailure(first);
guard.reportFailure(second);
failRequest(30_000);
expectBlocked();
expect(opened).toEqual([HOST]);
});
it('does not re-open on a sibling that settles after the window elapsed', () => {
// A and B start together; A plus a later request open the breaker. B is
// explicitly not counted as a strike, so it must not start a fresh
// 30-second window either — that would push the half-open trial past
// the intended cooldown.
const sibling = expectAllowed();
const first = expectAllowed();
advance(30_000);
guard.reportFailure(first);
failRequest(30_000);
expectBlocked();
expect(opened).toEqual([HOST]);
advance(OPEN_DURATION_MS);
guard.reportFailure(sibling);
const trial = expectAllowed();
expect(trial.trial).toBe(true);
expect(opened).toEqual([HOST]);
});
it('does not accumulate failures further apart than the streak window', () => {
failRequest(0);
advance(FAILURE_WINDOW_MS + 1);
failRequest(0);
expect(guard.check(HOST).allowed).toBe(true);
expect(opened).toEqual([]);
});
it('accumulates failures exactly at the streak window edge', () => {
failRequest(0);
advance(FAILURE_WINDOW_MS);
failRequest(0);
expectBlocked();
});
it('clears the record as soon as the host answers', () => {
failRequest(30_000);
const token = expectAllowed();
guard.reportSuccess(token);
failRequest(30_000);
expect(guard.check(HOST).allowed).toBe(true);
expect(opened).toEqual([]);
});
it('treats a success followed by an inconclusive report as a no-op', () => {
// The Stalker handler reports success on the response, then throws its
// own `HTTP Error 404` error, which is classified inconclusive.
failRequest(30_000);
const token = expectAllowed();
guard.reportSuccess(token);
guard.reportInconclusive(token);
failRequest(30_000);
expect(guard.check(HOST).allowed).toBe(true);
});
it('ignores an inconclusive failure entirely', () => {
const first = expectAllowed();
guard.reportInconclusive(first);
const second = expectAllowed();
guard.reportInconclusive(second);
expect(guard.check(HOST).allowed).toBe(true);
expect(opened).toEqual([]);
});
describe('half-open', () => {
beforeEach(() => {
failRequest(30_000);
failRequest(30_000);
expectBlocked();
advance(OPEN_DURATION_MS);
});
it('lets exactly one request through once the window elapses', () => {
const trial = expectAllowed();
expect(trial.trial).toBe(true);
expect(guard.check(HOST).allowed).toBe(false);
expect(guard.check(HOST).allowed).toBe(false);
});
it('closes the breaker when the trial succeeds', () => {
const trial = expectAllowed();
guard.reportSuccess(trial);
const next = expectAllowed();
expect(next.trial).toBe(false);
expect(guard.check(HOST).allowed).toBe(true);
});
it('re-opens immediately when the trial fails, without a second strike', () => {
const trial = expectAllowed();
advance(30_000);
guard.reportFailure(trial);
expect(guard.check(HOST)).toEqual({
allowed: false,
retryAfterMs: OPEN_DURATION_MS,
});
expect(opened).toEqual([HOST, HOST]);
});
it('offers the trial slot again when the trial is cancelled', () => {
const trial = expectAllowed();
guard.reportInconclusive(trial);
const retry = expectAllowed();
expect(retry.trial).toBe(true);
});
it('lets only the request holding the slot release it', () => {
// A trial can outlive its own 45 s window: the validated-redirect
// transport gives each of up to five hops its own 30 s budget. Once
// a replacement has been admitted, the abandoned request's late
// report must not hand the slot to a third request.
const abandoned = expectAllowed();
advance(TRIAL_TIMEOUT_MS);
const replacement = expectAllowed();
expect(replacement.trial).toBe(true);
guard.reportInconclusive(abandoned);
expectBlocked();
});
it('recovers when a trial never reports back', () => {
expectAllowed();
expectBlocked();
advance(TRIAL_TIMEOUT_MS);
const replacement = expectAllowed();
expect(replacement.trial).toBe(true);
});
});
describe('reset', () => {
it('reopens the host for requests immediately', () => {
failRequest(30_000);
failRequest(30_000);
expectBlocked();
guard.reset(HOST);
expect(expectAllowed().trial).toBe(false);
});
it('discards failures from requests that started before it', () => {
// The retry that cleared the breaker must not be poisoned by the
// 30-second stragglers it was waiting behind.
const straggler = expectAllowed();
const secondStraggler = expectAllowed();
advance(15_000);
guard.reset(HOST);
advance(15_000);
guard.reportFailure(straggler);
guard.reportFailure(secondStraggler);
expect(guard.check(HOST).allowed).toBe(true);
expect(opened).toEqual([]);
});
it('still counts failures from requests started after it', () => {
failRequest(30_000);
guard.reset(HOST);
failRequest(30_000);
failRequest(30_000);
expectBlocked();
});
it('supersedes in-flight requests for a host it has no record of', () => {
const token = expectAllowed();
guard.clear();
guard.reset(HOST);
advance(30_000);
guard.reportFailure(token);
failRequest(30_000);
expect(guard.check(HOST).allowed).toBe(true);
});
});
describe('bookkeeping bounds', () => {
it('forgets idle records instead of growing without bound', () => {
for (let index = 0; index < 300; index += 1) {
advance(1);
guard.check(`http://host-${index}.example.com`);
}
// The cap held, and a fresh host is still allowed through.
expect(guard.check(HOST).allowed).toBe(true);
});
it('keeps an open breaker while other hosts churn past the idle TTL', () => {
failRequest(30_000);
failRequest(30_000);
expectBlocked();
advance(1);
guard.check('http://noise.example.com');
expectBlocked();
});
});
describe('kill switch', () => {
afterEach(() => {
delete process.env[GUARD_DISABLED_ENV];
});
it('never blocks a request while disabled', () => {
process.env[GUARD_DISABLED_ENV] = '1';
failRequest(30_000);
failRequest(30_000);
failRequest(30_000);
expect(guard.check(HOST).allowed).toBe(true);
expect(opened).toEqual([]);
});
it('accepts the "true" spelling', () => {
process.env[GUARD_DISABLED_ENV] = 'true';
failRequest(30_000);
failRequest(30_000);
expect(guard.check(HOST).allowed).toBe(true);
});
it('is read per call, so an already open breaker stops blocking', () => {
failRequest(30_000);
failRequest(30_000);
expectBlocked();
process.env[GUARD_DISABLED_ENV] = '1';
expect(guard.check(HOST).allowed).toBe(true);
});
});
});
describe('reportGuardedHostFailure', () => {
const ENDPOINT = 'http://panel.example.com:8080';
const URL_ON_ENDPOINT = `${ENDPOINT}/player_api.php`;
let consoleWarnSpy: jest.SpyInstance;
const hopFailure = () =>
Object.assign(new Error('timeout of 30000ms exceeded'), {
code: 'ETIMEDOUT',
// A redirect hop lives on another origin and carries its own URL.
config: { url: 'http://cdn.example.com/player_api.php' },
});
const ownFailure = () =>
Object.assign(new Error('connect ECONNREFUSED'), {
code: 'ECONNREFUSED',
config: { url: URL_ON_ENDPOINT },
});
const attempt = (error: unknown): void => {
reportGuardedHostFailure(
beginGuardedHostRequest(URL_ON_ENDPOINT),
error,
// The handlers pass the URL they asked for; a failure on any other
// URL means a redirect answered first.
{ requestUrl: URL_ON_ENDPOINT }
);
};
beforeEach(() => {
resetHostConnectivityGuardForTests();
consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation();
});
afterEach(() => {
consoleWarnSpy.mockRestore();
resetHostConnectivityGuardForTests();
});
it('does not charge a same-origin redirect hop to the guarded endpoint', () => {
// `/player_api.php` -> 302 -> `/slow/player_api.php` on the same origin:
// the origin answered, so its record must clear even though the failing
// URL shares its origin. Otherwise two such requests fast-fail every
// other call to a portal that answers.
const sameOriginHop = () =>
Object.assign(new Error('timeout of 30000ms exceeded'), {
code: 'ETIMEDOUT',
config: { url: `${ENDPOINT}/slow/player_api.php` },
});
attempt(sameOriginHop());
attempt(sameOriginHop());
attempt(sameOriginHop());
expect(() => beginGuardedHostRequest(URL_ON_ENDPOINT)).not.toThrow();
});
it('does not charge a redirect hop on another origin to the guarded endpoint', () => {
// The guarded endpoint answered — it produced the redirect — so
// charging the downstream timeout to it would eventually fast-fail a
// working redirector.
attempt(hopFailure());
attempt(hopFailure());
attempt(hopFailure());
expect(() => beginGuardedHostRequest(URL_ON_ENDPOINT)).not.toThrow();
});
it('treats a reached redirect hop as proof the guarded endpoint answered', () => {
// Reaching a hop on another origin means the guarded endpoint returned a
// redirect, so its streak resets like after any other response —
// otherwise a single later timeout becomes the second strike against an
// endpoint that answered in between.
attempt(ownFailure());
attempt(hopFailure());
attempt(ownFailure());
expect(() => beginGuardedHostRequest(URL_ON_ENDPOINT)).not.toThrow();
});
it('still counts a failure the guarded endpoint itself produced', () => {
attempt(ownFailure());
attempt(ownFailure());
expect(() => beginGuardedHostRequest(URL_ON_ENDPOINT)).toThrow(
HostConnectivityGuardError
);
});
it('counts a failure that names no endpoint at all', () => {
// Not every transport error carries a config; absence must not become
// an excuse to ignore the evidence.
attempt(
Object.assign(new Error('socket hang up'), { code: 'ENOTFOUND' })
);
attempt(
Object.assign(new Error('socket hang up'), { code: 'ENOTFOUND' })
);
expect(() => beginGuardedHostRequest(URL_ON_ENDPOINT)).toThrow(
HostConnectivityGuardError
);
});
});
@@ -0,0 +1,624 @@
/**
* Per-host circuit breaker for portal requests.
*
* A dead portal costs the full axios timeout on every call (30 s for Xtream,
* 15 s for Stalker), and browsing a dead portal's catalog issues dozens of
* those back to back — 30-second spinners and a flooded main-process log. Once
* a host has refused to answer twice in a row there is nothing left to learn
* from waiting again, so subsequent requests fail immediately for a short
* while.
*
* The rules are deliberately timid, because being wrong here means refusing to
* talk to a portal that works:
*
* - Only connection-level evidence counts (see
* {@link classifyHostRequestFailure}). Any HTTP response — 200, 404, even
* 502 — proves the host is alive and clears the record.
* - The window is short (30 s), so a mistake costs one page of navigation, and
* a host that came back is retried on its own.
* - Half-open means exactly ONE request goes out, not a whole screenful.
* - An explicit reset (user retry, endpoint discovery) always wins, and
* failures from requests that started before it are discarded — otherwise a
* 30-second straggler settles right after the reset and re-opens the breaker
* underneath the retry that cleared it.
*
* Deliberately not persisted: process lifetime is the right scope for
* "unreachable right now".
*/
import { buildHostConnectivityFastFailMessage } from '@iptvnator/shared/interfaces';
/** Consecutive connection failures that trip the breaker. */
const FAILURE_THRESHOLD = 2;
/** How long requests fast-fail once the breaker is open. */
const OPEN_DURATION_MS = 30_000;
/**
* Failures further apart than this are unrelated, not a streak. Two sequential
* 30 s timeouts must fit inside it, hence comfortably above 60 s.
*/
const FAILURE_WINDOW_MS = 120_000;
/**
* Safety net for a half-open trial that never reports back. Above the longest
* request timeout (30 s) plus margin, so it only fires if a caller leaked the
* token — without it a lost report would keep the breaker open forever.
*/
const TRIAL_TIMEOUT_MS = 45_000;
/** Hosts tracked at once; portals per user are few, this is just a bound. */
const MAX_TRACKED_HOSTS = 256;
/** Idle records are forgotten; they hold nothing worth remembering. */
const IDLE_TTL_MS = 600_000;
const GUARD_DISABLED_ENV = 'IPTVNATOR_DISABLE_CONNECTIVITY_GUARD';
/**
* Error codes that prove the host itself did not answer.
*
* `ECONNRESET` is deliberately absent: a reset mid-transfer happens on hosts
* that are very much alive, and it is the one code a working stream can emit.
*/
const HOST_LEVEL_FAILURE_CODES = new Set([
'ECONNABORTED',
'ECONNREFUSED',
'EAI_AGAIN',
'EHOSTUNREACH',
'ENETUNREACH',
'ENOTFOUND',
'ETIMEDOUT',
]);
/**
* Thrown instead of making a request while the breaker is open.
*
* A real `Error`, because Electron serializes a rejected plain object to
* `[object Object]` and the renderer's classification would be lost. It
* carries NO `status` property on purpose — `getStalkerRequestErrorStatus`
* reads that field first, and a numeric status there would read as "the
* endpoint answered".
*/
export class HostConnectivityGuardError extends Error {
/** Scheme, host and port — see {@link portalEndpointKeyOf}. */
readonly endpoint: string;
constructor(endpoint: string) {
super(buildHostConnectivityFastFailMessage(endpoint));
this.name = 'HostConnectivityGuardError';
this.endpoint = endpoint;
}
}
/**
* Handed out by {@link HostConnectivityGuard.check} and passed back when the
* request settles. It records which attempt the report belongs to, so a reset
* or a parallel sibling cannot be mistaken for fresh evidence.
*/
export interface HostRequestToken {
/** Scheme, host and port — see {@link portalEndpointKeyOf}. */
readonly endpoint: string;
readonly epoch: number;
readonly startedAt: number;
/** Whether this request is the single probe allowed while half-open. */
readonly trial: boolean;
/**
* Which half-open slot this request holds, when it holds one.
*
* A trial can outlive its own window — `requestWithValidatedRedirects`
* gives each of up to five redirect hops its own 30 s budget — and once a
* replacement has been admitted, the abandoned request must not be able to
* free the replacement's slot. `trial: true` alone cannot tell the two
* apart, so the owner is identified explicitly.
*/
readonly trialId: number;
}
export type HostConnectivityCheck =
| { readonly allowed: true; readonly token: HostRequestToken }
| { readonly allowed: false; readonly retryAfterMs: number };
export type HostRequestOutcome = 'responded' | 'host-level' | 'inconclusive';
interface HostState {
consecutiveFailures: number;
/** When the last COUNTED failure was recorded. */
lastFailureAt: number;
openUntil: number;
trialStartedAt: number | null;
/** Monotonic id of the half-open slot; see `HostRequestToken.trialId`. */
trialId: number;
epoch: number;
lastTouchedAt: number;
}
export interface HostConnectivityGuardOptions {
now?: () => number;
onOpen?: (endpoint: string) => void;
}
/**
* Whether the guard is switched off. Read at call time, not at module load, so
* a debugging session can toggle it (matches `IPTVNATOR_ALLOW_INSECURE_TLS`).
*/
export function isHostConnectivityGuardDisabled(): boolean {
const value = process.env[GUARD_DISABLED_ENV]?.trim().toLowerCase();
return value === '1' || value === 'true';
}
/**
* What a failed request proves about the host.
*
* `responded` covers an error that still carries an HTTP response (5xx reaches
* the handlers as a rejection because `validateStatus` only tolerates <500) —
* the host is alive, so it clears the record. `inconclusive` covers everything
* that says nothing about reachability: a cancelled request, a URL rejected by
* the SSRF policy, a bug in our own code.
*/
export function classifyHostRequestFailure(error: unknown): HostRequestOutcome {
if (!error || typeof error !== 'object') {
return 'inconclusive';
}
if ((error as { response?: unknown }).response) {
return 'responded';
}
const code = (error as { code?: unknown }).code;
if (typeof code === 'string' && HOST_LEVEL_FAILURE_CODES.has(code)) {
return 'host-level';
}
return 'inconclusive';
}
/**
* The guard key for a request URL: its ORIGIN, i.e. scheme, host and port.
*
* Not `URL.host`, which omits a default port and therefore gives
* `http://panel.example` and `https://panel.example` the same key — two
* genuinely different endpoints, and a panel whose TLS listener is broken while
* plain HTTP works is a routine IPTV setup. Sharing one record there would let
* the dead one fast-fail the working one without ever contacting it.
*
* `URL.origin` also leaves out any `user:pass@` userinfo, so no credential
* reaches the key or the log line.
*/
export function portalEndpointKeyOf(url: string): string | null {
try {
const origin = new URL(url).origin;
// Opaque origins serialize as "null" and are not a usable key.
return origin && origin !== 'null' ? origin : null;
} catch {
return null;
}
}
export class HostConnectivityGuard {
private readonly states = new Map<string, HostState>();
private readonly now: () => number;
private readonly onOpen?: (endpoint: string) => void;
constructor(options: HostConnectivityGuardOptions = {}) {
this.now = options.now ?? (() => Date.now());
this.onOpen = options.onOpen;
}
/**
* Whether a request to `endpoint` may go out, and the token to report with.
*
* Fast-fails while the breaker is open, and while a half-open trial is
* still in flight — the point of the trial is that a screenful of requests
* does not all hang again to learn the same thing.
*/
check(endpoint: string): HostConnectivityCheck {
const now = this.now();
if (isHostConnectivityGuardDisabled()) {
return {
allowed: true,
token: {
endpoint,
epoch: 0,
startedAt: now,
trial: false,
trialId: 0,
},
};
}
const state = this.ensureState(endpoint, now);
state.lastTouchedAt = now;
if (state.openUntil > now) {
return { allowed: false, retryAfterMs: state.openUntil - now };
}
// Never opened, or the open window elapsed. The latter is half-open:
// let one request through and hold the rest back until it settles.
let trial = false;
if (state.openUntil > 0) {
const trialStale =
state.trialStartedAt !== null &&
now - state.trialStartedAt >= TRIAL_TIMEOUT_MS;
if (state.trialStartedAt !== null && !trialStale) {
return { allowed: false, retryAfterMs: 0 };
}
state.trialStartedAt = now;
state.trialId += 1;
trial = true;
}
return {
allowed: true,
token: {
endpoint,
epoch: state.epoch,
startedAt: now,
trial,
trialId: trial ? state.trialId : 0,
},
};
}
/**
* The host answered. Always clears the record, whatever the status was and
* whichever attempt it belonged to — a reachable host is a reachable host.
*/
reportSuccess(token: HostRequestToken): void {
const state = this.states.get(token.endpoint);
if (!state || isHostConnectivityGuardDisabled()) {
return;
}
state.consecutiveFailures = 0;
state.lastFailureAt = 0;
state.openUntil = 0;
state.trialStartedAt = null;
state.lastTouchedAt = this.now();
}
/** The host did not answer at all. */
reportFailure(token: HostRequestToken): void {
const state = this.states.get(token.endpoint);
if (!state || isHostConnectivityGuardDisabled()) {
return;
}
const now = this.now();
// Superseded by an explicit reset: the user asked for a fresh attempt
// and this verdict predates it.
if (state.epoch !== token.epoch) {
this.releaseTrial(state, token);
return;
}
// Only the request that still holds the slot ends the half-open state.
// An abandoned trial's late failure is ordinary evidence — it goes
// through the streak rules below and leaves the replacement alone.
const wasTrial = this.ownsTrial(state, token);
if (wasTrial) {
state.trialStartedAt = null;
}
state.lastTouchedAt = now;
if (
state.consecutiveFailures > 0 &&
now - state.lastFailureAt > FAILURE_WINDOW_MS
) {
state.consecutiveFailures = 0;
}
// A request already in flight when the previous failure was recorded is
// not the next link in a streak — it is a sibling of it. Catalog
// loading fans out several requests at once, and one hiccup failing all
// of them is one piece of evidence, not a trip. A request that started
// at or after that moment is a genuine new attempt.
let counted = false;
if (
state.consecutiveFailures === 0 ||
token.startedAt >= state.lastFailureAt
) {
state.consecutiveFailures += 1;
state.lastFailureAt = now;
counted = true;
}
// Only a failure this report actually counted may trip the threshold.
// A sibling arriving after the window elapsed would otherwise open a
// fresh one off the existing count, pushing the half-open trial past
// the intended cooldown. A failed trial still goes straight back to
// open: the host had its chance.
if (
wasTrial ||
(counted && state.consecutiveFailures >= FAILURE_THRESHOLD)
) {
const wasOpen = state.openUntil > now;
state.openUntil = now + OPEN_DURATION_MS;
if (!wasOpen) {
this.onOpen?.(token.endpoint);
}
}
}
/**
* The request failed for a reason that says nothing about the host. Only
* releases the half-open slot, so the next request can be the trial.
*/
reportInconclusive(token: HostRequestToken): void {
const state = this.states.get(token.endpoint);
if (!state || isHostConnectivityGuardDisabled()) {
return;
}
this.releaseTrial(state, token);
state.lastTouchedAt = this.now();
}
/**
* Forgets everything recorded for `endpoint`, and invalidates the reports of
* requests already in flight. Called when the user asks for a real attempt
* or when the URL may now point somewhere else.
*/
reset(endpoint: string): void {
const now = this.now();
const state = this.ensureState(endpoint, now);
state.consecutiveFailures = 0;
state.lastFailureAt = 0;
state.openUntil = 0;
state.trialStartedAt = null;
state.epoch += 1;
state.lastTouchedAt = now;
}
/**
* A token for a request that must NOT be policed or counted, but whose
* success still clears the record.
*
* Stalker endpoint discovery probes several candidate paths on one host and
* expects most of them to fail; counting those would let it declare a
* slow-but-alive portal unreachable. It does not take the half-open slot
* either — a probe is not the trial the guard is waiting for.
*/
observe(endpoint: string): HostRequestToken {
return {
endpoint,
epoch: this.states.get(endpoint)?.epoch ?? 0,
startedAt: this.now(),
trial: false,
trialId: 0,
};
}
/** Test seam: drops all recorded state. */
clear(): void {
this.states.clear();
}
/** Whether `token` still holds the current half-open slot. */
private ownsTrial(state: HostState, token: HostRequestToken): boolean {
return (
token.trial &&
state.trialStartedAt !== null &&
state.epoch === token.epoch &&
state.trialId === token.trialId
);
}
private releaseTrial(state: HostState, token: HostRequestToken): void {
if (this.ownsTrial(state, token)) {
state.trialStartedAt = null;
}
}
private ensureState(endpoint: string, now: number): HostState {
const existing = this.states.get(endpoint);
if (existing) {
return existing;
}
this.prune(now);
const state: HostState = {
consecutiveFailures: 0,
lastFailureAt: 0,
openUntil: 0,
trialStartedAt: null,
trialId: 0,
epoch: 0,
lastTouchedAt: now,
};
this.states.set(endpoint, state);
return state;
}
private prune(now: number): void {
for (const [endpoint, state] of this.states) {
if (
now - state.lastTouchedAt > IDLE_TTL_MS &&
state.openUntil <= now &&
state.trialStartedAt === null
) {
this.states.delete(endpoint);
}
}
// Still at the cap: forget the oldest records. Dropping one means
// contacting that host again, which is the safe direction to err in.
while (this.states.size >= MAX_TRACKED_HOSTS) {
const oldest = this.states.keys().next();
if (oldest.done) {
return;
}
this.states.delete(oldest.value);
}
}
}
let sharedGuard: HostConnectivityGuard | null = null;
/** The guard both portal IPC handlers share. */
export function getHostConnectivityGuard(): HostConnectivityGuard {
if (!sharedGuard) {
sharedGuard = new HostConnectivityGuard({
onOpen: (endpoint) =>
console.warn(
`[HostConnectivityGuard] ${endpoint} is not answering; skipping requests to it for ${
OPEN_DURATION_MS / 1000
}s`
),
});
}
return sharedGuard;
}
/** Test seam: forgets the shared guard so each spec starts clean. */
export function resetHostConnectivityGuardForTests(): void {
sharedGuard = null;
}
/**
* Reserves a slot for a request to `url` on the shared guard.
*
* Throws {@link HostConnectivityGuardError} instead of letting the request hang
* again while the breaker is open. Returns `null` for a URL with no usable
* origin — there is nothing to track then, and refusing the request over that
* would be worse than letting the transport report the real problem.
*/
export function beginGuardedHostRequest(url: string): HostRequestToken | null {
const endpoint = portalEndpointKeyOf(url);
if (!endpoint) {
return null;
}
const check = getHostConnectivityGuard().check(endpoint);
if (!check.allowed) {
throw new HostConnectivityGuardError(endpoint);
}
return check.token;
}
/** Untracked counterpart of {@link beginGuardedHostRequest}, see `observe`. */
export function observeGuardedHostRequest(
url: string
): HostRequestToken | null {
const endpoint = portalEndpointKeyOf(url);
return endpoint ? getHostConnectivityGuard().observe(endpoint) : null;
}
export function reportGuardedHostSuccess(token: HostRequestToken | null): void {
if (token) {
getHostConnectivityGuard().reportSuccess(token);
}
}
/**
* The URL a failed request was actually talking to, when the error says.
*
* Redirects are followed hop by hop, each with its own config, so a failure on
* a later hop carries THAT hop's URL rather than the one we asked for.
*/
function failedRequestUrlOf(error: unknown): string | null {
const url = (error as { config?: { url?: unknown } } | null)?.config?.url;
return typeof url === 'string' ? url : null;
}
function normalizedUrlOrNull(url: string): string | null {
try {
return new URL(url).toString();
} catch {
return null;
}
}
/**
* Whether the failure happened on a hop the guarded endpoint redirected us to.
*
* Reaching any later hop proves the guarded endpoint answered: the first hop is
* always the URL we asked for, and only a redirect status advances the chain.
* That holds for a same-origin redirect too, so comparing the whole URL — not
* just its origin — is what catches `/player_api.php` → `/slow/player_api.php`.
*
* Requires positive evidence: anything unparseable or unknown returns false and
* the failure is counted as usual, because guessing "redirect" here would stop
* the guard from ever tripping.
*/
function failedAfterRedirect(
error: unknown,
token: HostRequestToken,
requestUrl: string | undefined
): boolean {
const failedUrl = failedRequestUrlOf(error);
if (!failedUrl) {
return false;
}
const failedEndpoint = portalEndpointKeyOf(failedUrl);
if (failedEndpoint && failedEndpoint !== token.endpoint) {
return true;
}
if (!requestUrl) {
return false;
}
const failedNormalized = normalizedUrlOrNull(failedUrl);
const requestedNormalized = normalizedUrlOrNull(requestUrl);
return (
failedNormalized !== null &&
requestedNormalized !== null &&
failedNormalized !== requestedNormalized
);
}
/**
* Records what a failed request proved about its endpoint.
*
* `countFailures: false` is for requests exempt from the guard (endpoint
* discovery): their failures are expected and must not count, but an error that
* still carries an HTTP response proves the origin answered, and dropping that
* is what would let the breaker open in the middle of discovery.
*/
export function reportGuardedHostFailure(
token: HostRequestToken | null,
error: unknown,
options: { countFailures?: boolean; requestUrl?: string } = {}
): void {
if (!token) {
return;
}
const guard = getHostConnectivityGuard();
const countFailures = options.countFailures ?? true;
switch (classifyHostRequestFailure(error)) {
case 'host-level': {
if (!countFailures) {
break;
}
if (failedAfterRedirect(error, token, options.requestUrl)) {
// The guarded endpoint answered with a redirect, so this clears
// its record like any other response rather than merely
// declining to count the downstream failure. The failing hop is
// not guarded (it has no token of its own), so such a chain
// keeps costing a full timeout — a documented gap.
guard.reportSuccess(token);
break;
}
guard.reportFailure(token);
break;
}
case 'responded':
guard.reportSuccess(token);
break;
default:
guard.reportInconclusive(token);
break;
}
}
/** Forgets the recorded failures for the endpoint `url` points at. */
export function resetGuardedHost(url: string): boolean {
const endpoint = portalEndpointKeyOf(url);
if (!endpoint) {
return false;
}
getHostConnectivityGuard().reset(endpoint);
return true;
}
+2
View File
@@ -32,6 +32,7 @@ import { databaseWorkerClient } from './app/services/database-worker-client';
import WindowEvents from './app/events/window.events';
import { bootstrapWindowCloseGuard } from './app/services/window-close-guard.service';
import { registerStreamProbeHandlers } from './app/events/stream-probe';
import { registerConnectivityGuardHandlers } from './app/events/connectivity-guard.events';
import XtreamEvents from './app/events/xtream.events';
import { environment } from './environments/environment';
import {
@@ -163,6 +164,7 @@ export default class Main {
StalkerEvents.bootstrapStalkerEvents();
XtreamEvents.bootstrapXtreamEvents();
registerStreamProbeHandlers();
registerConnectivityGuardHandlers();
DatabaseEvents.bootstrapDatabaseEvents();
EpgEvents.bootstrapEpgEvents();
RemoteControlEvents.bootstrapRemoteControlEvents();