diff --git a/.changes/portals-dead-host-fast-fail.md b/.changes/portals-dead-host-fast-fail.md new file mode 100644 index 000000000..b2afcf095 --- /dev/null +++ b/.changes/portals-dead-host-fast-fail.md @@ -0,0 +1,10 @@ +--- +type: perf +area: portals +--- + +Unreachable Xtream and Stalker portals no longer cost a 30-second wait per +request. After a host fails to answer twice, further requests to it fail +immediately for 30 seconds, so browsing a dead portal stops filling the screen +with long spinners. Retrying, testing the connection, or editing the portal +address contacts it again right away. diff --git a/CLAUDE.md b/CLAUDE.md index 5569f14f0..14c3cc722 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -698,6 +698,7 @@ This project uses modern Angular signal-based APIs and patterns. **ALWAYS** use - `epg.events.ts` - EPG IPC registration; freshness/fetch orchestration lives in `epg-fetch.service.ts`, manual channel-mapping resolution and CRUD in `epg-mapping.service.ts`, worker lifecycle in `epg-worker.service.ts`, DB lookups in `epg-query.service.ts` - `xtream.events.ts` - Xtream Codes API - `stalker.events.ts` - Stalker portal API + - `connectivity-guard.events.ts` - `CONNECTIVITY_GUARD_RESET`: forgets the connection failures recorded for a portal host. Both portal handlers above run every request through the per-host circuit breaker in `util/host-connectivity-guard.ts` — after 2 consecutive connection-level failures (no HTTP response; `ETIMEDOUT`/`ENOTFOUND`/`ECONNREFUSED`/… but never `ECONNRESET`) requests to that endpoint fail immediately for 30 s. The key is `URL.origin`, not `URL.host`, which would give `http://panel` and `https://panel` one shared record and let a dead TLS listener fast-fail the working HTTP one instead of hanging the full 30 s/15 s axios timeout again, with one half-open trial request afterwards. Any HTTP response (4xx and 5xx included) clears the record. The refusal is a real `Error` whose wording is a renderer contract (`buildHostConnectivityFastFailMessage` in `libs/shared/interfaces`): it must carry no `HTTP Error `, no timeout wording and none of the auth phrases, or Stalker endpoint discovery misclassifies it and lazy portal repair fires against a host just declared dead. Discovery probes are exempt via the `skipConnectionGuard` payload flag (bypass + no failure counting, but successes still clear the record). Every user-driven retry/refresh that issues portal requests must reset BEFORE its first request, or the affordance fast-fails and looks broken; automatic and first-load paths deliberately do not reset. Current senders: Xtream content-gate Retry, Stalker catalog append retry (`retryContentPage`), Stalker search-page retry, `StalkerItvCacheService.refresh()` (Live TV refresh), both account-info dialogs' Retry, both destructive Xtream refresh implementations (`PlaylistRefreshActionService.refreshXtream()` and `RecentPlaylistsComponent.refreshXtreamPlaylist()`, before they delete the cached catalog), `StalkerPortalDiscoveryService.discover()`, and `PortalStatusService` on `skipCache`. Kill switch: `IPTVNATOR_DISABLE_CONNECTIVITY_GUARD=1`. Contract: `docs/architecture/host-connectivity-guard.md` - `player.events.ts` - External player IPC registration; MPV/VLC lifecycle logic lives in `mpv-session.service.ts`, `vlc-session.service.ts`, and shared `external-player-*` helpers - `settings.events.ts` - App settings - `electron.events.ts` - App version, etc. diff --git a/apps/electron-backend/src/app/api/main.preload.ts b/apps/electron-backend/src/app/api/main.preload.ts index 212f51604..c4fd4b4d5 100644 --- a/apps/electron-backend/src/app/api/main.preload.ts +++ b/apps/electron-backend/src/app/api/main.preload.ts @@ -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; diff --git a/apps/electron-backend/src/app/events/connectivity-guard.events.ts b/apps/electron-backend/src/app/events/connectivity-guard.events.ts new file mode 100644 index 000000000..fd3f95616 --- /dev/null +++ b/apps/electron-backend/src/app/events/connectivity-guard.events.ts @@ -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) }; + } + ); +} diff --git a/apps/electron-backend/src/app/events/stalker.events.spec.ts b/apps/electron-backend/src/app/events/stalker.events.spec.ts new file mode 100644 index 000000000..10fe5a800 --- /dev/null +++ b/apps/electron-backend/src/app/events/stalker.events.spec.ts @@ -0,0 +1,234 @@ +import { + STALKER_REQUEST, + buildHostConnectivityFastFailMessage, +} from '@iptvnator/shared/interfaces'; + +const registeredHandlers = new Map 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 = {}) => + requestHandler( + {}, + { + url: PORTAL_URL, + macAddress: MAC_ADDRESS, + params: { + type: 'itv', + action: 'get_all_channels', + JsHttpRequest: '1-xml', + }, + ...overrides, + } + ) as Promise; + + 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); + }); + }); +}); diff --git a/apps/electron-backend/src/app/events/stalker.events.ts b/apps/electron-backend/src/app/events/stalker.events.ts index 535cdbb8b..86bcf2a2d 100644 --- a/apps/electron-backend/src/app/events/stalker.events.ts +++ b/apps/electron-backend/src/app/events/stalker.events.ts @@ -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 | 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( diff --git a/apps/electron-backend/src/app/events/xtream.events.spec.ts b/apps/electron-backend/src/app/events/xtream.events.spec.ts index 9fd9830f1..0246dff3f 100644 --- a/apps/electron-backend/src/app/events/xtream.events.spec.ts +++ b/apps/electron-backend/src/app/events/xtream.events.spec.ts @@ -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 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; + + 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; + + 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); + }); +}); diff --git a/apps/electron-backend/src/app/events/xtream.events.ts b/apps/electron-backend/src/app/events/xtream.events.ts index 59072d371..0a0771ecd 100644 --- a/apps/electron-backend/src/app/events/xtream.events.ts +++ b/apps/electron-backend/src/app/events/xtream.events.ts @@ -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', diff --git a/apps/electron-backend/src/app/util/host-connectivity-guard.spec.ts b/apps/electron-backend/src/app/util/host-connectivity-guard.spec.ts new file mode 100644 index 000000000..3f61ff864 --- /dev/null +++ b/apps/electron-backend/src/app/util/host-connectivity-guard.spec.ts @@ -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 + ); + }); +}); diff --git a/apps/electron-backend/src/app/util/host-connectivity-guard.ts b/apps/electron-backend/src/app/util/host-connectivity-guard.ts new file mode 100644 index 000000000..c0c08e186 --- /dev/null +++ b/apps/electron-backend/src/app/util/host-connectivity-guard.ts @@ -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(); + 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; +} diff --git a/apps/electron-backend/src/main.ts b/apps/electron-backend/src/main.ts index 357a3f1f6..69b3d70b0 100644 --- a/apps/electron-backend/src/main.ts +++ b/apps/electron-backend/src/main.ts @@ -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(); diff --git a/apps/web/src/app/services/electron.service.ts b/apps/web/src/app/services/electron.service.ts index 4bef0c784..f495c2cb6 100644 --- a/apps/web/src/app/services/electron.service.ts +++ b/apps/web/src/app/services/electron.service.ts @@ -8,6 +8,7 @@ import { DataService, SettingsStore } from '@iptvnator/services'; import { AUTO_UPDATE_PLAYLISTS, AutoUpdatePlaylistsResult, + CONNECTIVITY_GUARD_RESET, ELECTRON_BRIDGE_SECURITY_ERROR_CODES, ERROR, normalizeHost, @@ -166,6 +167,14 @@ export class ElectronService extends DataService { )) as T; } + if (type === CONNECTIVITY_GUARD_RESET) { + const { url } = payload as { url?: string }; + if (url) { + await window.electron.resetHostConnectivityGuard(url); + } + return undefined as T; + } + if (type === 'OPEN_MPV_PLAYER') { const data = payload as PlayerLaunchPayload; try { @@ -279,6 +288,8 @@ export class ElectronService extends DataService { serialNumber?: string; /** Endpoint-discovery probes expect failures; no error snackbar. */ silent?: boolean; + /** Endpoint-discovery probes are exempt from the connectivity guard. */ + skipConnectionGuard?: boolean; }) { const context = createPortalDebugRequestContext({ provider: 'stalker', diff --git a/docs/architecture/host-connectivity-guard.md b/docs/architecture/host-connectivity-guard.md new file mode 100644 index 000000000..404614bfc --- /dev/null +++ b/docs/architecture/host-connectivity-guard.md @@ -0,0 +1,245 @@ +# Host Connectivity Guard + +Per-host circuit breaker for portal requests in the Electron main process. + +## The problem + +Every request to an unreachable portal costs its full axios timeout — 30 s for +`XTREAM_REQUEST`, 15 s for `STALKER_REQUEST` (30 s for `create_link`). Browsing a +dead portal's catalog issues dozens of those back to back, which shows up as +30-second spinners and a main-process log full of identical failures. Once a host +has refused to answer twice in a row there is nothing left to learn from waiting +again. + +## Where it lives + +`apps/electron-backend/src/app/util/host-connectivity-guard.ts` — a pure module +(no Electron imports) wired into both IPC handlers. + +The handlers are the choke point that sees _all_ traffic to a host. The +renderer's `executeStalkerRequest` is not: it has four documented bypasses +(auth, endpoint discovery, account info, the row-less stream resolver), and that +is exactly the traffic that hits dead hosts. + +**The PWA is deliberately not covered yet.** `apps/web-backend` sets no +per-request timeout at all, so a dead host there hangs on OS-level TCP timeouts +rather than a 15/30 s budget — a timeout-driven breaker would rarely trip. The +guard is written dependency-free so it can move into a `domain:shared-runtime` +library and be shared with `web-backend` when that gap is closed. + +## Rules + +Being wrong here means refusing to talk to a portal that works, so every rule +errs towards contacting the host: + +| | | +| ----------- | ---------------------------------------------------------------------- | +| Trip | 2 consecutive host-level failures within an inclusive 120 s window | +| Open for | 30 s (`OPEN_DURATION_MS`), matching the repo's other cooldowns | +| Half-open | exactly ONE trial request; the rest keep fast-failing until it settles | +| Reset | any HTTP response — 200, 404, even 502 — the host answered | +| Key | `URL.origin` — scheme, host **and** port (see below) | +| Kill switch | `IPTVNATOR_DISABLE_CONNECTIVITY_GUARD=1` (read per call) | + +**Host-level failure** means an error with no HTTP response whose code is one of +`ETIMEDOUT`, `ECONNABORTED`, `ENOTFOUND`, `EAI_AGAIN`, `ECONNREFUSED`, +`EHOSTUNREACH`, `ENETUNREACH`. `ECONNRESET` is deliberately excluded: a reset +mid-transfer happens on hosts that are very much alive. Cancelled requests +(`ERR_CANCELED`) and SSRF-policy refusals are `inconclusive` — they say nothing +about reachability and only release the half-open slot. + +**A failure is only charged to the endpoint that produced it.** Redirects are +followed hop by hop, each with its own config, so a failure on a later hop +carries that hop's URL in `error.config.url`. 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 — so it CLEARS the guarded endpoint's +record, exactly like any other response. Merely declining to count it would leave +an earlier direct failure standing, and a single later timeout would then +fast-fail an endpoint that answered in between. + +The comparison is against the whole request URL, not just its origin: a +same-origin redirect (`/player_api.php` → `/slow/player_api.php`) proves the +endpoint answered just as much as a cross-origin one, and charging it would +fast-fail every OTHER call to a portal that answers. Both handlers therefore pass +the URL they asked for. It requires positive evidence — anything unparseable or +unknown counts the failure as usual, because guessing "redirect" here would stop +the guard from ever tripping — and a failure that names no URL at all is still +counted. Round-tripping through `URL` is identity for both handlers' URL shapes +(including Stalker's hand-encoded `cmd`), and `requestWithValidatedRedirects` +normalizes hop 1 the same way, so the comparison is exact. + +Known gap: the failing hop is not guarded either (it has no token of its own), +so a permanently broken redirect chain keeps costing a full timeout. + +**The key is the origin, not the host.** `URL.host` omits a default port, so +`http://panel.example` and `https://panel.example` would share one record — +two genuinely different endpoints, and a panel whose TLS listener is broken +while plain HTTP works is a routine IPTV setup. Sharing state 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. + +Two more rules exist because of specific failure modes: + +- **Siblings are not a streak.** Catalog initialization fans out three category + requests at once; one network hiccup failing all three is one piece of + evidence, not a trip. A failure counts only if its request started at or after + the moment the previous failure was recorded. Timestamps are millisecond + coarse, so this only separates siblings once a request actually took time — + which is precisely the expensive case worth protecting. Only a failure that + was actually counted may trip the threshold: a sibling settling after the open + window elapsed would otherwise start a fresh one off the existing count and + push the half-open trial past the intended cooldown. +- **A reset invalidates reports already in flight.** `reset()` bumps a per-host + epoch instead of deleting the record, and a failure reported under an older + epoch is discarded. Without that, the 30-second stragglers a user was waiting + behind settle right after they press Retry and re-open the breaker underneath + the very retry that cleared it. + +A half-open trial that never reports back expires after 45 s, so a leaked token +cannot leave the breaker open forever. That expiry is why the slot has an +identity: a trial can genuinely outlive its 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's late report must not free +the replacement's slot and let a third request through. `trial: true` alone +cannot tell the two apart, so the token carries the slot id it owns; an +abandoned trial's failure is still counted as ordinary evidence. + +## The fast-fail error is a renderer contract + +`HostConnectivityGuardError` is a real `Error`: Electron serializes a rejected +plain object to `[object Object]`, which would destroy the renderer's +classification. It carries no `status` property, because +`getStalkerRequestErrorStatus` reads that field first. + +The message comes from `buildHostConnectivityFastFailMessage()` in +`libs/shared/interfaces/src/lib/host-connectivity.util.ts`. It names the full +endpoint, scheme included, so a user who imported the same panel over both HTTP +and HTTPS can tell which one was skipped. Its wording is load-bearing — the +Stalker renderer classifies transport failures purely from message text: + +| Must not contain | Otherwise | +| ---------------------------------------------------------------------------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------- | +| `HTTP Error ` | reads as "endpoint absent, probe the next candidate"; 404/401/403 also fire lazy portal repair against a host we just declared dead | +| `timed out`, `timeout of Nms`, `ETIMEDOUT` | discovery walks every candidate instead of aborting early | +| `authorization`, `unauthorized`, `access denied`, `invalid token`, `auth failed`, `handshake failed` | fires lazy portal repair (a bare `authorization` matches) | + +What remains is the "connection-level failure" slot the renderer already has for +ECONNREFUSED/ENOTFOUND: discovery stops probing and reports the host +unreachable, `shouldAttemptRepair` returns false, and the message reaches error +snackbars verbatim — which is why it reads like a sentence and names the +endpoint. + +**The endpoint in that sentence is user data, so the marker outranks the +heuristics.** A portal at `https://authorization.example` would otherwise make +its own fast-fail message match the broad auth phrase set and send an +unreachable host into lazy portal repair. `isStalkerAuthFailureMessage` and +`isStalkerProbeTimeout` therefore both return false for a message +`isHostConnectivityFastFailMessage` recognises, before their phrase matching +runs. (`getStalkerRequestErrorStatus` needs no such guard: `HTTP Error ` +contains a space, which a hostname cannot.) +`stalker-portal-discovery.utils.spec.ts` pins all three properties for the bare +and the IPC-wrapped form. + +## Exemption: endpoint discovery + +`STALKER_REQUEST` accepts `skipConnectionGuard`, set **only** by +`StalkerPortalDiscoveryService.probeContent`. Semantics: bypass the check, never +count failures, **but still report successes.** + +Discovery walks several candidate paths on one host and expects most of them to +fail; counting that would let it declare a slow-but-alive portal unreachable. +Reporting successes is equally load-bearing: `confirmFullPortal` runs the full +non-exempt authentication flow against auth-gated candidates, so without it two +hung handshakes could open the breaker mid-discovery and the next authenticate +would fast-fail into the "host unreachable" slot — abandoning a portal that +works. "Success" here means any observed response, including one attached to a +rejection: `validateStatus` lets 4xx through but rejects 5xx with +`error.response` set, and a 5xx proves the origin answered just as well as a +body does. That is why the exempt path reports through +`reportGuardedHostFailure(token, error, { countFailures: false })` rather than +skipping the report. + +## Explicit reset + +`CONNECTIVITY_GUARD_RESET` (`{ url }`) is handled by +`apps/electron-backend/src/app/events/connectivity-guard.events.ts`. One key +derivation is enough: both `normalizeXtreamServerUrl` and +`buildStalkerRequestUrl` rebuild their request URL from `URL.origin`, so the +origin a request ends up using is always the origin of the URL stored on the +playlist. + +**The rule: every user-driven retry or refresh that issues portal requests must +reset the guard before its first request.** The failures that opened the breaker +are usually the very ones the user is retrying, so a reset placed after the +request — or missing — makes the affordance do nothing until the window expires. +Automatic and first-load paths deliberately do NOT reset: only a user action +means "contact this host now", and clearing evidence the guard just collected +would defeat it. + +Call sites: + +- `retryContentInitialization` (`with-content.feature.ts`) — the Xtream + content-gate Retry button. The reset is the **first** awaited statement, before + the portal status check: a tripped guard fast-fails that check, its + `unavailable` verdict returns early, and a reset placed any later would never + run. +- `StalkerPortalDiscoveryService.discover()` — one site covering import, the Edit + dialog and lazy repair, which also guarantees a freshly edited address never + inherits a refusal recorded for the previous one. +- `retryContentPage` (`with-stalker-content.feature.ts`) — the Stalker grid + tail's append retry. +- `StalkerSearchComponent`'s search-page retry — the search results have their + own append error and retry, separate from the catalog's. +- `StalkerItvCacheService.refresh()` — the Live TV refresh button, on the same + path that already clears the cache's own error cooldown. This also covers + `refreshChannels()` in the live layout. +- Both account-info dialogs' Retry buttons (`AccountInfoComponent.reload()` for + Xtream, `StalkerAccountInfoComponent.reload()` for Stalker). Their automatic + load on open goes through a private `load()` that does not reset. +- The destructive Xtream refresh — before anything is deleted. It removes the + cached catalog and then bootstraps a re-import whose status request an open + guard would fast-fail, leaving the user with no catalog at all until the + cooldown expires. There are **two independent implementations** of this flow + and both need the reset: `PlaylistRefreshActionService.refreshXtream()` and + `RecentPlaylistsComponent.refreshXtreamPlaylist()` (the Workspace sources + page). They are near-duplicates of each other, which is exactly why the second + one was missed first time round. +- `PortalStatusService.checkPortalStatusDetails` when `skipCache` is set — the + user-initiated "Test Connection". + +Two retry paths deliberately have no reset: Xtream's `retryAppend()` is a no-op +because in-memory appends cannot fail, and the guard's own half-open trial is not +a user action. + +Where a retry clears a UI error flag, that flag is cleared **synchronously** +before awaiting the reset — otherwise the retry branch stays re-enterable and +the next `nearEnd` event fires a second retry. + +They all go through `resetHostConnectivityGuard()` +(`libs/services/src/lib/host-connectivity-reset.ts`), which holds the one rule +they share: the reset is best effort, because the guard only ever _delays_ a +request and a failed reset must not block the action that asked for it. In the +PWA the channel is unknown and `sendIpcEvent` no-ops, which is correct — +nothing there records per-host failures yet. + +## Interaction with VOD multi-source + +No exemption is needed. Multi-source resolution reaches `get_vod_info` over +`XTREAM_REQUEST`, but `VodSourceResolverService.loadVodDetails` already catches +that failure and falls back to its learned container cache or `null`, and +`probeSource` maps `null` to the verdict `unknown` — which +`VodSourceProbeCacheService` deliberately does not cache. A guard fast-fail +therefore degrades to a retryable "could not check", exactly matching the +module's contract that unreachable ≠ contacted-and-refused. The +`STREAM_PROBE_URL` reachability half never rejects and is untouched. + +## Tests + +- `apps/electron-backend/src/app/util/host-connectivity-guard.spec.ts` — the + state machine, with an injected clock. +- `apps/electron-backend/src/app/events/stalker.events.spec.ts` and + `xtream.events.spec.ts` — trip, fast-fail without contacting axios, reset, + exemption, and the absence of per-request log spam. +- `libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.utils.spec.ts` + and `stalker-portal-repair.service.spec.ts` — the message contract. diff --git a/docs/architecture/stalker-portal.md b/docs/architecture/stalker-portal.md index 5d0be80f9..31852d521 100644 --- a/docs/architecture/stalker-portal.md +++ b/docs/architecture/stalker-portal.md @@ -840,6 +840,13 @@ are logged and never retried or escalated. ## Request Transport and `cmd` Encoding +Requests to an unreachable portal are short-circuited by the main process' host +connectivity guard rather than hanging their full 15/30 s timeout again. That +guard's refusal is classified by the same message-text rules discovery uses (it +lands in the "connection-level failure" slot below), and endpoint-discovery +probes are exempt from it via `skipConnectionGuard` — see +[`host-connectivity-guard.md`](./host-connectivity-guard.md). + A real MAG/STB sends `cmd` unencoded: the portal's client JS concatenates raw `key=value` pairs, the browser URL layer escapes only what a URL cannot carry, and PHP's `$_GET` applies exactly one form-urldecode. The portal therefore sees diff --git a/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.spec.ts b/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.spec.ts index 6227ed3e2..e5fff96ee 100644 --- a/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.spec.ts +++ b/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.spec.ts @@ -17,6 +17,7 @@ import { } from '@iptvnator/services'; import { ChannelActions, PlaylistActions } from '@iptvnator/m3u-state'; import { + CONNECTIVITY_GUARD_RESET, ELECTRON_BRIDGE_SECURITY_ERROR_CODES, PLAYLIST_UPDATE, Playlist, @@ -688,4 +689,49 @@ describe('PlaylistRefreshActionService', () => { }) ); }); + + it('clears the connectivity guard before deleting the cached Xtream catalog', async () => { + // This refresh deletes the catalog and then forces a route bootstrap + // whose status request an open guard would fast-fail — leaving the user + // with no catalog at all until the cooldown expires. + const item = { + _id: 'xtream-guard', + title: 'Guarded Xtream', + serverUrl: 'http://panel.example:8080', + username: 'user', + password: 'secret', + macAddress: undefined, + portalUrl: undefined, + } as unknown as PlaylistMeta; + let confirmPromise: Promise | undefined; + const order: string[] = []; + dataService.sendIpcEvent.mockImplementation((event: string) => { + order.push(`ipc:${event}`); + return Promise.resolve({ success: true }); + }); + databaseService.deleteXtreamPlaylistContent.mockImplementation(() => { + order.push('deleteXtreamPlaylistContent'); + return Promise.resolve({ + hiddenCategories: [], + favorites: [], + recentlyViewed: [], + sourcePins: [], + }); + }); + dialogService.openConfirmDialog.mockImplementation( + ({ onConfirm }: { onConfirm?: () => Promise }) => { + confirmPromise = onConfirm?.(); + } + ); + + service.refresh(item); + await confirmPromise; + + expect(order[0]).toBe(`ipc:${CONNECTIVITY_GUARD_RESET}`); + expect(order).toContain('deleteXtreamPlaylistContent'); + expect(dataService.sendIpcEvent).toHaveBeenCalledWith( + CONNECTIVITY_GUARD_RESET, + { url: item.serverUrl } + ); + }); }); diff --git a/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.ts b/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.ts index dd911f832..97f1641e3 100644 --- a/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.ts +++ b/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.ts @@ -11,6 +11,7 @@ import { isDbAbortError, PlaybackPositionService, PlaylistRefreshService, + resetHostConnectivityGuard, RuntimeCapabilitiesService, SettingsStore, XtreamPendingRestoreService, @@ -124,6 +125,16 @@ export class PlaylistRefreshActionService { }); try { + // Before anything destructive: this refresh deletes the + // cached catalog and then forces a route bootstrap whose + // status request would be fast-failed by an open + // connectivity guard, leaving the user with no catalog at + // all until the cooldown expires. + await resetHostConnectivityGuard( + this.dataService, + item.serverUrl + ); + this.snackBar.open( this.translate.instant( 'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.STARTED' diff --git a/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.spec.ts b/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.spec.ts index b9de1ca97..13f94edaa 100644 --- a/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.spec.ts +++ b/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.spec.ts @@ -26,7 +26,11 @@ import { SortOrder, SortService, } from '@iptvnator/services'; -import { PLAYLIST_UPDATE, PlaylistMeta } from '@iptvnator/shared/interfaces'; +import { + CONNECTIVITY_GUARD_RESET, + PLAYLIST_UPDATE, + PlaylistMeta, +} from '@iptvnator/shared/interfaces'; import { RENDERER_PERFORMANCE_PHASE_HOOK_KEY, type RendererPerformancePhaseEvent, @@ -442,7 +446,11 @@ describe('RecentPlaylistsComponent busy state', () => { ); component.refreshXtreamPlaylist(item); - await Promise.resolve(); + // The refresh clears the connectivity guard before it starts deleting, + // so drain the microtask queue rather than counting exact ticks. + for (let index = 0; index < 6; index += 1) { + await Promise.resolve(); + } expect(component.isRefreshPending(item._id)).toBe(true); expect(component.getBusyMessage(item)).toBe( @@ -676,4 +684,47 @@ describe('RecentPlaylistsComponent busy state', () => { }) ); }); + + it('clears the connectivity guard before deleting the cached Xtream catalog', async () => { + // Second, independent implementation of the destructive refresh (the + // Workspace sources page). Same consequence as the shared action: the + // catalog is already gone by the time an open guard fast-fails the + // re-import bootstrap. + const item = { + _id: 'xtream-guard', + title: 'Guarded Xtream', + serverUrl: 'http://panel.example:8080', + } as PlaylistMeta; + const order: string[] = []; + let confirmPromise: Promise | undefined; + + dataService.sendIpcEvent.mockImplementation((event: string) => { + order.push(`ipc:${event}`); + return Promise.resolve({ success: true }); + }); + databaseService.deleteXtreamPlaylistContent.mockImplementation(() => { + order.push('deleteXtreamPlaylistContent'); + return Promise.resolve({ + hiddenCategories: [], + favorites: [], + recentlyViewed: [], + sourcePins: [], + }); + }); + dialogService.openConfirmDialog.mockImplementation( + ({ onConfirm }: { onConfirm?: () => Promise }) => { + confirmPromise = onConfirm?.(); + } + ); + + component.refreshXtreamPlaylist(item); + await confirmPromise; + + expect(order[0]).toBe(`ipc:${CONNECTIVITY_GUARD_RESET}`); + expect(order).toContain('deleteXtreamPlaylistContent'); + expect(dataService.sendIpcEvent).toHaveBeenCalledWith( + CONNECTIVITY_GUARD_RESET, + { url: item.serverUrl } + ); + }); }); diff --git a/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.ts b/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.ts index 5cb08a43b..0c75a608b 100644 --- a/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.ts +++ b/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.ts @@ -36,6 +36,7 @@ import { PlaybackPositionService, PlaylistDeleteActionService, PlaylistRefreshService, + resetHostConnectivityGuard, RuntimeCapabilitiesService, SortBy, SortService, @@ -366,6 +367,16 @@ export class RecentPlaylistsComponent { this.databaseService.createOperationId('xtream-refresh'); try { + // Before anything destructive, and for the same reason as + // the shared refresh action: this deletes the cached catalog + // and then navigates to re-import it, and an open + // connectivity guard would fast-fail that bootstrap — the + // user would be left with no catalog at all. + await resetHostConnectivityGuard( + this.dataService, + item.serverUrl + ); + // Show immediate feedback — deletion can take several seconds // for large playlists. this.snackBar.open( diff --git a/libs/portal/stalker/data-access/src/lib/stalker-itv-cache.service.spec.ts b/libs/portal/stalker/data-access/src/lib/stalker-itv-cache.service.spec.ts index 9f0434fbf..bc5ab51bd 100644 --- a/libs/portal/stalker/data-access/src/lib/stalker-itv-cache.service.spec.ts +++ b/libs/portal/stalker/data-access/src/lib/stalker-itv-cache.service.spec.ts @@ -1,6 +1,9 @@ import { TestBed } from '@angular/core/testing'; import { DataService } from '@iptvnator/services'; -import { PlaylistMeta } from '@iptvnator/shared/interfaces'; +import { + CONNECTIVITY_GUARD_RESET, + PlaylistMeta, +} from '@iptvnator/shared/interfaces'; import { StalkerItvCacheService } from './stalker-itv-cache.service'; import { StalkerSessionService } from './stalker-session.service'; @@ -50,7 +53,9 @@ function pageOf(items: unknown[], totalItems: number, pageSize = 14) { }; } -const UNSUPPORTED_ACTION = { js: { error: 'Unknown action: get_all_channels' } }; +const UNSUPPORTED_ACTION = { + js: { error: 'Unknown action: get_all_channels' }, +}; // Depth 10: executeStalkerRequest routes through an extra async hop for // the portal-repair pipeline, so page transitions settle a tick later. @@ -66,7 +71,12 @@ describe('StalkerItvCacheService', () => { function mockRequests(handlers: RequestHandlers): void { sendIpcEvent.mockImplementation( - async (_event: unknown, payload: unknown) => { + async (event: unknown, payload: unknown) => { + // An explicit refresh clears the main process' connectivity + // guard first; that call carries no `params`. + if (event === CONNECTIVITY_GUARD_RESET) { + return { success: true }; + } const params = (payload as { params: StalkerParams }).params; if (params['action'] === 'get_all_channels') { if (!handlers.allChannels) { @@ -282,6 +292,22 @@ describe('StalkerItvCacheService', () => { } }); + it('a refresh clears the connectivity guard before contacting the portal', async () => { + // Same reason it bypasses the local cooldown: the user asked for fresh + // channels, so a host the main process gave up on has to be contacted + // for real instead of fast-failed. + mockRequests({ + allChannels: () => pageOf([channel('1', 'News One', '5')], 1), + }); + + await service.refresh(PLAYLIST); + + expect(sendIpcEvent.mock.calls[0][0]).toBe(CONNECTIVITY_GUARD_RESET); + expect(sendIpcEvent.mock.calls[0][1]).toEqual({ + url: PLAYLIST.portalUrl, + }); + }); + it('a refresh bypasses the error cooldown', async () => { const nowSpy = jest.spyOn(Date, 'now').mockReturnValue(2_000_000); try { @@ -307,10 +333,7 @@ describe('StalkerItvCacheService', () => { it('deduplicates channels returned across crawl pages (portal ignoring the page param)', async () => { const samePage = () => pageOf( - [ - channel('1', 'News One', '5'), - channel('2', 'Sports HD', '9'), - ], + [channel('1', 'News One', '5'), channel('2', 'Sports HD', '9')], // Portal claims many items but returns the same two regardless // of `p`. 280 @@ -341,9 +364,10 @@ describe('StalkerItvCacheService', () => { await service.ensureLoaded(PLAYLIST); - expect( - service.getChannels(PLAYLIST)?.map((c) => c.id) - ).toEqual(['1', '2']); + expect(service.getChannels(PLAYLIST)?.map((c) => c.id)).toEqual([ + '1', + '2', + ]); }); it('deduplicates concurrent load requests', async () => { diff --git a/libs/portal/stalker/data-access/src/lib/stalker-itv-cache.service.ts b/libs/portal/stalker/data-access/src/lib/stalker-itv-cache.service.ts index 3400e8038..e462ccde1 100644 --- a/libs/portal/stalker/data-access/src/lib/stalker-itv-cache.service.ts +++ b/libs/portal/stalker/data-access/src/lib/stalker-itv-cache.service.ts @@ -1,6 +1,9 @@ import { Injectable, WritableSignal, inject, signal } from '@angular/core'; import { createLogger } from '@iptvnator/portal/shared/util'; -import { DataService } from '@iptvnator/services'; +import { + DataService, + resetHostConnectivityGuard, +} from '@iptvnator/services'; import { PlaylistMeta } from '@iptvnator/shared/interfaces'; import { StalkerItvChannel } from './models'; import { @@ -132,6 +135,14 @@ export class StalkerItvCacheService { return pending; } + // Same reason this clears its own error cooldown: the user asked for + // fresh channels, so a host the main process gave up on must be + // contacted for real rather than fast-failed. + await resetHostConnectivityGuard( + this.requestDeps.dataService, + playlist.portalUrl + ); + this.unsupportedKeys.delete(key); this.errorCooldownUntil.delete(key); await this.runLoad(key, playlist); diff --git a/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.service.spec.ts b/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.service.spec.ts index 32964a386..6855575de 100644 --- a/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.service.spec.ts +++ b/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.service.spec.ts @@ -1,5 +1,9 @@ import { TestBed } from '@angular/core/testing'; import { DataService } from '@iptvnator/services'; +import { + CONNECTIVITY_GUARD_RESET, + STALKER_REQUEST, +} from '@iptvnator/shared/interfaces'; import { StalkerPortalDiscoveryService } from './stalker-portal-discovery.service'; import { StalkerSessionService } from './stalker-session.service'; @@ -23,7 +27,13 @@ describe('StalkerPortalDiscoveryService', () => { function mockProbes( handlers: Record ): void { - sendIpcEvent.mockImplementation((_event, payload) => { + sendIpcEvent.mockImplementation((event, payload) => { + // Discovery clears the main process' connectivity guard first; the + // real handler answers with the outcome of that reset. + if (event === CONNECTIVITY_GUARD_RESET) { + return Promise.resolve({ success: true }); + } + const { url } = payload as { url: string }; const handler = handlers[url]; if (!handler) { @@ -36,6 +46,13 @@ describe('StalkerPortalDiscoveryService', () => { }); } + /** The portal requests only, excluding the guard reset discovery sends. */ + function probeCalls(): unknown[][] { + return sendIpcEvent.mock.calls.filter( + ([event]) => event === STALKER_REQUEST + ); + } + beforeEach(() => { sendIpcEvent = jest.fn(); authenticate = jest.fn(); @@ -66,7 +83,64 @@ describe('StalkerPortalDiscoveryService', () => { }); expect(authenticate).not.toHaveBeenCalled(); // The winning candidate ends discovery — no further probes. - expect(sendIpcEvent).toHaveBeenCalledTimes(1); + expect(probeCalls()).toHaveLength(1); + }); + + describe('host connectivity guard', () => { + it('clears the guard before probing, so an edited address is never refused', async () => { + // Import, an edited connection and lazy repair all come through + // here, and each is a deliberate "talk to this portal". + mockProbes({ + 'http://panel.example/portal.php': { + resolve: { js: [{ id: '1', title: 'News' }] }, + }, + }); + + await service.discover('http://panel.example/c', MAC); + + expect(sendIpcEvent.mock.calls[0]).toEqual([ + CONNECTIVITY_GUARD_RESET, + { url: 'http://panel.example/c' }, + ]); + }); + + it('exempts every probe from the guard', async () => { + // Candidates that do not exist fail BY DESIGN; counting those would + // let the guard declare a slow-but-alive host unreachable. + mockProbes({ + 'http://ministra.example/portal.php': { + reject: { message: 'HTTP Error: Not Found', status: 404 }, + }, + 'http://ministra.example/server/load.php': { + resolve: { js: [{ id: '1', title: 'News' }] }, + }, + }); + + await service.discover('http://ministra.example/c', MAC); + + expect(probeCalls()).toHaveLength(2); + for (const [, payload] of probeCalls()) { + expect(payload).toMatchObject({ skipConnectionGuard: true }); + } + }); + + it('probes anyway when clearing the guard fails', async () => { + mockProbes({ + 'http://panel.example/portal.php': { + resolve: { js: [{ id: '1', title: 'News' }] }, + }, + }); + sendIpcEvent.mockImplementationOnce(() => + Promise.reject(new Error('IPC unavailable')) + ); + + const outcome = await service.discover( + 'http://panel.example/c', + MAC + ); + + expect(outcome).toMatchObject({ status: 'resolved' }); + }); }); it('falls through a 404 portal.php to server/load.php and classifies by handshake', async () => { @@ -126,7 +200,7 @@ describe('StalkerPortalDiscoveryService', () => { portalUrl: 'http://ministra.example/server/load.php', isFullStalkerPortal: true, }); - expect(sendIpcEvent).toHaveBeenCalledTimes(1); + expect(probeCalls()).toHaveLength(1); }); it('classifies a token-enforcing portal.php panel as a full portal', async () => { @@ -287,7 +361,7 @@ describe('StalkerPortalDiscoveryService', () => { const outcome = await service.discover('http://down.example/c', MAC); expect(outcome).toEqual({ status: 'unreachable' }); - expect(sendIpcEvent).toHaveBeenCalledTimes(1); + expect(probeCalls()).toHaveLength(1); }); it('keeps probing past a candidate timeout — one handler can hang while siblings work', async () => { @@ -374,13 +448,16 @@ describe('StalkerPortalDiscoveryService', () => { ); const discovery = service.discover('http://slow.example/c', MAC); - // Let the probe resolve so the auth attempt actually starts. - await Promise.resolve(); - await Promise.resolve(); - await Promise.resolve(); + // Let the guard reset and the probe resolve so the auth attempt + // actually starts. + for (let i = 0; i < 6; i += 1) { + await Promise.resolve(); + } jest.advanceTimersByTime(65_000); - await Promise.resolve(); + for (let i = 0; i < 6; i += 1) { + await Promise.resolve(); + } // The abandoned attempt is cancelled, so its get_profile never // goes out to adopt the MAC's token behind a later candidate. diff --git a/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.service.ts b/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.service.ts index d3903e685..ccdf7cd92 100644 --- a/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.service.ts +++ b/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.service.ts @@ -1,5 +1,5 @@ import { Injectable, inject } from '@angular/core'; -import { DataService } from '@iptvnator/services'; +import { DataService, resetHostConnectivityGuard } from '@iptvnator/services'; import { STALKER_REQUEST } from '@iptvnator/shared/interfaces'; import { createLogger } from '@iptvnator/portal/shared/util'; import { @@ -153,6 +153,11 @@ export class StalkerPortalDiscoveryService { const candidates = buildStalkerEndpointCandidates(rawUrl); let authRejection: StalkerPortalDiscoveryRejection | null = null; + // Discovery runs on import, on an edited connection and on lazy + // repair — every one of them is a deliberate "talk to this portal" + // and must not inherit a refusal recorded for the previous address. + await resetHostConnectivityGuard(this.dataService, rawUrl); + for (const candidate of candidates) { let probeResponse: unknown; try { @@ -352,6 +357,11 @@ export class StalkerPortalDiscoveryService { // Probing absent endpoints fails BY DESIGN — the // transport services skip their error snackbar for us. silent: true, + // For the same reason these failures must not feed the + // main process' connectivity guard: several candidates on + // one host are expected to fail, and treating that as "the + // host is dead" would abandon a slow-but-alive portal. + skipConnectionGuard: true, }) ), PROBE_TIMEOUT_MS diff --git a/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.utils.spec.ts b/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.utils.spec.ts index 4c430faa7..53cd2d9c3 100644 --- a/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.utils.spec.ts +++ b/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.utils.spec.ts @@ -1,3 +1,4 @@ +import { buildHostConnectivityFastFailMessage } from '@iptvnator/shared/interfaces'; import { buildStalkerEndpointCandidates, classifyStalkerProbeResponse, @@ -5,6 +6,7 @@ import { isStalkerAuthFailureBody, isStalkerAuthFailureMessage, isStalkerAuthFailureResponse, + isStalkerProbeTimeout, legacyTransformStalkerPortalUrl, normalizeStalkerPortalInputUrl, } from './stalker-portal-discovery.utils'; @@ -367,6 +369,65 @@ describe('getStalkerRequestErrorStatus', () => { }); }); +describe('host connectivity guard fast-fail message', () => { + // The main process refuses requests to a host that stopped answering, and + // only the message survives `ipcRenderer.invoke`. These assertions are the + // contract: the wording must land in the "connection-level failure" slot, + // because that is what a tripped guard actually means. Any other reading + // has a consequence — a parsed status makes discovery walk every candidate + // and can fire lazy repair against a host we just declared dead; timeout + // wording loses the abort-early semantics; an auth phrase fires repair too. + const message = buildHostConnectivityFastFailMessage( + 'http://portal.example:8080' + ); + const ipcWrapped = new Error( + `Error invoking remote method 'STALKER_REQUEST': ${message}` + ); + + it('carries no HTTP status, bare or IPC-wrapped', () => { + expect( + getStalkerRequestErrorStatus(new Error(message)) + ).toBeUndefined(); + expect(getStalkerRequestErrorStatus(ipcWrapped)).toBeUndefined(); + }); + + it('does not read as a timeout', () => { + expect(isStalkerProbeTimeout(new Error(message))).toBe(false); + expect(isStalkerProbeTimeout(ipcWrapped)).toBe(false); + }); + + it('does not read as an authorization failure', () => { + expect(isStalkerAuthFailureMessage(message)).toBe(false); + expect(isStalkerAuthFailureMessage(ipcWrapped.message)).toBe(false); + }); + + it('survives an endpoint whose hostname reads like a classifier keyword', () => { + // The endpoint is interpolated user data: a portal at + // `https://authorization.example` would otherwise match the broad auth + // phrase set and send an unreachable host into lazy portal repair. + const keywordHost = buildHostConnectivityFastFailMessage( + 'https://authorization.example' + ); + + expect(isStalkerAuthFailureMessage(keywordHost)).toBe(false); + expect(isStalkerProbeTimeout(new Error(keywordHost))).toBe(false); + expect( + getStalkerRequestErrorStatus(new Error(keywordHost)) + ).toBeUndefined(); + + const unauthorizedHost = buildHostConnectivityFastFailMessage( + 'http://unauthorized.example:8080' + ); + expect(isStalkerAuthFailureMessage(unauthorizedHost)).toBe(false); + }); + + it('names the endpoint so the snackbar it reaches says something useful', () => { + // Scheme included: the same panel can be imported over both HTTP and + // HTTPS, and the user needs to know which one was skipped. + expect(message).toContain('http://portal.example:8080'); + }); +}); + describe('legacyTransformStalkerPortalUrl', () => { it('uses portal.php for an offline bare host', () => { expect(legacyTransformStalkerPortalUrl('http://x.example')).toBe( diff --git a/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.utils.ts b/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.utils.ts index b5e793438..57098961d 100644 --- a/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.utils.ts +++ b/libs/portal/stalker/data-access/src/lib/stalker-portal-discovery.utils.ts @@ -9,7 +9,10 @@ * it responds, instead of guessing from the URL shape (#850, #686, #755). */ -import { isStalkerAuthFailureResponse } from '@iptvnator/shared/interfaces'; +import { + isHostConnectivityFastFailMessage, + isStalkerAuthFailureResponse, +} from '@iptvnator/shared/interfaces'; /** * Candidate API endpoints for a pasted portal URL, in probe order. @@ -199,6 +202,13 @@ export function isStalkerProbeTimeout(error: unknown): boolean { } const message = String((error as { message?: unknown }).message ?? ''); + // The connectivity guard's refusal names the portal endpoint, and a + // hostname is user data. A guard refusal is a connection-level failure, so + // it must not be read as one candidate hanging while its siblings answer. + if (isHostConnectivityFastFailMessage(message)) { + return false; + } + return /timed out|timeout of \d+\s*ms|ETIMEDOUT/i.test(message); } diff --git a/libs/portal/stalker/data-access/src/lib/stalker-portal-repair.service.spec.ts b/libs/portal/stalker/data-access/src/lib/stalker-portal-repair.service.spec.ts index 849ddce26..63ca6fd44 100644 --- a/libs/portal/stalker/data-access/src/lib/stalker-portal-repair.service.spec.ts +++ b/libs/portal/stalker/data-access/src/lib/stalker-portal-repair.service.spec.ts @@ -2,7 +2,10 @@ import { TestBed } from '@angular/core/testing'; import { of, Subject, throwError } from 'rxjs'; import type { Playlist } from '@iptvnator/shared/interfaces'; import { PlaylistsService } from '@iptvnator/services'; -import { PlaylistMeta } from '@iptvnator/shared/interfaces'; +import { + PlaylistMeta, + buildHostConnectivityFastFailMessage, +} from '@iptvnator/shared/interfaces'; import { StalkerPortalDiscoveryService } from './stalker-portal-discovery.service'; import { StalkerPortalRepairService } from './stalker-portal-repair.service'; import { StalkerSessionService } from './stalker-session.service'; @@ -228,6 +231,22 @@ describe('StalkerPortalRepairService', () => { ); }); + it('never triggers on the connectivity guard refusing a dead host', () => { + // Repair means "re-probe every candidate endpoint". Doing that + // because the host stopped answering would spend a full discovery + // run on a host we already know is not there. + expect( + service.shouldAttemptRepair( + MISCLASSIFIED, + new Error( + buildHostConnectivityFastFailMessage( + 'http://portal.example:8080' + ) + ) + ) + ).toBe(false); + }); + it('never triggers for playlists without portal coordinates', () => { expect( service.shouldAttemptRepair( diff --git a/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-content.feature.spec.ts b/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-content.feature.spec.ts index 532afcffd..c9f01cc48 100644 --- a/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-content.feature.spec.ts +++ b/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-content.feature.spec.ts @@ -3,7 +3,11 @@ import { TestBed } from '@angular/core/testing'; import { patchState, signalStore, withMethods, withState } from '@ngrx/signals'; import { TranslateService } from '@ngx-translate/core'; import { DataService } from '@iptvnator/services'; -import { PlaylistMeta, StalkerPortalActions } from '@iptvnator/shared/interfaces'; +import { + CONNECTIVITY_GUARD_RESET, + PlaylistMeta, + StalkerPortalActions, +} from '@iptvnator/shared/interfaces'; import { StalkerItvChannel } from '../../models'; import { StalkerItvCacheService } from '../../stalker-itv-cache.service'; import { StalkerSessionService } from '../../stalker-session.service'; @@ -107,7 +111,9 @@ async function waitForCondition( * Controllable stand-in for StalkerItvCacheService so the legacy paged flow * can be tested in isolation and the cache-served flow deterministically. */ -function createItvCacheMock(initialChannels: StalkerItvChannel[] | null = null) { +function createItvCacheMock( + initialChannels: StalkerItvChannel[] | null = null +) { const version = signal(0); let channels = initialChannels; @@ -395,9 +401,11 @@ describe('withStalkerContent failure states', () => { store.setPage(1); await waitForCondition(() => store.getPaginatedContent().length === 3); - expect( - store.getPaginatedContent().map((item) => item.name) - ).toEqual(['Movie page 1', 'Shared Movie', 'Movie page 2']); + expect(store.getPaginatedContent().map((item) => item.name)).toEqual([ + 'Movie page 1', + 'Shared Movie', + 'Movie page 2', + ]); expect(store.hasMoreContent()).toBe(false); }); @@ -492,19 +500,27 @@ describe('withStalkerContent failure states', () => { ); // The failed append left page 1 on screen, not the empty state. - expect( - store.getPaginatedContent().map((item) => item.name) - ).toEqual(['Movie page 1']); + expect(store.getPaginatedContent().map((item) => item.name)).toEqual([ + 'Movie page 1', + ]); expect(store.contentError()).toBeNull(); failPageTwo = false; - store.retryContentPage(); + dataService.sendIpcEvent.mockClear(); + await store.retryContentPage(); await waitForCondition(() => store.getPaginatedContent().length === 2); + // Two failed appends are exactly what opens the main process' + // connectivity guard, so the retry has to clear it FIRST — otherwise + // this button fast-fails without a request and looks broken. + const [firstChannel] = dataService.sendIpcEvent.mock.calls[0]; + expect(firstChannel).toBe(CONNECTIVITY_GUARD_RESET); + expect(store.hasContentAppendError()).toBe(false); - expect( - store.getPaginatedContent().map((item) => item.name) - ).toEqual(['Movie page 1', 'Movie page 2']); + expect(store.getPaginatedContent().map((item) => item.name)).toEqual([ + 'Movie page 1', + 'Movie page 2', + ]); }); it('falls back to a synthetic all-radio category when radio categories are unavailable', async () => { @@ -650,7 +666,12 @@ describe('withStalkerContent failure states', () => { describe('withStalkerContent full ITV channel list cache', () => { const CACHED_CHANNELS: StalkerItvChannel[] = [ { id: '1', cmd: 'ffrt http://x/1', name: 'News One', tv_genre_id: '5' }, - { id: '2', cmd: 'ffrt http://x/2', name: 'Sports HD', tv_genre_id: '9' }, + { + id: '2', + cmd: 'ffrt http://x/2', + name: 'Sports HD', + tv_genre_id: '9', + }, { id: '3', cmd: 'ffrt http://x/3', name: 'News Two', tv_genre_id: '5' }, ]; diff --git a/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-content.feature.ts b/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-content.feature.ts index e4f8bf695..8119854ae 100644 --- a/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-content.feature.ts +++ b/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-content.feature.ts @@ -9,7 +9,7 @@ import { } from '@ngrx/signals'; import { TranslateService } from '@ngx-translate/core'; import { createLogger } from '@iptvnator/portal/shared/util'; -import { DataService } from '@iptvnator/services'; +import { DataService, resetHostConnectivityGuard } from '@iptvnator/services'; import { StalkerCategoryItem, StalkerContentItem, @@ -197,9 +197,7 @@ function buildEmptyContentPatch( * Portals can shift items between pages while the list is being appended — * a duplicate id would render the same card twice and break `track` hints. */ -function dedupeContentById( - items: StalkerContentItem[] -): StalkerContentItem[] { +function dedupeContentById(items: StalkerContentItem[]): StalkerContentItem[] { const seenIds = new Set(); return items.filter((item) => { const id = @@ -848,50 +846,62 @@ export function withStalkerContent() { const storeContext = store as typeof store & StalkerContentResourceStoreContract; const itvCache = inject(StalkerItvCacheService); + const dataService = inject(DataService); return { - /** - * Kicks off the full ITV channel list load as soon as the Live TV - * section is entered (instead of waiting for the first category - * click), so the all-channels view and category count badges are - * available immediately. Safe to call repeatedly — the cache - * de-duplicates in-flight loads and memoizes unsupported portals. - */ - preloadItvChannels(): void { - void itvCache.ensureLoaded(storeContext.currentPlaylist()); - }, - /** - * Re-runs the content loader with unchanged params — the retry - * for a failed append page. - */ - retryContentPage(): void { - patchState(store, { appendError: null }); - storeContext.getContentResource.reload(); - }, - async refreshItvChannels(): Promise { - await itvCache.refresh(storeContext.currentPlaylist()); - }, - setCategories( - type: StalkerContentType, - categories: StalkerCategoryItem[] - ) { - patchState(store, buildCategoryPatch(type, categories)); - }, - resetCategories() { - patchState(store, { - vodCategories: [], - seriesCategories: [], - itvCategories: [], - radioCategories: [], - categoryError: null, - }); - }, - setItvChannels(channels: StalkerItvChannel[]) { - patchState(store, { itvChannels: channels }); - }, - setRadioChannels(channels: StalkerItvChannel[]) { - patchState(store, { radioChannels: channels }); - }, + /** + * Kicks off the full ITV channel list load as soon as the Live TV + * section is entered (instead of waiting for the first category + * click), so the all-channels view and category count badges are + * available immediately. Safe to call repeatedly — the cache + * de-duplicates in-flight loads and memoizes unsupported portals. + */ + preloadItvChannels(): void { + void itvCache.ensureLoaded(storeContext.currentPlaylist()); + }, + /** + * Re-runs the content loader with unchanged params — the retry + * for a failed append page. + * + * The guard reset comes FIRST: two failed appends are exactly what + * opens the breaker, so without it this button would fast-fail + * without a request and appear to do nothing for 30 seconds. + */ + async retryContentPage(): Promise { + // Clear the flag synchronously, before the await: otherwise + // this stays re-enterable and the next `nearEnd` event + // fires a second retry. + patchState(store, { appendError: null }); + await resetHostConnectivityGuard( + dataService, + storeContext.currentPlaylist()?.portalUrl + ); + storeContext.getContentResource.reload(); + }, + async refreshItvChannels(): Promise { + await itvCache.refresh(storeContext.currentPlaylist()); + }, + setCategories( + type: StalkerContentType, + categories: StalkerCategoryItem[] + ) { + patchState(store, buildCategoryPatch(type, categories)); + }, + resetCategories() { + patchState(store, { + vodCategories: [], + seriesCategories: [], + itvCategories: [], + radioCategories: [], + categoryError: null, + }); + }, + setItvChannels(channels: StalkerItvChannel[]) { + patchState(store, { itvChannels: channels }); + }, + setRadioChannels(channels: StalkerItvChannel[]) { + patchState(store, { radioChannels: channels }); + }, }; }) ); diff --git a/libs/portal/stalker/feature/src/lib/stalker-account-info/stalker-account-info.component.spec.ts b/libs/portal/stalker/feature/src/lib/stalker-account-info/stalker-account-info.component.spec.ts index 48d2187e7..789380c13 100644 --- a/libs/portal/stalker/feature/src/lib/stalker-account-info/stalker-account-info.component.spec.ts +++ b/libs/portal/stalker/feature/src/lib/stalker-account-info/stalker-account-info.component.spec.ts @@ -4,14 +4,18 @@ import { NoopAnimationsModule } from '@angular/platform-browser/animations'; import { TranslateModule } from '@ngx-translate/core'; import { of, throwError } from 'rxjs'; import { StalkerAccountInfoService } from '@iptvnator/portal/stalker/data-access'; -import { PlaylistsService } from '@iptvnator/services'; -import { PlaylistMeta } from '@iptvnator/shared/interfaces'; +import { DataService, PlaylistsService } from '@iptvnator/services'; +import { + CONNECTIVITY_GUARD_RESET, + PlaylistMeta, +} from '@iptvnator/shared/interfaces'; import { StalkerAccountInfoComponent } from './stalker-account-info.component'; describe('StalkerAccountInfoComponent', () => { let fixture: ComponentFixture; let component: StalkerAccountInfoComponent; let accountInfoService: { fetchAccountInfo: jest.Mock }; + let dataService: { sendIpcEvent: jest.Mock }; let playlistsService: { getPlaylist: jest.Mock }; const playlist = { @@ -48,6 +52,7 @@ describe('StalkerAccountInfoComponent', () => { TranslateModule.forRoot(), ], providers: [ + { provide: DataService, useValue: dataService }, { provide: MAT_DIALOG_DATA, useValue: { playlist } }, { provide: StalkerAccountInfoService, @@ -71,6 +76,9 @@ describe('StalkerAccountInfoComponent', () => { accountInfoService = { fetchAccountInfo: jest.fn().mockResolvedValue(freshSnapshot), }; + dataService = { + sendIpcEvent: jest.fn().mockResolvedValue({ success: true }), + }; playlistsService = { getPlaylist: jest.fn().mockReturnValue(of(null)), }; @@ -207,6 +215,10 @@ describe('StalkerAccountInfoComponent', () => { TranslateModule.forRoot(), ], providers: [ + { + provide: DataService, + useValue: { sendIpcEvent: jest.fn() }, + }, { provide: MAT_DIALOG_DATA, useValue: { @@ -256,4 +268,25 @@ describe('StalkerAccountInfoComponent', () => { .find((row) => row.labelKey === 'STALKER.ACCOUNT_INFO.STATUS') ).toBeUndefined(); }); + + it('clears the connectivity guard before the Retry button re-reads the profile', async () => { + // The profile request is exactly what a tripped guard fast-fails, so a + // reset placed after it would leave Retry doing nothing for 30 seconds. + // Opening the dialog deliberately does not reset. + await createComponent(); + expect(dataService.sendIpcEvent).not.toHaveBeenCalled(); + accountInfoService.fetchAccountInfo.mockClear(); + + await component.reload(); + + expect(dataService.sendIpcEvent).toHaveBeenCalledWith( + CONNECTIVITY_GUARD_RESET, + { url: playlist.portalUrl } + ); + expect( + dataService.sendIpcEvent.mock.invocationCallOrder[0] + ).toBeLessThan( + accountInfoService.fetchAccountInfo.mock.invocationCallOrder[0] + ); + }); }); diff --git a/libs/portal/stalker/feature/src/lib/stalker-account-info/stalker-account-info.component.ts b/libs/portal/stalker/feature/src/lib/stalker-account-info/stalker-account-info.component.ts index 3d914ae6d..eb95f5d1a 100644 --- a/libs/portal/stalker/feature/src/lib/stalker-account-info/stalker-account-info.component.ts +++ b/libs/portal/stalker/feature/src/lib/stalker-account-info/stalker-account-info.component.ts @@ -17,7 +17,11 @@ import { StalkerAccountSnapshot, } from '@iptvnator/portal/stalker/data-access'; import { createLogger } from '@iptvnator/portal/shared/util'; -import { PlaylistsService } from '@iptvnator/services'; +import { + DataService, + PlaylistsService, + resetHostConnectivityGuard, +} from '@iptvnator/services'; import { PlaylistMeta, type StalkerAccountInfoDialogData, @@ -56,6 +60,7 @@ export class StalkerAccountInfoComponent { MAT_DIALOG_DATA ); private readonly accountInfoService = inject(StalkerAccountInfoService); + private readonly dataService = inject(DataService); private readonly playlistsService = inject(PlaylistsService); private readonly logger = createLogger('StalkerAccountInfo'); @@ -222,7 +227,17 @@ export class StalkerAccountInfoComponent { void this.load(); } + /** + * The template's Retry button. Clears the connectivity guard first, for the + * same reason as every other explicit retry: the profile request is what a + * tripped guard fast-fails. Opening the dialog does not reset — only a user + * action means "contact this host now". + */ async reload(): Promise { + await resetHostConnectivityGuard( + this.dataService, + this.playlist?.portalUrl + ); await this.load(); } diff --git a/libs/portal/stalker/feature/src/lib/stalker-catalog-facade.service.ts b/libs/portal/stalker/feature/src/lib/stalker-catalog-facade.service.ts index bdd768867..58f96a0b0 100644 --- a/libs/portal/stalker/feature/src/lib/stalker-catalog-facade.service.ts +++ b/libs/portal/stalker/feature/src/lib/stalker-catalog-facade.service.ts @@ -205,7 +205,7 @@ export class StalkerCatalogFacadeService implements StalkerPortalCatalogFacade< } retryAppend(): void { - this.stalkerStore.retryContentPage(); + void this.stalkerStore.retryContentPage(); } saveScrollPosition(scrollTop: number): void { diff --git a/libs/portal/stalker/feature/src/lib/stalker-search/stalker-search.component.spec.ts b/libs/portal/stalker/feature/src/lib/stalker-search/stalker-search.component.spec.ts index 6f4068d8d..3f957d545 100644 --- a/libs/portal/stalker/feature/src/lib/stalker-search/stalker-search.component.spec.ts +++ b/libs/portal/stalker/feature/src/lib/stalker-search/stalker-search.component.spec.ts @@ -19,6 +19,7 @@ import { } from '@iptvnator/portal/stalker/data-access'; import { createPlaybackSessionKey } from '@iptvnator/playback/util'; import { DataService, PlaylistsService } from '@iptvnator/services'; +import { CONNECTIVITY_GUARD_RESET } from '@iptvnator/shared/interfaces'; import type { ResolvedPortalPlayback } from '@iptvnator/shared/interfaces'; import { StalkerSearchComponent } from './stalker-search.component'; @@ -43,6 +44,12 @@ class StubStalkerInlineDetailComponent { readonly inlinePlaybackClosed = output(); } +async function flushMicrotasks(): Promise { + for (let index = 0; index < 4; index += 1) { + await Promise.resolve(); + } +} + describe('StalkerSearchComponent playback session key', () => { let fixture: ComponentFixture; const playlist = signal({ @@ -274,6 +281,7 @@ describe('StalkerSearchComponent playback session key', () => { describe('StalkerSearchComponent result paging', () => { let component: StalkerSearchComponent; + let dataService: { sendIpcEvent: jest.Mock }; const activePlaylist = signal({ _id: 'playlist|one', title: 'Search portal', @@ -289,6 +297,9 @@ describe('StalkerSearchComponent result paging', () => { } beforeEach(() => { + dataService = { + sendIpcEvent: jest.fn().mockResolvedValue({ success: true }), + }; activePlaylist.set({ _id: 'playlist|one', title: 'Search portal', @@ -309,7 +320,7 @@ describe('StalkerSearchComponent result paging', () => { }, }, { provide: Location, useValue: { back: jest.fn() } }, - { provide: DataService, useValue: {} }, + { provide: DataService, useValue: dataService }, { provide: PlaylistContextFacade, useValue: { activePlaylist }, @@ -406,7 +417,7 @@ describe('StalkerSearchComponent result paging', () => { expect(component.searchHasMore()).toBe(false); }); - it('keeps accumulated pages on a failed append and retries the SAME page', () => { + it('keeps accumulated pages on a failed append and retries the SAME page', async () => { component.applySearchPageSuccess(1, searchItems('page1', 3), 6); expect(component.searchHasMore()).toBe(true); @@ -429,8 +440,19 @@ describe('StalkerSearchComponent result paging', () => { const pageBefore = component.searchPage(); component.loadMoreSearchResults(); expect(component.searchPage()).toBe(pageBefore); + // Cleared synchronously, so a second near-end cannot re-enter the retry + // while the connectivity-guard reset is still in flight. expect(component.searchAppendError()).toBe(false); + + // The reload itself waits for that reset: a tripped guard would + // otherwise fast-fail this retry without contacting the portal. + expect(reload).not.toHaveBeenCalled(); + await flushMicrotasks(); expect(reload).toHaveBeenCalledTimes(1); + expect(dataService.sendIpcEvent).toHaveBeenCalledWith( + CONNECTIVITY_GUARD_RESET, + { url: activePlaylist().portalUrl } + ); // With the error cleared, the following near-end advances normally. component.loadMoreSearchResults(); @@ -514,11 +536,7 @@ describe('StalkerSearchComponent result paging', () => { configurable: true, value: { isLoading: () => false, reload: jest.fn(() => true) }, }); - component.applySearchPageSuccess( - 1, - searchItems('portalA', 3), - 6 - ); + component.applySearchPageSuccess(1, searchItems('portalA', 3), 6); component.loadMoreSearchResults(); expect(component.searchPage()).toBe(2); diff --git a/libs/portal/stalker/feature/src/lib/stalker-search/stalker-search.component.ts b/libs/portal/stalker/feature/src/lib/stalker-search/stalker-search.component.ts index 8f055dea1..8cde3321d 100644 --- a/libs/portal/stalker/feature/src/lib/stalker-search/stalker-search.component.ts +++ b/libs/portal/stalker/feature/src/lib/stalker-search/stalker-search.component.ts @@ -22,7 +22,11 @@ import { StalkerPortalRepairService, StalkerSessionService, } from '@iptvnator/portal/stalker/data-access'; -import { DataService, PlaylistsService } from '@iptvnator/services'; +import { + DataService, + PlaylistsService, + resetHostConnectivityGuard, +} from '@iptvnator/services'; import { PlaybackPositionData, ResolvedPortalPlayback, @@ -404,14 +408,30 @@ export class StalkerSearchComponent { if (this.searchAppendError()) { // Retry the SAME page — advancing would permanently omit it. - this.searchAppendError.set(false); - this.searchResultsResource.reload(); + void this.retrySearchPage(); return; } this.searchPage.update((page) => page + 1); } + /** + * Two failed search pages are exactly what opens the main process' + * connectivity guard, so the reset has to precede the reload — otherwise + * this retry fast-fails without contacting a portal that may have + * recovered, and keeps repeating the same error until the window expires. + */ + private async retrySearchPage(): Promise { + // Clear the flag synchronously: awaiting first would leave this branch + // re-enterable, and the next `nearEnd` event would fire a second retry. + this.searchAppendError.set(false); + await resetHostConnectivityGuard( + this.dataService, + this.currentPlaylist()?.portalUrl + ); + this.searchResultsResource.reload(); + } + readonly isSelectedVodFavorite = signal(false); constructor() { diff --git a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.connectivity-guard.spec.ts b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.connectivity-guard.spec.ts new file mode 100644 index 000000000..d89a96c51 --- /dev/null +++ b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.connectivity-guard.spec.ts @@ -0,0 +1,125 @@ +import { TestBed } from '@angular/core/testing'; +import { CONNECTIVITY_GUARD_RESET } from '@iptvnator/shared/interfaces'; +import { PortalStatusType } from '../../xtream-state'; +import { + createContentTestProviders, + createContentTestStore, + createPendingRestoreServiceMock, + TEST_PLAYLIST, +} from './with-content.feature.spec-helpers'; + +jest.mock('@iptvnator/portal/shared/util', () => ({ + createLogger: () => ({ + debug: jest.fn(), + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + }), +})); + +let checkPortalStatusMock: jest.Mock, []>; + +const TestContentStore = createContentTestStore(() => checkPortalStatusMock()); + +describe('withContent host connectivity guard', () => { + let store: InstanceType; + let dataService: { sendIpcEvent: jest.Mock }; + + beforeEach(() => { + localStorage.clear(); + checkPortalStatusMock = jest.fn().mockResolvedValue('active'); + dataService = { + sendIpcEvent: jest.fn().mockResolvedValue({ success: true }), + }; + + TestBed.configureTestingModule({ + providers: createContentTestProviders(TestContentStore, { + dataSource: { + getCategories: jest.fn().mockResolvedValue([]), + getCachedCategories: jest.fn().mockResolvedValue([]), + getContent: jest.fn().mockResolvedValue([]), + getCachedContent: jest.fn().mockResolvedValue([]), + hasCategories: jest.fn().mockResolvedValue(true), + hasContent: jest.fn().mockResolvedValue(false), + restoreUserData: jest.fn().mockResolvedValue(undefined), + }, + databaseService: { + clearXtreamImportCache: jest.fn().mockResolvedValue(true), + cancelOperation: jest.fn().mockResolvedValue(true), + createOperationId: jest + .fn() + .mockImplementation( + (prefix?: string) => `${prefix ?? 'db-op'}-1` + ), + getXtreamImportStatus: jest + .fn() + .mockResolvedValue('completed'), + setXtreamImportStatus: jest.fn().mockResolvedValue(true), + supportsDbOperationCancellation: jest + .fn() + .mockReturnValue(true), + }, + xtreamApiService: { + cancelSession: jest.fn().mockResolvedValue(true), + }, + pendingRestoreService: createPendingRestoreServiceMock(), + dataService, + }), + }); + + store = TestBed.inject(TestContentStore); + }); + + afterEach(() => { + localStorage.clear(); + }); + + it('clears the guard before the portal status check the retry depends on', async () => { + // Ordering is the whole point. A tripped guard fast-fails that status + // check, which resolves to 'unavailable' and returns early — so a reset + // placed after it would never run, and the Retry button would silently + // do nothing until the guard's window expired on its own. + const order: string[] = []; + dataService.sendIpcEvent.mockImplementation((event: string) => { + order.push(`ipc:${event}`); + return Promise.resolve({ success: true }); + }); + checkPortalStatusMock.mockImplementation(() => { + order.push('checkPortalStatus'); + return Promise.resolve('unavailable' as PortalStatusType); + }); + + await store.retryContentInitialization(); + + expect(order).toEqual([ + `ipc:${CONNECTIVITY_GUARD_RESET}`, + 'checkPortalStatus', + ]); + expect(dataService.sendIpcEvent).toHaveBeenCalledWith( + CONNECTIVITY_GUARD_RESET, + { url: TEST_PLAYLIST.serverUrl } + ); + }); + + it('retries anyway when clearing the guard fails', async () => { + dataService.sendIpcEvent.mockRejectedValue( + new Error('IPC unavailable') + ); + + await store.retryContentInitialization(); + + expect(checkPortalStatusMock).toHaveBeenCalled(); + expect(store.isContentInitialized()).toBe(true); + }); + + it('leaves the guard alone for an ordinary initialization', async () => { + // Only a user-driven retry means "contact this host now"; the automatic + // first load must not clear evidence the guard just collected. + await store.initializeContent(); + + expect(dataService.sendIpcEvent).not.toHaveBeenCalledWith( + CONNECTIVITY_GUARD_RESET, + expect.anything() + ); + }); +}); diff --git a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.import-phase.spec.ts b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.import-phase.spec.ts index 1da2256ae..0b468d510 100644 --- a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.import-phase.spec.ts +++ b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.import-phase.spec.ts @@ -1,13 +1,8 @@ import { TestBed } from '@angular/core/testing'; -import { - DatabaseService, - XtreamPendingRestoreService, -} from '@iptvnator/services'; -import { XTREAM_DATA_SOURCE } from '../../data-sources/xtream-data-source.interface'; -import { XtreamApiService } from '../../services/xtream-api.service'; import { PortalStatusType } from '../../xtream-state'; import { createAbortError, + createContentTestProviders, createContentTestStore, createDeferred, createPendingRestoreServiceMock, @@ -87,25 +82,15 @@ describe('withContent loading-cached phase', () => { checkPortalStatusMock = jest.fn().mockResolvedValue('active'); TestBed.configureTestingModule({ - providers: [ - TestContentStore, - { - provide: XTREAM_DATA_SOURCE, - useValue: dataSource, + providers: createContentTestProviders(TestContentStore, { + dataSource, + databaseService, + xtreamApiService: { + cancelSession: jest.fn().mockResolvedValue(true), }, - { - provide: DatabaseService, - useValue: databaseService, - }, - { - provide: XtreamApiService, - useValue: { cancelSession: jest.fn().mockResolvedValue(true) }, - }, - { - provide: XtreamPendingRestoreService, - useValue: createPendingRestoreServiceMock(), - }, - ], + pendingRestoreService: createPendingRestoreServiceMock(), + dataService: { sendIpcEvent: jest.fn() }, + }), }); store = TestBed.inject(TestContentStore); diff --git a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec-helpers.ts b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec-helpers.ts index 8f4d2c204..27b3d631e 100644 --- a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec-helpers.ts +++ b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec-helpers.ts @@ -4,8 +4,17 @@ import { RENDERER_PERFORMANCE_PHASE_HOOK_KEY, type RendererPerformancePhaseEvent, } from '@iptvnator/shared/logging'; -import { XtreamPendingRestoreSnapshot } from '@iptvnator/services'; -import { XtreamPlaylistData } from '../../data-sources/xtream-data-source.interface'; +import { + DataService, + DatabaseService, + XtreamPendingRestoreService, + XtreamPendingRestoreSnapshot, +} from '@iptvnator/services'; +import { + XTREAM_DATA_SOURCE, + XtreamPlaylistData, +} from '../../data-sources/xtream-data-source.interface'; +import { XtreamApiService } from '../../services/xtream-api.service'; import { PortalStatusType } from '../../xtream-state'; import { withContent } from './with-content.feature'; @@ -48,10 +57,7 @@ export function createPendingRestoreServiceMock() { } const serializedState = JSON.stringify(state); - if ( - !activeSnapshot || - activeSerializedState !== serializedState - ) { + if (!activeSnapshot || activeSerializedState !== serializedState) { activeSnapshot = { playlistId, revision: ++nextRevision, @@ -68,16 +74,13 @@ export function createPendingRestoreServiceMock() { expectedSnapshot: XtreamPendingRestoreSnapshot, apply: (state: XtreamPendingRestoreState) => Promise ) => { - const currentSnapshot = - mock.getSnapshotOrThrow(playlistId); + const currentSnapshot = mock.getSnapshotOrThrow(playlistId); if (!currentSnapshot) { return lastConsumedRevision === expectedSnapshot.revision ? 'consumed' : 'superseded'; } - if ( - currentSnapshot.revision !== expectedSnapshot.revision - ) { + if (currentSnapshot.revision !== expectedSnapshot.revision) { return 'superseded'; } @@ -116,6 +119,34 @@ export function createContentTestStore( ); } +/** + * The provider list every `withContent` spec needs. Shared so the dependency + * set lives in one place: the feature injects all five, and a spec that omits + * one fails with a NullInjector error rather than a useful message. + */ +export function createContentTestProviders( + store: unknown, + mocks: { + dataSource: unknown; + databaseService: unknown; + xtreamApiService: unknown; + pendingRestoreService: unknown; + dataService: unknown; + } +) { + return [ + store, + { provide: XTREAM_DATA_SOURCE, useValue: mocks.dataSource }, + { provide: DatabaseService, useValue: mocks.databaseService }, + { provide: XtreamApiService, useValue: mocks.xtreamApiService }, + { + provide: XtreamPendingRestoreService, + useValue: mocks.pendingRestoreService, + }, + { provide: DataService, useValue: mocks.dataService }, + ]; +} + const performanceHookSymbol = Symbol.for(RENDERER_PERFORMANCE_PHASE_HOOK_KEY); export function setPerformanceHook( diff --git a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec.ts b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec.ts index 8acb45a06..6dd3ec5c5 100644 --- a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec.ts +++ b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec.ts @@ -1,13 +1,8 @@ import { TestBed } from '@angular/core/testing'; -import { - DatabaseService, - XtreamPendingRestoreService, -} from '@iptvnator/services'; -import { XTREAM_DATA_SOURCE } from '../../data-sources/xtream-data-source.interface'; -import { XtreamApiService } from '../../services/xtream-api.service'; import { PortalStatusType } from '../../xtream-state'; import { createAbortError, + createContentTestProviders, createContentTestStore, createDeferred, createPendingRestoreServiceMock, @@ -60,6 +55,7 @@ describe('withContent import state', () => { let pendingRestoreService: ReturnType< typeof createPendingRestoreServiceMock >; + let dataService: { sendIpcEvent: jest.Mock }; beforeEach(() => { localStorage.clear(); @@ -95,28 +91,19 @@ describe('withContent import state', () => { cancelSession: jest.fn().mockResolvedValue(true), }; pendingRestoreService = createPendingRestoreServiceMock(); + dataService = { + sendIpcEvent: jest.fn().mockResolvedValue({ success: true }), + }; checkPortalStatusMock = jest.fn().mockResolvedValue('active'); TestBed.configureTestingModule({ - providers: [ - TestContentStore, - { - provide: XTREAM_DATA_SOURCE, - useValue: dataSource, - }, - { - provide: DatabaseService, - useValue: databaseService, - }, - { - provide: XtreamApiService, - useValue: xtreamApiService, - }, - { - provide: XtreamPendingRestoreService, - useValue: pendingRestoreService, - }, - ], + providers: createContentTestProviders(TestContentStore, { + dataSource, + databaseService, + xtreamApiService, + pendingRestoreService, + dataService, + }), }); store = TestBed.inject(TestContentStore); diff --git a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.ts b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.ts index 6e30ed48c..ac222a5ac 100644 --- a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.ts +++ b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.ts @@ -18,9 +18,11 @@ import { } from '@iptvnator/shared/logging'; import { createLogger } from '@iptvnator/portal/shared/util'; import { + DataService, DatabaseService, DbOperationEvent, isDbAbortError, + resetHostConnectivityGuard, XtreamPendingRestoreService, XtreamImportStatus, } from '@iptvnator/services'; @@ -217,6 +219,7 @@ export function withContent() { withMethods((store) => { const dataSource = inject(XTREAM_DATA_SOURCE); + const dataService = inject(DataService); const databaseService = inject(DatabaseService); const pendingRestoreService = inject(XtreamPendingRestoreService); const xtreamApiService = inject(XtreamApiService); @@ -957,8 +960,7 @@ export function withContent() { pendingState, { onEvent: trackImportEvent, - operationId: - restoreOperationId, + operationId: restoreOperationId, } ); throwIfImportCancelled(importSessionId); @@ -1102,7 +1104,8 @@ export function withContent() { 'live', { sessionId: options?.sessionId, - onPhaseChange: publishTypedImportPhase('live'), + onPhaseChange: + publishTypedImportPhase('live'), } ), dataSource.getCategories( @@ -1111,7 +1114,8 @@ export function withContent() { 'vod', { sessionId: options?.sessionId, - onPhaseChange: publishTypedImportPhase('vod'), + onPhaseChange: + publishTypedImportPhase('vod'), } ), dataSource.getCategories( @@ -1120,7 +1124,8 @@ export function withContent() { 'series', { sessionId: options?.sessionId, - onPhaseChange: publishTypedImportPhase('series'), + onPhaseChange: + publishTypedImportPhase('series'), } ), ]); @@ -1278,7 +1283,8 @@ export function withContent() { operationId: seriesOperationId, sessionId: options?.sessionId, onEvent: trackImportEvent, - onPhaseChange: publishTypedImportPhase('series'), + onPhaseChange: + publishTypedImportPhase('series'), } )) as XtreamSerieItem[]; throwIfImportCancelled(options?.importSessionId); @@ -1368,6 +1374,16 @@ export function withContent() { }, async retryContentInitialization(): Promise { + // FIRST, before the status check below: that check is the + // one request a tripped connectivity guard would fast-fail, + // and its 'unavailable' verdict returns early — so a reset + // placed any later would never run and this button would + // silently do nothing for the guard's whole window. + await resetHostConnectivityGuard( + dataService, + getCredentialsFromStore()?.credentials.serverUrl + ); + const portalStatus = (await getPortalStore().checkPortalStatus?.()) ?? getPortalStore().portalStatus?.() ?? diff --git a/libs/portal/xtream/data-access/src/lib/stores/xtream.store.spec.ts b/libs/portal/xtream/data-access/src/lib/stores/xtream.store.spec.ts index 5feb04d4e..f56a199b7 100644 --- a/libs/portal/xtream/data-access/src/lib/stores/xtream.store.spec.ts +++ b/libs/portal/xtream/data-access/src/lib/stores/xtream.store.spec.ts @@ -1,5 +1,6 @@ import { TestBed } from '@angular/core/testing'; import { + DataService, DatabaseService, PlaybackPositionRuntimeBridgeService, PlaylistsService, @@ -96,6 +97,7 @@ function configureStore( useValue: { isEnabled: jest.fn(() => false) }, }, { provide: XtreamPendingRestoreService, useValue: {} }, + { provide: DataService, useValue: { sendIpcEvent: jest.fn() } }, { provide: XtreamUrlService, useValue: {} }, { provide: XtreamXmltvFallbackService, useValue: {} }, { provide: PORTAL_PLAYER, useValue: {} }, diff --git a/libs/portal/xtream/feature/src/lib/account-info/account-info.component.spec.ts b/libs/portal/xtream/feature/src/lib/account-info/account-info.component.spec.ts index 9ee8d80c6..197cab058 100644 --- a/libs/portal/xtream/feature/src/lib/account-info/account-info.component.spec.ts +++ b/libs/portal/xtream/feature/src/lib/account-info/account-info.component.spec.ts @@ -8,6 +8,8 @@ import { XtreamApiService, XtreamStore, } from '@iptvnator/portal/xtream/data-access'; +import { DataService } from '@iptvnator/services'; +import { CONNECTIVITY_GUARD_RESET } from '@iptvnator/shared/interfaces'; import { AccountInfoComponent } from './account-info.component'; describe('AccountInfoComponent', () => { @@ -17,6 +19,7 @@ describe('AccountInfoComponent', () => { getAccountInfo: jest.Mock; }; let currentPlaylist: WritableSignal; + let dataService: { sendIpcEvent: jest.Mock }; beforeEach(async () => { xtreamApiService = { @@ -35,6 +38,9 @@ describe('AccountInfoComponent', () => { }), }; currentPlaylist = signal(null); + dataService = { + sendIpcEvent: jest.fn().mockResolvedValue({ success: true }), + }; await TestBed.configureTestingModule({ imports: [ @@ -43,6 +49,10 @@ describe('AccountInfoComponent', () => { TranslateModule.forRoot(), ], providers: [ + { + provide: DataService, + useValue: dataService, + }, { provide: MAT_DIALOG_DATA, useValue: { @@ -112,4 +122,24 @@ describe('AccountInfoComponent', () => { expect(component.isActive()).toBe(true); expect(component.userDetails()[0]?.tone).toBe('positive'); }); + + it('clears the connectivity guard before the Retry button re-reads the account', async () => { + // The account request is exactly what a tripped guard fast-fails, so a + // reset placed after it would leave Retry doing nothing for 30 seconds. + // The automatic first load deliberately does not reset. + expect(dataService.sendIpcEvent).not.toHaveBeenCalled(); + xtreamApiService.getAccountInfo.mockClear(); + + await component.reload(); + + expect(dataService.sendIpcEvent).toHaveBeenCalledWith( + CONNECTIVITY_GUARD_RESET, + { url: 'https://dialog.example.test' } + ); + expect( + dataService.sendIpcEvent.mock.invocationCallOrder[0] + ).toBeLessThan( + xtreamApiService.getAccountInfo.mock.invocationCallOrder[0] + ); + }); }); diff --git a/libs/portal/xtream/feature/src/lib/account-info/account-info.component.ts b/libs/portal/xtream/feature/src/lib/account-info/account-info.component.ts index 1611fbf82..bbe911862 100644 --- a/libs/portal/xtream/feature/src/lib/account-info/account-info.component.ts +++ b/libs/portal/xtream/feature/src/lib/account-info/account-info.component.ts @@ -15,6 +15,10 @@ import { XtreamStore, } from '@iptvnator/portal/xtream/data-access'; import { createLogger } from '@iptvnator/portal/shared/util'; +import { + DataService, + resetHostConnectivityGuard, +} from '@iptvnator/services'; import { resolveXtreamPortalStatus, type XtreamAccountInfoDialogData, @@ -54,6 +58,7 @@ export class AccountInfoComponent { inject(MAT_DIALOG_DATA, { optional: true, }) ?? {}; + private readonly dataService = inject(DataService); private readonly xtreamApiService = inject(XtreamApiService); private readonly xtreamStore = inject(XtreamStore); private readonly logger = createLogger('XtreamAccountInfo'); @@ -217,10 +222,25 @@ export class AccountInfoComponent { ]); constructor() { - void this.reload(); + void this.load(); } + /** + * The template's Retry button. Clears the main process' connectivity guard + * first: the account request is exactly what a tripped guard fast-fails, so + * without this the button would do nothing for the guard's whole window. + * The automatic first load deliberately does not reset — only an explicit + * user action means "contact this host now". + */ async reload(): Promise { + await resetHostConnectivityGuard( + this.dataService, + this.currentPlaylist()?.serverUrl + ); + await this.load(); + } + + private async load(): Promise { const playlist = this.currentPlaylist(); if (!playlist?.serverUrl || !playlist.username || !playlist.password) { diff --git a/libs/services/src/index.ts b/libs/services/src/index.ts index 162cc9a0b..b0e108bae 100644 --- a/libs/services/src/index.ts +++ b/libs/services/src/index.ts @@ -1,6 +1,7 @@ export * from './lib/catalog-title-match.service'; export * from './lib/cross-portal-similar.service'; export * from './lib/data.service'; +export * from './lib/host-connectivity-reset'; export * from './lib/database-electron.service'; export * from './lib/downloads.service'; export * from './lib/playback-position-runtime-bridge.service'; diff --git a/libs/services/src/lib/host-connectivity-reset.spec.ts b/libs/services/src/lib/host-connectivity-reset.spec.ts new file mode 100644 index 000000000..a35568949 --- /dev/null +++ b/libs/services/src/lib/host-connectivity-reset.spec.ts @@ -0,0 +1,52 @@ +import { CONNECTIVITY_GUARD_RESET } from '@iptvnator/shared/interfaces'; +import { resetHostConnectivityGuard } from './host-connectivity-reset'; + +describe('resetHostConnectivityGuard', () => { + let sendIpcEvent: jest.Mock; + + beforeEach(() => { + sendIpcEvent = jest.fn().mockResolvedValue({ success: true }); + }); + + it('asks the main process to forget the failures recorded for the host', async () => { + await resetHostConnectivityGuard( + { sendIpcEvent }, + 'http://portal.example:8080' + ); + + expect(sendIpcEvent).toHaveBeenCalledWith(CONNECTIVITY_GUARD_RESET, { + url: 'http://portal.example:8080', + }); + }); + + it.each([undefined, null, ''])( + 'sends nothing when there is no URL to reset (%p)', + async (url) => { + // A playlist can be missing its address entirely; there is no host + // to clear then, and an empty reset would be a pointless round-trip. + await resetHostConnectivityGuard({ sendIpcEvent }, url); + + expect(sendIpcEvent).not.toHaveBeenCalled(); + } + ); + + it('never propagates a transport failure to the caller', async () => { + // The guard only ever delays a request, so a failed reset must not + // block the retry, discovery run or status check that asked for it. + sendIpcEvent.mockRejectedValue(new Error('IPC unavailable')); + + await expect( + resetHostConnectivityGuard({ sendIpcEvent }, 'http://x.example') + ).resolves.toBeUndefined(); + }); + + it('tolerates the PWA no-op, where the channel is unknown', async () => { + // PwaService.sendIpcEvent returns undefined synchronously for channels + // it does not implement; nothing there records per-host failures yet. + sendIpcEvent.mockReturnValue(undefined); + + await expect( + resetHostConnectivityGuard({ sendIpcEvent }, 'http://x.example') + ).resolves.toBeUndefined(); + }); +}); diff --git a/libs/services/src/lib/host-connectivity-reset.ts b/libs/services/src/lib/host-connectivity-reset.ts new file mode 100644 index 000000000..a05c61a35 --- /dev/null +++ b/libs/services/src/lib/host-connectivity-reset.ts @@ -0,0 +1,33 @@ +import { CONNECTIVITY_GUARD_RESET } from '@iptvnator/shared/interfaces'; +import type { DataService } from './data.service'; + +/** + * Asks the main process to forget the connection failures it recorded for the + * host `url` points at, so the next request contacts it for real instead of + * being fast-failed by the connectivity guard. + * + * Send this whenever the user asked for a real attempt (portal retry, "test + * connection") or handed over an address that may now point somewhere else + * (import, edited connection, lazy repair). + * + * Always best effort: the guard only ever delays a request, so a failure here + * must never block the action that asked for it. In the PWA the channel is + * unknown and `sendIpcEvent` no-ops, which is correct — nothing there records + * per-host failures yet. + */ +export async function resetHostConnectivityGuard( + dataService: Pick, + url: string | null | undefined +): Promise { + if (!url) { + return; + } + + try { + await Promise.resolve( + dataService.sendIpcEvent(CONNECTIVITY_GUARD_RESET, { url }) + ); + } catch { + // Deliberately silent: the caller's own request reports the real state. + } +} diff --git a/libs/services/src/lib/portal-status.service.spec.ts b/libs/services/src/lib/portal-status.service.spec.ts index 4541ffe5a..bb149a9a8 100644 --- a/libs/services/src/lib/portal-status.service.spec.ts +++ b/libs/services/src/lib/portal-status.service.spec.ts @@ -4,6 +4,7 @@ import { createEnvironmentInjector, runInInjectionContext, } from '@angular/core'; +import { CONNECTIVITY_GUARD_RESET } from '@iptvnator/shared/interfaces'; import { DataService } from './data.service'; import { PortalStatusService } from './portal-status.service'; @@ -252,4 +253,56 @@ describe('PortalStatusService', () => { service.checkPortalStatus('https://example.com', 'user', 'pass') ).resolves.toBe('active'); }); + + describe('host connectivity guard', () => { + beforeEach(() => { + dataService.sendIpcEvent.mockResolvedValue({ + payload: { user_info: { auth: 1, status: 'Active' } }, + }); + }); + + it('clears the guard first when the user asked to test the connection', async () => { + // Otherwise "Test Connection" on a panel that just came back would + // report it unavailable without contacting it at all. + await service.checkPortalStatusDetails( + 'http://example.com', + 'user', + 'pass', + { skipCache: true } + ); + + expect(dataService.sendIpcEvent.mock.calls[0]).toEqual([ + CONNECTIVITY_GUARD_RESET, + { url: 'http://example.com' }, + ]); + }); + + it('leaves the guard alone for passive background checks', async () => { + await service.checkPortalStatusDetails( + 'http://example.com', + 'user', + 'pass' + ); + + expect(dataService.sendIpcEvent).not.toHaveBeenCalledWith( + CONNECTIVITY_GUARD_RESET, + expect.anything() + ); + }); + + it('still runs the check when clearing the guard fails', async () => { + dataService.sendIpcEvent.mockRejectedValueOnce( + new Error('IPC unavailable') + ); + + await expect( + service.checkPortalStatusDetails( + 'http://example.com', + 'user', + 'pass', + { skipCache: true } + ) + ).resolves.toMatchObject({ status: 'active' }); + }); + }); }); diff --git a/libs/services/src/lib/portal-status.service.ts b/libs/services/src/lib/portal-status.service.ts index 7074af1dc..a6c0c9ccc 100644 --- a/libs/services/src/lib/portal-status.service.ts +++ b/libs/services/src/lib/portal-status.service.ts @@ -6,13 +6,10 @@ import { XtreamPortalStatusResponseLike, } from '@iptvnator/shared/interfaces'; import { DataService } from './data.service'; +import { resetHostConnectivityGuard } from './host-connectivity-reset'; export type PortalStatus = - | 'active' - | 'inactive' - | 'expired' - | 'unavailable' - | 'checking'; + 'active' | 'inactive' | 'expired' | 'unavailable' | 'checking'; /** * Status plus the parsed account expiration from the same round-trip the @@ -135,6 +132,14 @@ export class PortalStatusService { if (pending) { return pending; } + } else { + // Skipping the cache means the user asked to test this portal + // right now, so the main process must forget any connection + // failures it recorded for the host and contact it for real. + await resetHostConnectivityGuard( + this.dataService, + connection.serverUrl + ); } const request = this.fetchPortalStatus( diff --git a/libs/shared/interfaces/src/index.ts b/libs/shared/interfaces/src/index.ts index 1abab7749..ee3d210a9 100644 --- a/libs/shared/interfaces/src/index.ts +++ b/libs/shared/interfaces/src/index.ts @@ -14,6 +14,7 @@ export * from './lib/epg-program.model'; export * from './lib/external-player-arguments.utils'; export * from './lib/external-player-session.interface'; export * from './lib/global-search-result.interface'; +export * from './lib/host-connectivity.util'; export * from './lib/indexed-db.config'; export * from './lib/ipc-command.class'; export * from './lib/ipc-commands'; diff --git a/libs/shared/interfaces/src/lib/electron-api.interface.ts b/libs/shared/interfaces/src/lib/electron-api.interface.ts index 87eb3e982..03ff1594b 100644 --- a/libs/shared/interfaces/src/lib/electron-api.interface.ts +++ b/libs/shared/interfaces/src/lib/electron-api.interface.ts @@ -283,6 +283,13 @@ export interface ElectronBridgeStalkerRequestPayload { token?: string; serialNumber?: string; requestId?: string; + /** + * Endpoint-discovery probes only: exempt this request from the main + * process' per-host connectivity guard. Discovery walks several candidate + * paths on one host and expects most of them to fail, so its failures must + * neither count towards the guard nor be fast-failed by it. + */ + skipConnectionGuard?: boolean; } export interface ElectronBridgeXtreamRequestPayload { @@ -774,6 +781,11 @@ export interface ElectronBridgeApi { stalkerRequest: ( payload: ElectronBridgeStalkerRequestPayload ) => Promise>; + /** + * Forgets the connection failures recorded for the host `url` points at, so + * the next request contacts it for real instead of being fast-failed. + */ + resetHostConnectivityGuard: (url: string) => Promise; xtreamRequest: ( payload: ElectronBridgeXtreamRequestPayload ) => Promise; diff --git a/libs/shared/interfaces/src/lib/host-connectivity.util.ts b/libs/shared/interfaces/src/lib/host-connectivity.util.ts new file mode 100644 index 000000000..625df0519 --- /dev/null +++ b/libs/shared/interfaces/src/lib/host-connectivity.util.ts @@ -0,0 +1,49 @@ +/** + * Cross-process contract for the main-process host connectivity guard. + * + * The guard fast-fails requests to a portal host that just failed to answer + * several times in a row, instead of hanging the full axios timeout again. + * + * Only the MESSAGE of a rejected IPC handler survives the trip to the + * renderer — `ipcRenderer.invoke` strips custom properties and re-wraps the + * value — and the Stalker renderer classifies transport failures purely from + * that text. The wording below is therefore a contract, not a label: + * + * - no `HTTP Error `: `getStalkerRequestErrorStatus` would report a + * status, so endpoint discovery would read "this endpoint is absent, probe + * the next candidate", and 404/401/403 would additionally fire lazy portal + * repair against a host we just declared dead; + * - no timeout wording: `isStalkerProbeTimeout` would keep discovery walking + * every candidate instead of aborting on the first one; + * - none of the auth phrases `isStalkerAuthFailureMessage` accepts (a bare + * `authorization` is one of them) — that would also trigger lazy repair. + * + * What is left is the "connection-level failure" slot the renderer already + * has for ECONNREFUSED/ENOTFOUND: discovery stops probing and reports the + * host unreachable, which is exactly what a tripped guard means. The message + * also reaches error snackbars verbatim, so it has to read like a sentence. + * + * `endpoint` is the request origin (scheme, host and port), not a bare host: a + * user with the same panel imported over both HTTP and HTTPS needs to see which + * of the two was skipped. A scheme is safe here — none of the forbidden + * substrings above can appear in one. + */ +export function buildHostConnectivityFastFailMessage(endpoint: string): string { + return `Portal ${endpoint} is not responding; skipped after repeated connection failures`; +} + +const FAST_FAIL_MESSAGE_PATTERN = + /is not responding; skipped after repeated connection failures/; + +/** + * Whether a message came from the guard, bare or in the form Electron wraps + * it into (`Error invoking remote method 'STALKER_REQUEST': `). + * + * Exported so tests and any future renderer-side handling share one + * definition of the wording instead of re-typing the sentence. + */ +export function isHostConnectivityFastFailMessage(message: unknown): boolean { + return ( + typeof message === 'string' && FAST_FAIL_MESSAGE_PATTERN.test(message) + ); +} diff --git a/libs/shared/interfaces/src/lib/ipc-commands.ts b/libs/shared/interfaces/src/lib/ipc-commands.ts index 10f882375..36aebbcfd 100644 --- a/libs/shared/interfaces/src/lib/ipc-commands.ts +++ b/libs/shared/interfaces/src/lib/ipc-commands.ts @@ -100,6 +100,14 @@ export const STALKER_REQUEST = 'STALKER_REQUEST'; export const STALKER_RESPONSE = 'STALKER_RESPONSE'; export const PORTAL_DEBUG_EVENT = 'PORTAL_DEBUG_EVENT'; +/** + * Forgets the main process' recorded connection failures for a portal host, so + * the next request contacts it for real. Sent whenever the user asks for a + * fresh attempt (portal retry, "test connection") or hands over a possibly + * different portal (endpoint discovery on import, edit, or lazy repair). + */ +export const CONNECTIVITY_GUARD_RESET = 'CONNECTIVITY_GUARD_RESET'; + // Settings export const SETTINGS_UPDATE = 'SETTINGS_UPDATE'; export const DELETE_ALL_PLAYLISTS = 'DELETE_ALL_PLAYLISTS'; diff --git a/libs/shared/interfaces/src/lib/stalker-auth-failure.util.ts b/libs/shared/interfaces/src/lib/stalker-auth-failure.util.ts index 39927e0db..3b8cc86ff 100644 --- a/libs/shared/interfaces/src/lib/stalker-auth-failure.util.ts +++ b/libs/shared/interfaces/src/lib/stalker-auth-failure.util.ts @@ -21,6 +21,8 @@ * URL builders were centralised here. Two copies would drift. */ +import { isHostConnectivityFastFailMessage } from './host-connectivity.util'; + /** The exact plain-text bodies the stock middleware emits. */ export const STALKER_AUTH_FAILURE_BODIES = [ 'Authorization failed.', @@ -122,9 +124,19 @@ function isStalkerJsonAuthFailurePhrase(value: string): boolean { * arbitrary portal body, where the same breadth would false-positive. */ export function isStalkerAuthFailureMessage(message: unknown): boolean { - return ( - typeof message === 'string' && isStalkerJsonAuthFailurePhrase(message) - ); + if (typeof message !== 'string') { + return false; + } + + // The connectivity guard's refusal names the portal endpoint, and a + // hostname is user data: `https://authorization.example` would otherwise + // match the broad phrase set below and send an unreachable host into lazy + // portal repair. An explicit marker outranks a heuristic. + if (isHostConnectivityFastFailMessage(message)) { + return false; + } + + return isStalkerJsonAuthFailurePhrase(message); } /**