refactor(epg): split the current-programs lookup out of EpgService

Finishes the size finding rather than arguing it: `epg.service.ts` is now
247 lines and has left `max-lines-baseline.mjs` — the baseline the repo says
must only shrink is one file shorter.

The whole "what is on air" machinery moved out unchanged:

- `EpgSingleProgramLookup` — one channel, and the scope ladder it walks.
- `EpgCurrentProgramsLookup` — the batch ladder, backed by the per-channel
  lookups so both shapes answer identically.
- `EpgLookupContext` — the five collaborators they need from the service.
- `epg-lookup-normalization.util.ts` — the two pure id/URL normalizers the
  service's metadata path shares with them.

The service keeps the import pipeline, the capability guards and the two
public entry points that delegate. Note for readers: the cache field must be
declared before the context that captures it — class fields initialise in
declaration order, and a context built too early carries an undefined cache.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
4grayandClaude Opus 5 committed 2026-09-20 20:39:00 +02:00
1 parent 9d879ad0a4
commit 3a7fb9d4be
6 files changed
+616 -511

No files matched your search

@@ -0,0 +1,420 @@
import { Observable, forkJoin, from, of } from 'rxjs';
import {
catchError,
finalize,
map,
shareReplay,
switchMap,
tap,
timeout,
} from 'rxjs/operators';
import { EpgProgram } from '@iptvnator/shared/interfaces';
import { EpgLookupOptions } from './epg-runtime-bridge.service';
import {
normalizeLookupChannelIds,
normalizeLookupSourceUrls,
} from './epg-lookup-normalization.util';
import type { EpgLookupContext } from './epg-lookup-context';
import { EpgSingleProgramLookup } from './epg-single-program.lookup';
/**
* "What is on air right now", for one channel or a batch of them.
*
* The whole source-scope ladder lives here: the caller's playlist sources
* first, then the Settings-managed global ones, and — only when the caller
* opts in with `anySourceFallback` — a final retry across every imported
* guide. Extracted from `EpgService` to keep both files inside the
* repository's production size limit.
*/
export class EpgCurrentProgramsLookup {
constructor(
private readonly ctx: EpgLookupContext,
private readonly single: EpgSingleProgramLookup
) {}
/**
* Gets current programs for multiple channels (batch operation)
* @param channelIds Array of channel IDs
* @returns Observable of Map with channelId -> current program
*/
forChannels(
channelIds: string[],
options?: EpgLookupOptions
): Observable<Map<string, EpgProgram | null>> {
if (!this.ctx.bridge.supportsProgramLookup) {
return of(new Map());
}
if (!channelIds || channelIds.length === 0) {
return of(new Map());
}
const scoped$ = this.getSourceScopedCurrentProgramsForChannels(
channelIds,
options
);
if (!scoped$) {
// Nothing declares a scope — neither the caller nor Settings — so
// the pool of every imported source is the only answer there is.
return this.getUnscopedCurrentProgramsForChannels(channelIds);
}
if (!options?.anySourceFallback) {
return scoped$;
}
return scoped$.pipe(
switchMap((scopedMap) => this.fillFromAnySource(scopedMap))
);
}
/**
* The scoped lookup: the caller's playlist sources first (with the global
* sources as fallback), else the Settings-managed global sources alone.
* `null` when no scope applies and the unscoped pool is the only answer.
*
* The ladder is the same whether or not the bridge has the batch
* endpoint. Collapsing a legacy preload straight into the source-less
* lookup would drop the caller's scope entirely, which is exactly what
* the scopes exist to prevent: two imported guides reusing one XMLTV id
* would answer each other's channels.
*/
private getSourceScopedCurrentProgramsForChannels(
channelIds: string[],
options?: EpgLookupOptions
): Observable<Map<string, EpgProgram | null>> | null {
const sourceUrls = normalizeLookupSourceUrls(options);
if (sourceUrls.length > 0) {
return this.scopedCurrentProgramsForChannels(
channelIds,
sourceUrls,
this.ctx.globalSourceUrls(sourceUrls)
);
}
const globalSourceUrls = this.ctx.globalSourceUrls();
if (globalSourceUrls.length > 0) {
return this.scopedCurrentProgramsForChannels(
channelIds,
globalSourceUrls,
[]
);
}
return null;
}
/** One scoped batch, or its per-channel equivalent on an older preload. */
private scopedCurrentProgramsForChannels(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls: string[]
): Observable<Map<string, EpgProgram | null>> {
if (this.ctx.bridge.supportsCurrentProgramBatch) {
return this.getScopedCurrentProgramsForChannels(
channelIds,
sourceUrls,
fallbackSourceUrls
);
}
const normalizedChannelIds = normalizeLookupChannelIds(channelIds);
if (normalizedChannelIds.length === 0) {
return of(new Map());
}
// `getScopedCurrentProgramForChannel` walks the same scope ->
// fallback-scope ladder the batch query does, and caches under the
// same scoped key.
return forkJoin(
normalizedChannelIds.map((channelId) =>
this.single
.scoped(channelId, sourceUrls, fallbackSourceUrls)
.pipe(
this.ctx.guard(),
timeout(5000),
map((program) => ({ channelId, program })),
catchError(() => of({ channelId, program: null }))
)
)
).pipe(
this.ctx.guard(),
map(
(results) =>
new Map(
results.map((result) => [
result.channelId,
result.program,
])
)
)
);
}
/**
* `anySourceFallback`: keys the scope answered with `null` are retried
* against every imported source. A scoped miss is kept as the answer
* when the pool has nothing either, so the merged map still names every
* requested key.
*/
private fillFromAnySource(
scopedMap: Map<string, EpgProgram | null>
): Observable<Map<string, EpgProgram | null>> {
const unresolvedIds = Array.from(scopedMap.entries())
.filter(([, program]) => !program)
.map(([channelId]) => channelId);
if (unresolvedIds.length === 0) {
return of(scopedMap);
}
return this.getUnscopedCurrentProgramsForChannels(unresolvedIds).pipe(
map((anySourceMap) => {
const mergedMap = new Map(scopedMap);
anySourceMap.forEach((program, channelId) => {
if (program) {
mergedMap.set(channelId, program);
}
});
return mergedMap;
})
);
}
/** Lookup across every imported source, cached under the source-less key. */
private getUnscopedCurrentProgramsForChannels(
channelIds: string[]
): Observable<Map<string, EpgProgram | null>> {
const resultMap = new Map<string, EpgProgram | null>();
const channelsToFetch: string[] = [];
const now = Date.now();
const offsetMinutes = this.ctx.offsetMinutes();
// Check cache for each channel
channelIds.forEach((channelId) => {
const cached = this.ctx.cache.getFresh(channelId, now);
if (cached && cached.offsetMinutes === offsetMinutes) {
resultMap.set(channelId, cached.program);
} else {
channelsToFetch.push(channelId);
}
});
// If all channels were cached, return immediately
if (channelsToFetch.length === 0) {
return of(resultMap);
}
// Single batched IPC + SQL query when the backend supports it.
// Replaces the legacy N+1 forkJoin where each channel fired its own
// GET_CHANNEL_PROGRAMS round-trip.
if (this.ctx.bridge.supportsCurrentProgramBatch) {
return from(
this.ctx.bridge.getCurrentProgramsBatch(channelsToFetch, {
nowMs: this.ctx.clockMs(),
})
).pipe(
this.ctx.guard(),
timeout(5000),
map((batchResult) => {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
const program = batchResult?.[channelId] ?? null;
resultMap.set(channelId, program);
this.ctx.cache.set(
channelId,
program,
offsetMinutes,
cacheTimestamp
);
});
return resultMap;
}),
catchError((err) => {
console.error('EPG batch current programs error:', err);
return of(resultMap);
})
);
}
// Fallback for older preload bundles without the batch endpoint.
// Deliberately the unscoped per-channel lookup: everything reaching
// here wants the source-less pool — either nothing declared a scope,
// or this is the any-source retry after the scoped pass. The public
// `getCurrentProgramForChannel` would put the Settings scope back on
// and make the retry re-ask the question the scoped pass answered.
const fetchObservables = channelsToFetch.map((channelId) =>
this.single.unscoped(channelId).pipe(
this.ctx.guard(),
timeout(5000),
map((program) => ({ channelId, program })),
catchError(() => of({ channelId, program: null }))
)
);
return forkJoin(fetchObservables).pipe(
this.ctx.guard(),
map((results) => {
results.forEach((result) => {
resultMap.set(result.channelId, result.program);
});
return resultMap;
})
);
}
private getScopedCurrentProgramsForChannels(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls = this.ctx.globalSourceUrls(sourceUrls)
): Observable<Map<string, EpgProgram | null>> {
const normalizedChannelIds = normalizeLookupChannelIds(channelIds);
if (normalizedChannelIds.length === 0) {
return of(new Map());
}
const resultMap = new Map<string, EpgProgram | null>();
const channelsToFetch: string[] = [];
normalizedChannelIds.forEach((channelId) => {
const cached = this.ctx.cache.get(
this.ctx.cache.keyFor(channelId, sourceUrls)
);
if (cached) {
resultMap.set(channelId, cached.program);
} else {
channelsToFetch.push(channelId);
}
});
if (channelsToFetch.length === 0) {
return of(resultMap);
}
const batchCacheKey = this.ctx.cache.batchKeyFor(
channelsToFetch,
sourceUrls,
fallbackSourceUrls
);
const existingRequest = this.ctx.cache.batchInFlight(batchCacheKey);
// Tag entries with the offset the verdict was computed with, not the
// one current when the response lands.
const offsetMinutes = this.ctx.offsetMinutes();
const request$ =
existingRequest ??
this.fetchScopedCurrentProgramsBatch(
channelsToFetch,
sourceUrls,
fallbackSourceUrls
).pipe(
this.ctx.guard(),
tap((fetchedMap) => {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
this.ctx.cache.set(
this.ctx.cache.keyFor(channelId, sourceUrls),
fetchedMap.get(channelId) ?? null,
offsetMinutes,
cacheTimestamp
);
});
}),
finalize(() =>
this.ctx.cache.releaseBatch(batchCacheKey, request$)
),
shareReplay({ bufferSize: 1, refCount: false })
);
if (!existingRequest) {
this.ctx.cache.registerBatch(batchCacheKey, request$);
}
return request$.pipe(
this.ctx.guard(),
map((fetchedMap) => {
const mergedResultMap = new Map(resultMap);
channelsToFetch.forEach((channelId) => {
mergedResultMap.set(
channelId,
fetchedMap.get(channelId) ?? null
);
});
return mergedResultMap;
})
);
}
private fetchScopedCurrentProgramsBatch(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls: string[]
): Observable<Map<string, EpgProgram | null>> {
const nowMs = this.ctx.clockMs();
return from(
this.ctx.bridge.getCurrentProgramsBatch(channelIds, {
sourceUrls,
nowMs,
})
).pipe(
this.ctx.guard(),
timeout(5000),
switchMap((scopedResult) => {
const resultMap = new Map<string, EpgProgram | null>();
const fallbackChannelIds: string[] = [];
channelIds.forEach((channelId) => {
const program = scopedResult?.[channelId] ?? null;
resultMap.set(channelId, program);
if (!program) {
fallbackChannelIds.push(channelId);
}
});
if (fallbackChannelIds.length === 0) {
return of(resultMap);
}
if (fallbackSourceUrls.length === 0) {
return of(resultMap);
}
return from(
this.ctx.bridge.getCurrentProgramsBatch(
fallbackChannelIds,
{
sourceUrls: fallbackSourceUrls,
nowMs,
}
)
).pipe(
this.ctx.guard(),
timeout(5000),
map((globalResult) => {
fallbackChannelIds.forEach((channelId) => {
resultMap.set(
channelId,
globalResult?.[channelId] ?? null
);
});
return resultMap;
}),
catchError((err) => {
console.error(
'EPG global fallback current programs error:',
err
);
return of(resultMap);
})
);
}),
catchError((err) => {
console.error('EPG scoped batch current programs error:', err);
return of(this.createNullProgramMap(channelIds));
})
);
}
private createNullProgramMap(
channelIds: string[]
): Map<string, EpgProgram | null> {
return new Map(channelIds.map((channelId) => [channelId, null]));
}
}
@@ -0,0 +1,17 @@
import { MonoTypeOperatorFunction } from 'rxjs';
import { EpgProgramCache } from './epg-program-cache';
import { EpgRuntimeBridgeService } from './epg-runtime-bridge.service';
/** Everything the lookups need from `EpgService`, and nothing else. */
export interface EpgLookupContext {
readonly bridge: EpgRuntimeBridgeService;
readonly cache: EpgProgramCache;
/** Retires work that a changed XMLTV source set has invalidated. */
guard<T>(): MonoTypeOperatorFunction<T>;
/** Current EPG display offset, in minutes. */
offsetMinutes(): number;
/** Wall-clock now, expressed in the provider's uncorrected EPG clock. */
clockMs(): number;
/** Settings-managed XMLTV URLs, optionally minus a caller's own scope. */
globalSourceUrls(excluding?: string[]): string[];
}
@@ -0,0 +1,20 @@
import { normalizeEpgUrls } from '@iptvnator/shared/m3u-utils';
import { EpgLookupOptions } from './epg-runtime-bridge.service';
/** Trimmed, de-duplicated, non-empty channel ids, in request order. */
export function normalizeLookupChannelIds(channelIds: string[]): string[] {
return Array.from(
new Set(
channelIds
.map((channelId) => channelId.trim())
.filter((channelId) => channelId.length > 0)
)
);
}
/** The caller's declared XMLTV scope, normalized like every stored URL. */
export function normalizeLookupSourceUrls(
options?: EpgLookupOptions
): string[] {
return normalizeEpgUrls(options?.sourceUrls ?? []);
}
@@ -0,0 +1,129 @@
import { Observable, from, of } from 'rxjs';
import { catchError, map, switchMap, timeout } from 'rxjs/operators';
import { EpgProgram } from '@iptvnator/shared/interfaces';
import { EpgLookupOptions } from './epg-runtime-bridge.service';
import { normalizeEpgPrograms } from './epg-program-normalization.util';
import { normalizeLookupSourceUrls } from './epg-lookup-normalization.util';
import type { EpgLookupContext } from './epg-lookup-context';
/**
* "What is on air" for a single channel, and the scope ladder one channel
* walks: the caller's sources, then the Settings-managed ones, then — for a
* caller that asked for it — the whole imported pool. The batch lookup backs
* its own per-channel work with these, so both shapes answer identically.
*/
export class EpgSingleProgramLookup {
constructor(private readonly ctx: EpgLookupContext) {}
/**
* Gets the current EPG program for a specific channel (with caching)
* @param channelId Channel ID (tvg-id or channel name)
* @returns Observable of current program or null
*/
forChannel(
channelId: string,
options?: EpgLookupOptions
): Observable<EpgProgram | null> {
if (!this.ctx.bridge.supportsProgramLookup || !channelId) {
return of(null);
}
const sourceUrls = normalizeLookupSourceUrls(options);
if (sourceUrls.length > 0) {
return this.scoped(
channelId,
sourceUrls,
this.ctx.globalSourceUrls(sourceUrls)
);
}
const globalSourceUrls = this.ctx.globalSourceUrls();
if (globalSourceUrls.length > 0) {
return this.scoped(channelId, globalSourceUrls, []);
}
return this.unscoped(channelId);
}
/**
* The lookup across every imported source, cached under the source-less
* key. Kept separate from `getCurrentProgramForChannel`, which re-applies
* the Settings-managed scope whenever global URLs exist: a caller that
* has already decided it wants the unscoped pool — the `anySourceFallback`
* retry on a preload without the batch endpoint — must not have that
* scope put back on.
*/
unscoped(channelId: string): Observable<EpgProgram | null> {
// Check cache first
const cacheKey = this.ctx.cache.keyFor(channelId);
// Fetch from backend
return this.ctx.cache.getOrFetch(cacheKey, () =>
from(this.ctx.bridge.getChannelPrograms(channelId)).pipe(
this.ctx.guard(),
map((programs) => normalizeEpgPrograms(programs ?? [])),
map((programs: EpgProgram[]) => this.findCurrent(programs)),
catchError((err) => {
console.error('EPG get current program error:', err);
return of(null);
})
)
);
}
/**
* Finds the current program from a list of programs
*/
findCurrent(programs: EpgProgram[]): EpgProgram | null {
const now = this.ctx.clockMs();
return (
programs.find((program) => {
const start = new Date(program.start).getTime();
const stop = new Date(program.stop).getTime();
return start <= now && now <= stop;
}) || null
);
}
scoped(
channelId: string,
sourceUrls: string[],
fallbackSourceUrls: string[]
): Observable<EpgProgram | null> {
const cacheKey = this.ctx.cache.keyFor(channelId, sourceUrls);
return this.ctx.cache.getOrFetch(cacheKey, () =>
from(
this.ctx.bridge.getChannelPrograms(channelId, { sourceUrls })
).pipe(
this.ctx.guard(),
timeout(3000),
map((programs) => normalizeEpgPrograms(programs ?? [])),
switchMap((programs) => {
const currentProgram = this.findCurrent(programs);
if (currentProgram) {
return of(currentProgram);
}
return this.fallbackFor(channelId, fallbackSourceUrls);
}),
catchError((err) => {
console.error('EPG scoped current program error:', err);
return this.fallbackFor(channelId, fallbackSourceUrls);
})
)
);
}
private fallbackFor(
channelId: string,
sourceUrls: string[]
): Observable<EpgProgram | null> {
if (sourceUrls.length === 0) {
return of(null);
}
return this.scoped(channelId, sourceUrls, []);
}
}
+30 -510
View File
@@ -1,16 +1,8 @@
import { inject, Injectable } from '@angular/core';
import { MatSnackBar } from '@angular/material/snack-bar';
import { TranslateService } from '@ngx-translate/core';
import { BehaviorSubject, forkJoin, from, Observable, of } from 'rxjs';
import {
catchError,
finalize,
map,
shareReplay,
switchMap,
tap,
timeout,
} from 'rxjs/operators';
import { BehaviorSubject, from, Observable, of } from 'rxjs';
import { catchError, map, tap, timeout } from 'rxjs/operators';
import {
createDevLogger,
EpgChannelMetadata,
@@ -25,6 +17,13 @@ import {
import { normalizeEpgPrograms } from './epg-program-normalization.util';
import { EpgProgramCache } from './epg-program-cache';
import { EpgChannelMetadataLookup } from './epg-channel-metadata.lookup';
import { EpgCurrentProgramsLookup } from './epg-current-programs.lookup';
import { EpgSingleProgramLookup } from './epg-single-program.lookup';
import type { EpgLookupContext } from './epg-lookup-context';
import {
normalizeLookupChannelIds,
normalizeLookupSourceUrls,
} from './epg-lookup-normalization.util';
import { normalizeEpgUrls } from '@iptvnator/shared/m3u-utils';
const debugEpgService = createDevLogger('EpgService');
@@ -63,6 +62,23 @@ export class EpgService {
() => this.sourceSettings.guard()
);
/** The whole "what is on air" scope ladder, single channel or batch. */
private readonly lookupContext: EpgLookupContext = {
bridge: this.epgBridge,
cache: this.programCache,
guard: () => this.sourceSettings.guard(),
offsetMinutes: () => this.epgOffsetMinutes(),
clockMs: () => this.epgClockMs(),
globalSourceUrls: (excluding) => this.getGlobalEpgSourceUrls(excluding),
};
private readonly singleProgram = new EpgSingleProgramLookup(
this.lookupContext
);
private readonly currentPrograms = new EpgCurrentProgramsLookup(
this.lookupContext,
this.singleProgram
);
/** Display offset every "currently airing" decision in here is made with. */
private epgOffsetMinutes(): number {
return this.settingsStore.resolvedEpgOffsetMinutes();
@@ -175,71 +191,7 @@ export class EpgService {
if (!this.epgBridge.supportsProgramLookup || !channelId) {
return of(null);
}
const sourceUrls = this.normalizeSourceUrls(options);
if (sourceUrls.length > 0) {
return this.getScopedCurrentProgramForChannel(
channelId,
sourceUrls,
this.getGlobalEpgSourceUrls(sourceUrls)
);
}
const globalSourceUrls = this.getGlobalEpgSourceUrls();
if (globalSourceUrls.length > 0) {
return this.getScopedCurrentProgramForChannel(
channelId,
globalSourceUrls,
[]
);
}
return this.getUnscopedCurrentProgramForChannel(channelId);
}
/**
* The lookup across every imported source, cached under the source-less
* key. Kept separate from `getCurrentProgramForChannel`, which re-applies
* the Settings-managed scope whenever global URLs exist: a caller that
* has already decided it wants the unscoped pool — the `anySourceFallback`
* retry on a preload without the batch endpoint — must not have that
* scope put back on.
*/
private getUnscopedCurrentProgramForChannel(
channelId: string
): Observable<EpgProgram | null> {
// Check cache first
const cacheKey = this.programCache.keyFor(channelId);
// Fetch from backend
return this.programCache.getOrFetch(cacheKey, () =>
from(this.epgBridge.getChannelPrograms(channelId)).pipe(
this.sourceSettings.guard(),
map((programs) => normalizeEpgPrograms(programs ?? [])),
map((programs: EpgProgram[]) =>
this.findCurrentProgram(programs)
),
catchError((err) => {
console.error('EPG get current program error:', err);
return of(null);
})
)
);
}
/**
* Finds the current program from a list of programs
*/
private findCurrentProgram(programs: EpgProgram[]): EpgProgram | null {
const now = this.epgClockMs();
return (
programs.find((program) => {
const start = new Date(program.start).getTime();
const stop = new Date(program.stop).getTime();
return start <= now && now <= stop;
}) || null
);
return this.singleProgram.forChannel(channelId, options);
}
/**
@@ -254,371 +206,10 @@ export class EpgService {
if (!this.epgBridge.supportsProgramLookup) {
return of(new Map());
}
if (!channelIds || channelIds.length === 0) {
return of(new Map());
}
const scoped$ = this.getSourceScopedCurrentProgramsForChannels(
channelIds,
options
);
if (!scoped$) {
// Nothing declares a scope — neither the caller nor Settings — so
// the pool of every imported source is the only answer there is.
return this.getUnscopedCurrentProgramsForChannels(channelIds);
}
if (!options?.anySourceFallback) {
return scoped$;
}
return scoped$.pipe(
switchMap((scopedMap) => this.fillFromAnySource(scopedMap))
);
}
/**
* The scoped lookup: the caller's playlist sources first (with the global
* sources as fallback), else the Settings-managed global sources alone.
* `null` when no scope applies and the unscoped pool is the only answer.
*
* The ladder is the same whether or not the bridge has the batch
* endpoint. Collapsing a legacy preload straight into the source-less
* lookup would drop the caller's scope entirely, which is exactly what
* the scopes exist to prevent: two imported guides reusing one XMLTV id
* would answer each other's channels.
*/
private getSourceScopedCurrentProgramsForChannels(
channelIds: string[],
options?: EpgLookupOptions
): Observable<Map<string, EpgProgram | null>> | null {
const sourceUrls = this.normalizeSourceUrls(options);
if (sourceUrls.length > 0) {
return this.scopedCurrentProgramsForChannels(
channelIds,
sourceUrls,
this.getGlobalEpgSourceUrls(sourceUrls)
);
}
const globalSourceUrls = this.getGlobalEpgSourceUrls();
if (globalSourceUrls.length > 0) {
return this.scopedCurrentProgramsForChannels(
channelIds,
globalSourceUrls,
[]
);
}
return null;
}
/** One scoped batch, or its per-channel equivalent on an older preload. */
private scopedCurrentProgramsForChannels(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls: string[]
): Observable<Map<string, EpgProgram | null>> {
if (this.epgBridge.supportsCurrentProgramBatch) {
return this.getScopedCurrentProgramsForChannels(
channelIds,
sourceUrls,
fallbackSourceUrls
);
}
const normalizedChannelIds = this.normalizeChannelIds(channelIds);
if (normalizedChannelIds.length === 0) {
return of(new Map());
}
// `getScopedCurrentProgramForChannel` walks the same scope ->
// fallback-scope ladder the batch query does, and caches under the
// same scoped key.
return forkJoin(
normalizedChannelIds.map((channelId) =>
this.getScopedCurrentProgramForChannel(
channelId,
sourceUrls,
fallbackSourceUrls
).pipe(
this.sourceSettings.guard(),
timeout(5000),
map((program) => ({ channelId, program })),
catchError(() => of({ channelId, program: null }))
)
)
).pipe(
this.sourceSettings.guard(),
map(
(results) =>
new Map(
results.map((result) => [
result.channelId,
result.program,
])
)
)
);
}
/**
* `anySourceFallback`: keys the scope answered with `null` are retried
* against every imported source. A scoped miss is kept as the answer
* when the pool has nothing either, so the merged map still names every
* requested key.
*/
private fillFromAnySource(
scopedMap: Map<string, EpgProgram | null>
): Observable<Map<string, EpgProgram | null>> {
const unresolvedIds = Array.from(scopedMap.entries())
.filter(([, program]) => !program)
.map(([channelId]) => channelId);
if (unresolvedIds.length === 0) {
return of(scopedMap);
}
return this.getUnscopedCurrentProgramsForChannels(unresolvedIds).pipe(
map((anySourceMap) => {
const mergedMap = new Map(scopedMap);
anySourceMap.forEach((program, channelId) => {
if (program) {
mergedMap.set(channelId, program);
}
});
return mergedMap;
})
);
}
/** Lookup across every imported source, cached under the source-less key. */
private getUnscopedCurrentProgramsForChannels(
channelIds: string[]
): Observable<Map<string, EpgProgram | null>> {
const resultMap = new Map<string, EpgProgram | null>();
const channelsToFetch: string[] = [];
const now = Date.now();
const offsetMinutes = this.epgOffsetMinutes();
// Check cache for each channel
channelIds.forEach((channelId) => {
const cached = this.programCache.getFresh(channelId, now);
if (cached && cached.offsetMinutes === offsetMinutes) {
resultMap.set(channelId, cached.program);
} else {
channelsToFetch.push(channelId);
}
});
// If all channels were cached, return immediately
if (channelsToFetch.length === 0) {
return of(resultMap);
}
// Single batched IPC + SQL query when the backend supports it.
// Replaces the legacy N+1 forkJoin where each channel fired its own
// GET_CHANNEL_PROGRAMS round-trip.
if (this.epgBridge.supportsCurrentProgramBatch) {
return from(
this.epgBridge.getCurrentProgramsBatch(channelsToFetch, {
nowMs: this.epgClockMs(),
})
).pipe(
this.sourceSettings.guard(),
timeout(5000),
map((batchResult) => {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
const program = batchResult?.[channelId] ?? null;
resultMap.set(channelId, program);
this.programCache.set(
channelId,
program,
offsetMinutes,
cacheTimestamp
);
});
return resultMap;
}),
catchError((err) => {
console.error('EPG batch current programs error:', err);
return of(resultMap);
})
);
}
// Fallback for older preload bundles without the batch endpoint.
// Deliberately the unscoped per-channel lookup: everything reaching
// here wants the source-less pool — either nothing declared a scope,
// or this is the any-source retry after the scoped pass. The public
// `getCurrentProgramForChannel` would put the Settings scope back on
// and make the retry re-ask the question the scoped pass answered.
const fetchObservables = channelsToFetch.map((channelId) =>
this.getUnscopedCurrentProgramForChannel(channelId).pipe(
this.sourceSettings.guard(),
timeout(5000),
map((program) => ({ channelId, program })),
catchError(() => of({ channelId, program: null }))
)
);
return forkJoin(fetchObservables).pipe(
this.sourceSettings.guard(),
map((results) => {
results.forEach((result) => {
resultMap.set(result.channelId, result.program);
});
return resultMap;
})
);
}
private getScopedCurrentProgramsForChannels(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls = this.getGlobalEpgSourceUrls(sourceUrls)
): Observable<Map<string, EpgProgram | null>> {
const normalizedChannelIds = this.normalizeChannelIds(channelIds);
if (normalizedChannelIds.length === 0) {
return of(new Map());
}
const resultMap = new Map<string, EpgProgram | null>();
const channelsToFetch: string[] = [];
normalizedChannelIds.forEach((channelId) => {
const cached = this.programCache.get(
this.programCache.keyFor(channelId, sourceUrls)
);
if (cached) {
resultMap.set(channelId, cached.program);
} else {
channelsToFetch.push(channelId);
}
});
if (channelsToFetch.length === 0) {
return of(resultMap);
}
const batchCacheKey = this.programCache.batchKeyFor(
channelsToFetch,
sourceUrls,
fallbackSourceUrls
);
const existingRequest = this.programCache.batchInFlight(batchCacheKey);
// Tag entries with the offset the verdict was computed with, not the
// one current when the response lands.
const offsetMinutes = this.epgOffsetMinutes();
const request$ =
existingRequest ??
this.fetchScopedCurrentProgramsBatch(
channelsToFetch,
sourceUrls,
fallbackSourceUrls
).pipe(
this.sourceSettings.guard(),
tap((fetchedMap) => {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
this.programCache.set(
this.programCache.keyFor(channelId, sourceUrls),
fetchedMap.get(channelId) ?? null,
offsetMinutes,
cacheTimestamp
);
});
}),
finalize(() =>
this.programCache.releaseBatch(batchCacheKey, request$)
),
shareReplay({ bufferSize: 1, refCount: false })
);
if (!existingRequest) {
this.programCache.registerBatch(batchCacheKey, request$);
}
return request$.pipe(
this.sourceSettings.guard(),
map((fetchedMap) => {
const mergedResultMap = new Map(resultMap);
channelsToFetch.forEach((channelId) => {
mergedResultMap.set(
channelId,
fetchedMap.get(channelId) ?? null
);
});
return mergedResultMap;
})
);
}
private fetchScopedCurrentProgramsBatch(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls: string[]
): Observable<Map<string, EpgProgram | null>> {
const nowMs = this.epgClockMs();
return from(
this.epgBridge.getCurrentProgramsBatch(channelIds, {
sourceUrls,
nowMs,
})
).pipe(
this.sourceSettings.guard(),
timeout(5000),
switchMap((scopedResult) => {
const resultMap = new Map<string, EpgProgram | null>();
const fallbackChannelIds: string[] = [];
channelIds.forEach((channelId) => {
const program = scopedResult?.[channelId] ?? null;
resultMap.set(channelId, program);
if (!program) {
fallbackChannelIds.push(channelId);
}
});
if (fallbackChannelIds.length === 0) {
return of(resultMap);
}
if (fallbackSourceUrls.length === 0) {
return of(resultMap);
}
return from(
this.epgBridge.getCurrentProgramsBatch(fallbackChannelIds, {
sourceUrls: fallbackSourceUrls,
nowMs,
})
).pipe(
this.sourceSettings.guard(),
timeout(5000),
map((globalResult) => {
fallbackChannelIds.forEach((channelId) => {
resultMap.set(
channelId,
globalResult?.[channelId] ?? null
);
});
return resultMap;
}),
catchError((err) => {
console.error(
'EPG global fallback current programs error:',
err
);
return of(resultMap);
})
);
}),
catchError((err) => {
console.error('EPG scoped batch current programs error:', err);
return of(this.createNullProgramMap(channelIds));
})
);
return this.currentPrograms.forChannels(channelIds, options);
}
getChannelMetadataForChannels(
@@ -629,88 +220,17 @@ export class EpgService {
return of(new Map());
}
const normalizedChannelIds = this.normalizeChannelIds(channelIds);
const normalizedChannelIds = normalizeLookupChannelIds(channelIds);
if (normalizedChannelIds.length === 0) {
return of(new Map());
}
return this.channelMetadata.forChannels(
normalizedChannelIds,
this.normalizeSourceUrls(options)
normalizeLookupSourceUrls(options)
);
}
private normalizeChannelIds(channelIds: string[]): string[] {
return Array.from(
new Set(
channelIds
.map((channelId) => channelId.trim())
.filter((channelId) => channelId.length > 0)
)
);
}
private normalizeSourceUrls(options?: EpgLookupOptions): string[] {
return normalizeEpgUrls(options?.sourceUrls ?? []);
}
private getScopedCurrentProgramForChannel(
channelId: string,
sourceUrls: string[],
fallbackSourceUrls: string[]
): Observable<EpgProgram | null> {
const cacheKey = this.programCache.keyFor(channelId, sourceUrls);
return this.programCache.getOrFetch(cacheKey, () =>
from(
this.epgBridge.getChannelPrograms(channelId, { sourceUrls })
).pipe(
this.sourceSettings.guard(),
timeout(3000),
map((programs) => normalizeEpgPrograms(programs ?? [])),
switchMap((programs) => {
const currentProgram = this.findCurrentProgram(programs);
if (currentProgram) {
return of(currentProgram);
}
return this.getFallbackCurrentProgramForChannel(
channelId,
fallbackSourceUrls
);
}),
catchError((err) => {
console.error('EPG scoped current program error:', err);
return this.getFallbackCurrentProgramForChannel(
channelId,
fallbackSourceUrls
);
})
)
);
}
private getFallbackCurrentProgramForChannel(
channelId: string,
sourceUrls: string[]
): Observable<EpgProgram | null> {
if (sourceUrls.length === 0) {
return of(null);
}
return this.getScopedCurrentProgramForChannel(
channelId,
sourceUrls,
[]
);
}
private createNullProgramMap(
channelIds: string[]
): Map<string, EpgProgram | null> {
return new Map(channelIds.map((channelId) => [channelId, null]));
}
private getGlobalEpgSourceUrls(excluding: string[] = []): string[] {
const excludedUrls = new Set(excluding);
return normalizeEpgUrls(
-1
View File
@@ -23,7 +23,6 @@ export const maxLinesBaseline = [
'apps/web/src/app/services/electron.service.ts',
'apps/web/src/app/services/pwa.service.ts',
'apps/xtream-mock-server/src/app/generators/marketing.generator.ts',
'libs/epg/data-access/src/lib/epg.service.ts',
'libs/m3u-state/src/lib/effects.ts',
'libs/playlist/m3u/feature-player/src/lib/video-player/video-player.component.ts',
'libs/playlist/shared/ui/src/lib/playlist-switcher/playlist-switcher.component.ts',