From c5282e328d30e39cf703d815050e6fffcaaadde1 Mon Sep 17 00:00:00 2001 From: 4gray Date: Sun, 20 Sep 2026 20:44:57 +0200 Subject: [PATCH] refactor(epg): split the scoped batch request out of the current-programs lookup MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Greptile counts raw lines where the ESLint rule skips blanks and comments, so the extracted lookup still read as 420 lines against the repo's 400-line directive. Split along the real seam instead of shaving comments: - `EpgScopedBatchLookup` — how ONE scope is asked: the per-channel cache in front of the request, the order-insensitive deduplication of identical batches in flight, and the retry against the fallback scope. - `EpgCurrentProgramsLookup` — which scope to ask, and the any-source retry. Every file in the lib is now comfortably inside the limit: service 249, current-programs 256, scoped-batch 180, single-program 129, cache 166, channel-metadata 101. Co-Authored-By: Claude Opus 5 --- .../src/lib/epg-current-programs.lookup.ts | 178 +---------------- .../src/lib/epg-scoped-batch.lookup.ts | 180 ++++++++++++++++++ libs/epg/data-access/src/lib/epg.service.ts | 4 +- 3 files changed, 190 insertions(+), 172 deletions(-) create mode 100644 libs/epg/data-access/src/lib/epg-scoped-batch.lookup.ts diff --git a/libs/epg/data-access/src/lib/epg-current-programs.lookup.ts b/libs/epg/data-access/src/lib/epg-current-programs.lookup.ts index 55e5a4064..38c296596 100644 --- a/libs/epg/data-access/src/lib/epg-current-programs.lookup.ts +++ b/libs/epg/data-access/src/lib/epg-current-programs.lookup.ts @@ -1,13 +1,5 @@ import { Observable, forkJoin, from, of } from 'rxjs'; -import { - catchError, - finalize, - map, - shareReplay, - switchMap, - tap, - timeout, -} from 'rxjs/operators'; +import { catchError, map, switchMap, timeout } from 'rxjs/operators'; import { EpgProgram } from '@iptvnator/shared/interfaces'; import { EpgLookupOptions } from './epg-runtime-bridge.service'; import { @@ -16,6 +8,7 @@ import { } from './epg-lookup-normalization.util'; import type { EpgLookupContext } from './epg-lookup-context'; import { EpgSingleProgramLookup } from './epg-single-program.lookup'; +import { EpgScopedBatchLookup } from './epg-scoped-batch.lookup'; /** * "What is on air right now" for a batch of channels. @@ -29,7 +22,8 @@ import { EpgSingleProgramLookup } from './epg-single-program.lookup'; export class EpgCurrentProgramsLookup { constructor( private readonly ctx: EpgLookupContext, - private readonly single: EpgSingleProgramLookup + private readonly single: EpgSingleProgramLookup, + private readonly scopedBatch: EpgScopedBatchLookup ) {} /** @@ -110,7 +104,7 @@ export class EpgCurrentProgramsLookup { fallbackSourceUrls: string[] ): Observable> { if (this.ctx.bridge.supportsCurrentProgramBatch) { - return this.getScopedCurrentProgramsForChannels( + return this.scopedBatch.forChannels( channelIds, sourceUrls, fallbackSourceUrls @@ -122,9 +116,8 @@ export class EpgCurrentProgramsLookup { return of(new Map()); } - // `getScopedCurrentProgramForChannel` walks the same scope -> - // fallback-scope ladder the batch query does, and caches under the - // same scoped key. + // The per-channel lookup 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 @@ -260,161 +253,4 @@ export class EpgCurrentProgramsLookup { }) ); } - - private getScopedCurrentProgramsForChannels( - channelIds: string[], - sourceUrls: string[], - fallbackSourceUrls = this.ctx.globalSourceUrls(sourceUrls) - ): Observable> { - const normalizedChannelIds = normalizeLookupChannelIds(channelIds); - if (normalizedChannelIds.length === 0) { - return of(new Map()); - } - - const resultMap = new Map(); - 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> { - 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(); - 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 { - return new Map(channelIds.map((channelId) => [channelId, null])); - } } diff --git a/libs/epg/data-access/src/lib/epg-scoped-batch.lookup.ts b/libs/epg/data-access/src/lib/epg-scoped-batch.lookup.ts new file mode 100644 index 000000000..9080f9036 --- /dev/null +++ b/libs/epg/data-access/src/lib/epg-scoped-batch.lookup.ts @@ -0,0 +1,180 @@ +import { Observable, from, of } from 'rxjs'; +import { + catchError, + finalize, + map, + shareReplay, + switchMap, + tap, + timeout, +} from 'rxjs/operators'; +import { EpgProgram } from '@iptvnator/shared/interfaces'; +import { normalizeLookupChannelIds } from './epg-lookup-normalization.util'; +import type { EpgLookupContext } from './epg-lookup-context'; + +/** + * How ONE scope is asked: the per-channel cache in front of the request, the + * order-insensitive deduplication of identical batches still in flight, and + * the retry of whatever that scope could not answer against the fallback + * scope. `EpgCurrentProgramsLookup` decides WHICH scope to ask. + */ +export class EpgScopedBatchLookup { + constructor(private readonly ctx: EpgLookupContext) {} + + forChannels( + channelIds: string[], + sourceUrls: string[], + fallbackSourceUrls = this.ctx.globalSourceUrls(sourceUrls) + ): Observable> { + const normalizedChannelIds = normalizeLookupChannelIds(channelIds); + if (normalizedChannelIds.length === 0) { + return of(new Map()); + } + + const resultMap = new Map(); + 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.fetchBatch( + 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 fetchBatch( + channelIds: string[], + sourceUrls: string[], + fallbackSourceUrls: string[] + ): Observable> { + 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(); + 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.nullProgramMap(channelIds)); + }) + ); + } + + private nullProgramMap( + channelIds: string[] + ): Map { + return new Map(channelIds.map((channelId) => [channelId, null])); + } +} diff --git a/libs/epg/data-access/src/lib/epg.service.ts b/libs/epg/data-access/src/lib/epg.service.ts index 0458d7f57..525e308da 100644 --- a/libs/epg/data-access/src/lib/epg.service.ts +++ b/libs/epg/data-access/src/lib/epg.service.ts @@ -19,6 +19,7 @@ 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 { EpgScopedBatchLookup } from './epg-scoped-batch.lookup'; import type { EpgLookupContext } from './epg-lookup-context'; import { normalizeLookupChannelIds, @@ -76,7 +77,8 @@ export class EpgService { ); private readonly currentPrograms = new EpgCurrentProgramsLookup( this.lookupContext, - this.singleProgram + this.singleProgram, + new EpgScopedBatchLookup(this.lookupContext) ); /** Display offset every "currently airing" decision in here is made with. */