diff --git a/apps/master-api/src/app.ts b/apps/master-api/src/app.ts index 3ac2c22..d5806f1 100644 --- a/apps/master-api/src/app.ts +++ b/apps/master-api/src/app.ts @@ -18,6 +18,7 @@ import { startGooglePlayReconciliationPoller } from "./google-play/service"; import { nodeConnectionManager } from "./websocket/manager"; import { MetadataService } from "./metadata/service"; import { toErrorResponse } from "./security/errors"; +import { startIptvEpgWarmup } from "./viewer/iptv-client"; async function main() { const config = loadConfig(); @@ -81,6 +82,7 @@ async function main() { const googlePlay = registerGooglePlayRoutes(app, config); registerSettingsRoutes(app, config, googlePlay); startGooglePlayReconciliationPoller(googlePlay, config.GOOGLE_PLAY_SYNC_INTERVAL_MS); + startIptvEpgWarmup(config); app.get("/api/v1/node/connect", { websocket: true }, (socket) => { void nodeConnectionManager.handleConnection(socket); diff --git a/apps/master-api/src/viewer/iptv-client.ts b/apps/master-api/src/viewer/iptv-client.ts index eaf0123..8aef9d0 100644 --- a/apps/master-api/src/viewer/iptv-client.ts +++ b/apps/master-api/src/viewer/iptv-client.ts @@ -105,6 +105,8 @@ function sortChannels(streams: XtreamLiveStream[]): XtreamLiveStream[] { return [...streams].sort((a, b) => (a.num ?? a.stream_id) - (b.num ?? b.stream_id)); } +type EpgTitlePair = { now: string | null; next: string | null }; + export async function getViewerIptvChannels( viewerId: string, config: Config, @@ -123,19 +125,25 @@ export async function getViewerIptvChannels( : await xt.client.getLiveStreams(categoryId); streams = sortChannels(streams); - const nowTitles = includeEpg - ? await getCachedCategoryNowTitles(tokenId, categoryId, streams, xt.client) - : null; + // Bij includeEpg=false toch warm cache meeleveren (instant Nu-titels). + // Bij includeEpg=true: SWR — direct cache, achtergrond-refresh; bij cold miss wachten. + const epgTitles = includeEpg + ? await getCategoryEpgTitles(tokenId, categoryId, streams, xt.client, { waitIfCold: true }) + : peekCategoryEpgTitles(tokenId, categoryId); return { - channels: streams.map((s) => ({ - streamId: s.stream_id, - name: s.name.trim(), - logoUrl: s.stream_icon || null, - number: s.num ?? null, - categoryId: s.category_id ? String(s.category_id) : null, - nowTitle: nowTitles?.get(s.stream_id) ?? null, - })), + channels: streams.map((s) => { + const pair = epgTitles?.get(s.stream_id); + return { + streamId: s.stream_id, + name: s.name.trim(), + logoUrl: s.stream_icon || null, + number: s.num ?? null, + categoryId: s.category_id ? String(s.category_id) : null, + nowTitle: pair?.now ?? null, + nextTitle: pair?.next ?? null, + }; + }), }; } @@ -195,8 +203,20 @@ export async function getViewerIptvEpg(viewerId: string, config: Config, streamI const ch = streams.find((s) => s.stream_id === streamId); if (!ch) throw new AppError("NOT_FOUND", "Zender niet gevonden", 404); + const cached = getCachedStreamEpg(tokenId, streamId); + if (cached) { + return { + streamId, + channelName: ch.name, + logoUrl: ch.stream_icon || null, + now: cached.now, + next: cached.next, + }; + } + const epg = await xt.client.getShortEpg(streamId, 8); const { now: nowProgram, next: nextProgram } = parseNowNextPrograms(epg); + setCachedStreamEpg(tokenId, streamId, nowProgram, nextProgram); return { streamId, @@ -207,13 +227,46 @@ export async function getViewerIptvEpg(viewerId: string, config: Config, streamI }; } -const EPG_NOW_CACHE_TTL_MS = 90_000; -const epgNowCache = new Map }>(); +/** Fris tot 5 min; stale tot 20 min (SWR). */ +const EPG_FRESH_MS = 5 * 60_000; +const EPG_STALE_MS = 20 * 60_000; +const EPG_DETAIL_TTL_MS = 15 * 60_000; + +type CategoryEpgCacheEntry = { at: number; titles: Map }; +const epgCategoryCache = new Map(); +const epgCategoryRefreshing = new Set(); + +type StreamEpgCacheEntry = { + at: number; + now: { title: string; start: string; end: string | null } | null; + next: { title: string; start: string; end: string | null } | null; +}; +const epgStreamCache = new Map(); function epgCacheKey(tokenId: string, categoryId?: string): string { return `${tokenId}:${categoryId?.trim() || "__all__"}`; } +function streamEpgKey(tokenId: string, streamId: number): string { + return `${tokenId}:s:${streamId}`; +} + +function getCachedStreamEpg(tokenId: string, streamId: number): StreamEpgCacheEntry | null { + const hit = epgStreamCache.get(streamEpgKey(tokenId, streamId)); + if (!hit) return null; + if (Date.now() - hit.at > EPG_DETAIL_TTL_MS) return null; + return hit; +} + +function setCachedStreamEpg( + tokenId: string, + streamId: number, + now: StreamEpgCacheEntry["now"], + next: StreamEpgCacheEntry["next"] +) { + epgStreamCache.set(streamEpgKey(tokenId, streamId), { at: Date.now(), now, next }); +} + export function parseNowNextPrograms(epg: XtreamEpgListing[]): { now: { title: string; start: string; end: string | null } | null; next: { title: string; start: string; end: string | null } | null; @@ -240,33 +293,83 @@ export function parseNowNextPrograms(epg: XtreamEpgListing[]): { return { now: nowProgram, next: nextProgram }; } -async function getCachedCategoryNowTitles( +function peekCategoryEpgTitles( + tokenId: string, + categoryId: string | undefined +): Map | null { + const key = epgCacheKey(tokenId, categoryId); + const cached = epgCategoryCache.get(key); + if (!cached) return null; + if (Date.now() - cached.at > EPG_STALE_MS) return null; + return cached.titles; +} + +async function getCategoryEpgTitles( + tokenId: string, + categoryId: string | undefined, + streams: XtreamLiveStream[], + client: XtreamClient, + opts: { waitIfCold: boolean } +): Promise> { + const key = epgCacheKey(tokenId, categoryId); + const cached = epgCategoryCache.get(key); + const age = cached ? Date.now() - cached.at : Number.POSITIVE_INFINITY; + + if (cached && age < EPG_STALE_MS) { + if (age >= EPG_FRESH_MS) { + void refreshCategoryEpgInBackground(tokenId, categoryId, streams, client); + } + return cached.titles; + } + + if (!opts.waitIfCold) { + void refreshCategoryEpgInBackground(tokenId, categoryId, streams, client); + return cached?.titles ?? new Map(); + } + + const titles = await fetchNowNextTitlesParallel( + client, + streams.map((s) => s.stream_id), + 14, + tokenId + ); + epgCategoryCache.set(key, { at: Date.now(), titles }); + return titles; +} + +function refreshCategoryEpgInBackground( tokenId: string, categoryId: string | undefined, streams: XtreamLiveStream[], client: XtreamClient -): Promise> { +) { const key = epgCacheKey(tokenId, categoryId); - const cached = epgNowCache.get(key); - if (cached && Date.now() - cached.at < EPG_NOW_CACHE_TTL_MS) { - return cached.titles; - } - - const titles = await fetchNowTitlesParallel( - client, - streams.map((s) => s.stream_id), - 14 - ); - epgNowCache.set(key, { at: Date.now(), titles }); - return titles; + if (epgCategoryRefreshing.has(key)) return; + epgCategoryRefreshing.add(key); + void (async () => { + try { + const titles = await fetchNowNextTitlesParallel( + client, + streams.map((s) => s.stream_id), + 14, + tokenId + ); + epgCategoryCache.set(key, { at: Date.now(), titles }); + } catch { + // stale cache blijft bruikbaar + } finally { + epgCategoryRefreshing.delete(key); + } + })(); } -async function fetchNowTitlesParallel( +async function fetchNowNextTitlesParallel( client: XtreamClient, streamIds: number[], - concurrency: number -): Promise> { - const result = new Map(); + concurrency: number, + tokenId?: string +): Promise> { + const result = new Map(); if (streamIds.length === 0) return result; const queue = [...streamIds]; @@ -276,8 +379,14 @@ async function fetchNowTitlesParallel( if (streamId == null) return; try { const epg = await client.getShortEpg(streamId, 4); - const { now } = parseNowNextPrograms(epg); - if (now?.title) result.set(streamId, now.title); + const { now, next } = parseNowNextPrograms(epg); + result.set(streamId, { + now: now?.title ?? null, + next: next?.title ?? null, + }); + if (tokenId) { + setCachedStreamEpg(tokenId, streamId, now, next); + } } catch { // zender zonder EPG overslaan } @@ -287,3 +396,57 @@ async function fetchNowTitlesParallel( await Promise.all(workers); return result; } + +/** + * Warm populaire categorieën op zodat TV-lijsten meteen Nu/Straks hebben. + * Draait periodiek op de master-api. + */ +export function startIptvEpgWarmup(config: Config) { + const WARM_CAT_LIMIT = 10; + const WARM_STREAM_CAP = 80; + const INTERVAL_MS = 12 * 60_000; + + const tick = async () => { + try { + const lines = await prisma.iptvLine.findMany({ + where: { enableLive: true, lastError: null, addonToken: { revoked: false } }, + select: { addonTokenId: true }, + take: 4, + orderBy: { updatedAt: "desc" }, + }); + + for (const line of lines) { + const xt = await loadXtreamForToken(line.addonTokenId, config); + if (!xt?.line.enableLive) continue; + + const cats = await xt.client.getLiveCategories().catch(() => []); + const sorted = [...cats].sort((a, b) => + friendlyCategoryName(a.category_name).localeCompare( + friendlyCategoryName(b.category_name), + "nl" + ) + ); + + for (const cat of sorted.slice(0, WARM_CAT_LIMIT)) { + const catId = String(cat.category_id); + const key = epgCacheKey(line.addonTokenId, catId); + const cached = epgCategoryCache.get(key); + if (cached && Date.now() - cached.at < EPG_FRESH_MS) continue; + + let streams = await xt.client.getLiveStreams(catId).catch(() => []); + streams = sortChannels(streams).slice(0, WARM_STREAM_CAP); + if (streams.length === 0) continue; + + await getCategoryEpgTitles(line.addonTokenId, catId, streams, xt.client, { + waitIfCold: true, + }); + } + } + } catch (err) { + console.warn("[iptv-epg-warmup]", err); + } + }; + + setTimeout(() => void tick(), 20_000); + setInterval(() => void tick(), INTERVAL_MS); +} diff --git a/apps/master-api/src/viewer/service.ts b/apps/master-api/src/viewer/service.ts index d0eae14..11d32c3 100644 --- a/apps/master-api/src/viewer/service.ts +++ b/apps/master-api/src/viewer/service.ts @@ -449,6 +449,7 @@ export class ViewerService { year: number | null; posterUrl: string | null; backdropUrl: string | null; + logoUrl: string | null; mediaFileId: string | null; addedAt: number; badge: string; @@ -478,6 +479,7 @@ export class ViewerService { year: m.year, posterUrl: m.posterUrl, backdropUrl: m.backdropUrl, + logoUrl: m.logoUrl, mediaFileId: m.mediaFiles[0]?.id ?? null, addedAt: added.get(m.id) ?? 0, badge: "Nieuw", @@ -501,6 +503,7 @@ export class ViewerService { year: s.year, posterUrl: s.posterUrl, backdropUrl: s.backdropUrl, + logoUrl: s.logoUrl, mediaFileId: null, addedAt: added.get(s.id) ?? 0, badge: "Nieuwe aflevering",