import { prisma } from "../database/client"; import type { Config } from "../config"; import { AppError } from "../security/errors"; import { loadXtreamForToken } from "../iptv/service"; import { decodeEpgText, epgToIso, type XtreamClient, type XtreamEpgListing, type XtreamLiveStream, type XmltvNowNext, } from "../iptv/xtream"; const AGGREGATE_LABELS = new Set([ "alles", "all", "alle", "alle zenders", "all channels", "all channel", "all categories", "all cats", "everything", ]); function normCatLabel(s: string): string { return s.toLowerCase().replace(/\s+/g, " ").trim(); } /** Maak Xtream-categorienamen leesbaarder voor TV-gebruikers. */ export function friendlyCategoryName(raw: string): string { let s = raw.trim().replace(/^[\s★☆●▪•\-–—]+/, ""); const m = s.match(/^([A-Z]{2,3})\s*[|·:]\s*(.+)$/i); if (m) { const rest = m[2].trim(); // "NL | ALL" is een landpakket — niet reduceren tot "ALL" (dat verdween daarna als verzamelbak) if (AGGREGATE_LABELS.has(normCatLabel(rest))) { return m[1].toUpperCase(); } s = rest; } return s.trim() || raw.trim() || "Overig"; } export async function resolveViewerIptvToken(viewerId: string): Promise { const viewer = await prisma.viewerUser.findUnique({ where: { id: viewerId }, select: { iptvAccess: true, iptvAddonTokenId: true, enabled: true, appAccess: true }, }); if (!viewer?.enabled || !viewer.appAccess || !viewer.iptvAccess) return null; // Gekoppelde lijn altijd proberen — lastError mag TV niet stilleggen if (viewer.iptvAddonTokenId) { const linked = await prisma.addonToken.findFirst({ where: { id: viewer.iptvAddonTokenId, revoked: false, iptvLine: { enableLive: true }, }, select: { id: true }, }); if (linked) return linked.id; } const fallbackOk = await prisma.iptvLine.findFirst({ where: { enableLive: true, lastError: null, addonToken: { revoked: false }, }, select: { addonTokenId: true }, orderBy: { updatedAt: "desc" }, }); if (fallbackOk) return fallbackOk.addonTokenId; const fallbackAny = await prisma.iptvLine.findFirst({ where: { enableLive: true, addonToken: { revoked: false }, }, select: { addonTokenId: true }, orderBy: { updatedAt: "desc" }, }); return fallbackAny?.addonTokenId ?? null; } export async function getViewerIptvStatus(viewerId: string, config: Config) { const tokenId = await resolveViewerIptvToken(viewerId); if (!tokenId) return { enabled: false as const }; const xt = await loadXtreamForToken(tokenId, config); if (!xt?.line.enableLive) return { enabled: false as const }; const cats = await xt.client.getLiveCategories().catch(() => []); return { enabled: true as const, categoryCount: cats.length, }; } export async function getViewerIptvCategories(viewerId: string, config: Config) { const tokenId = await resolveViewerIptvToken(viewerId); if (!tokenId) throw new AppError("FORBIDDEN", "Geen TV-toegang", 403); const xt = await loadXtreamForToken(tokenId, config); if (!xt?.line.enableLive) throw new AppError("FORBIDDEN", "Live TV niet beschikbaar", 403); // Categorieën + (korte) zenderlijst parallel — categorieën nooit laten falen op trage streams const [catsRaw, streamsRaw] = await Promise.all([ xt.client.getLiveCategories().catch(() => []), Promise.race([ xt.client.getLiveStreams().catch(() => []), new Promise((resolve) => setTimeout(() => resolve([]), 6_000)), ]), ]); const cats = Array.isArray(catsRaw) ? catsRaw : []; const streams = Array.isArray(streamsRaw) ? streamsRaw : []; const counts = new Map(); for (const s of streams) { const key = normalizeIptvCategoryId(s.category_id); if (!key) continue; counts.set(key, (counts.get(key) ?? 0) + 1); } let categories = cats .map((c) => { const id = normalizeIptvCategoryId(c.category_id); const rawName = String(c.category_name ?? ""); return { id, name: friendlyCategoryName(rawName), rawName, channelCount: counts.get(id) ?? 0, }; }) // Alleen pure provider-dump "Alles/All" weg — niet "NL | ALL" landpakketten .filter((c) => c.id.length > 0 && !isAggregateCategoryName(c.rawName)) .sort((a, b) => a.name.localeCompare(b.name, "nl")); // Fallback: bouquets uit zenders afleiden als get_live_categories leeg/fout is if (categories.length === 0 && counts.size > 0) { categories = [...counts.entries()] .map(([id, channelCount]) => ({ id, name: `Groep ${id}`, rawName: `Groep ${id}`, channelCount, })) .sort((a, b) => a.name.localeCompare(b.name, "nl")); } // Één eigen verzamelknop categories.unshift({ id: "__all__", name: "Alle zenders", rawName: "Alle zenders", channelCount: streams.length, }); return { categories }; } /** * Alleen pure verzamelbakken (Alles / All), géén landpakketten zoals "NL | ALL". * Match op ruwe naam (zonder landcode te strippen). */ export function isAggregateCategoryName(raw: string): boolean { const n = normCatLabel(String(raw ?? "").replace(/^[\s★☆●▪•\-–—]+/, "")); return AGGREGATE_LABELS.has(n); } /** Xtream geeft category_id soms als number, "12" of "12.0". */ export function normalizeIptvCategoryId(raw: unknown): string { if (raw == null) return ""; const s = String(raw).trim(); if (!s) return ""; if (/^\d+(\.0+)?$/.test(s)) return String(parseInt(s, 10)); return s; } 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, categoryId?: string, includeEpg = false ) { const tokenId = await resolveViewerIptvToken(viewerId); if (!tokenId) throw new AppError("FORBIDDEN", "Geen TV-toegang", 403); const xt = await loadXtreamForToken(tokenId, config); if (!xt?.line.enableLive) throw new AppError("FORBIDDEN", "Live TV niet beschikbaar", 403); let streams = !categoryId || categoryId === "__all__" ? await xt.client.getLiveStreams() : await xt.client.getLiveStreams(categoryId); if (!Array.isArray(streams)) streams = []; streams = sortChannels(streams); const epgTitles = includeEpg ? await resolveCategoryEpgTitles(tokenId, categoryId, streams, xt.client, { waitForXmltv: true, }) : peekCategoryEpgTitles(tokenId, categoryId); return { 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, }; }), }; } export async function getViewerIptvPlayUrl( viewerId: string, config: Config, streamId: number ) { const tokenId = await resolveViewerIptvToken(viewerId); if (!tokenId) throw new AppError("FORBIDDEN", "Geen TV-toegang", 403); const xt = await loadXtreamForToken(tokenId, config); if (!xt?.line.enableLive) throw new AppError("FORBIDDEN", "Live TV niet beschikbaar", 403); const streams = await xt.client.getLiveStreams(); const ch = streams.find((s) => s.stream_id === streamId); if (!ch) throw new AppError("NOT_FOUND", "Zender niet gevonden", 404); const resolved = await xt.client.resolveLivePlayUrls(streamId); const classicTs = resolved.classic.find((u) => /\.ts$/i.test(u)); const classicHls = resolved.classic.find((u) => /\.m3u8$/i.test(u) || u.includes("m3u8")); const playlistHls = resolved.hls ?? resolved.llhls; const playlistTs = resolved.ts; const candidates = [classicTs, playlistTs, classicHls, playlistHls, ...resolved.classic].filter( (u): u is string => Boolean(u) ); const unique = [...new Set(candidates)]; const primary = unique[0]; if (!primary) throw new AppError("UNAVAILABLE", "Geen afspeel-URL voor deze zender", 503); const fallback = unique[1]; const formatFor = (url: string): "hls" | "ts" => url.includes(".m3u8") || url.includes("m3u8") ? "hls" : "ts"; return { streamId, name: ch.name, logoUrl: ch.stream_icon || null, streamUrl: primary, format: formatFor(primary), fallbackStreamUrl: fallback ?? null, fallbackFormat: fallback ? formatFor(fallback) : null, }; } export async function getViewerIptvEpg(viewerId: string, config: Config, streamId: number) { const tokenId = await resolveViewerIptvToken(viewerId); if (!tokenId) throw new AppError("FORBIDDEN", "Geen TV-toegang", 403); const xt = await loadXtreamForToken(tokenId, config); if (!xt?.line.enableLive) throw new AppError("FORBIDDEN", "Live TV niet beschikbaar", 403); const streams = await xt.client.getLiveStreams(); 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 && (cached.now || cached.next)) { return { streamId, channelName: ch.name, logoUrl: ch.stream_icon || null, now: cached.now, next: cached.next, }; } // Probeer eerst XMLTV (bulk, al warm) via epg_channel_id const epgId = ch.epg_channel_id?.trim(); if (epgId) { try { const xmlMap = await xt.client.getXmltvNowNextByChannelId(); const hit = xmlMap.get(epgId) ?? xmlMap.get(epgId.toLowerCase()); if (hit && (hit.now || hit.next)) { setCachedStreamEpg(tokenId, streamId, hit.now, hit.next); return { streamId, channelName: ch.name, logoUrl: ch.stream_icon || null, now: hit.now, next: hit.next, }; } } catch { // fall through naar short_epg } } const epg = await xt.client.getShortEpg(streamId, 12); const { now: nowProgram, next: nextProgram } = parseNowNextPrograms(epg); setCachedStreamEpg(tokenId, streamId, nowProgram, nextProgram); return { streamId, channelName: ch.name, logoUrl: ch.stream_icon || null, now: nowProgram, next: nextProgram, }; } /** Fris tot 10 min; stale tot 30 min (SWR). */ const EPG_FRESH_MS = 10 * 60_000; const EPG_STALE_MS = 30 * 60_000; const EPG_DETAIL_TTL_MS = 20 * 60_000; type CategoryEpgCacheEntry = { at: number; titles: Map; /** true als XMLTV + short_epg-pass voor missing klaar is / voldoende dekking */ complete: boolean; }; const epgCategoryCache = new Map(); const epgCategoryRefreshing = new Set(); const xmltvLineRefreshing = 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 }); } function isNowPlayingFlag(v: XtreamEpgListing["now_playing"]): boolean { return v === 1 || v === true || v === "1"; } export function parseNowNextPrograms(epg: XtreamEpgListing[]): { now: { title: string; start: string; end: string | null } | null; next: { title: string; start: string; end: string | null } | null; } { const now = Date.now(); type Prog = { title: string; start: string; end: string | null }; let nowProgram: Prog | null = null; let nextProgram: Prog | null = null; let flaggedNow: Prog | null = null; for (const listing of epg) { const start = epgToIso(listing, "start"); const end = epgToIso(listing, "end"); const title = decodeEpgText(listing.title) || "Programma"; if (isNowPlayingFlag(listing.now_playing) && start) { flaggedNow = { title, start, end: end ?? null }; } if (!start) continue; const startMs = Date.parse(start); const endMs = end ? Date.parse(end) : startMs + 30 * 60_000; if (!Number.isFinite(startMs)) continue; if (startMs <= now && now < endMs) { nowProgram = { title, start, end: end ?? null }; } else if (startMs > now && !nextProgram) { nextProgram = { title, start, end: end ?? null }; } } // Tijd-match heeft voorrang; now_playing als fallback als tijd mist/scheef zit if (!nowProgram && flaggedNow) nowProgram = flaggedNow; return { now: nowProgram, next: nextProgram }; } 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; } function coverageRatio(titles: Map, streamIds: number[]): number { if (streamIds.length === 0) return 1; let hit = 0; for (const id of streamIds) { const t = titles.get(id); if (t && (t.now || t.next)) hit++; } return hit / streamIds.length; } function mergeTitles( into: Map, from: Map ): Map { for (const [id, pair] of from) { const prev = into.get(id); if (!prev) { into.set(id, pair); continue; } into.set(id, { now: pair.now ?? prev.now, next: pair.next ?? prev.next, }); } return into; } function titlesFromXmltv( streams: XtreamLiveStream[], xmlMap: Map, tokenId: string ): Map { const titles = new Map(); for (const s of streams) { const epgId = s.epg_channel_id?.trim(); if (!epgId) continue; const hit = xmlMap.get(epgId) ?? xmlMap.get(epgId.toLowerCase()) ?? xmlMap.get(epgId.toUpperCase()); if (!hit) continue; titles.set(s.stream_id, { now: hit.now?.title ?? null, next: hit.next?.title ?? null, }); if (hit.now || hit.next) { setCachedStreamEpg(tokenId, s.stream_id, hit.now, hit.next); } } return titles; } async function resolveCategoryEpgTitles( tokenId: string, categoryId: string | undefined, streams: XtreamLiveStream[], client: XtreamClient, opts: { waitForXmltv: boolean } ): Promise> { const key = epgCacheKey(tokenId, categoryId); const streamIds = streams.map((s) => s.stream_id); const cached = epgCategoryCache.get(key); const age = cached ? Date.now() - cached.at : Number.POSITIVE_INFINITY; const cov = cached ? coverageRatio(cached.titles, streamIds) : 0; // Goede, recente cache met fatsoenlijke dekking → direct terug, optioneel refresh if (cached && age < EPG_STALE_MS && cov >= 0.55 && cached.complete) { if (age >= EPG_FRESH_MS) { void refreshCategoryEpgInBackground(tokenId, categoryId, streams, client); } return cached.titles; } // Incomplete of cold: XMLTV bulk laden (hoofdbron) let titles = new Map(cached?.titles ?? []); const applyXmltvMap = (xmlMap: Map) => { titles = mergeTitles(titles, titlesFromXmltv(streams, xmlMap, tokenId)); }; const xmlPromise = client.getXmltvNowNextByChannelId(); if (opts.waitForXmltv) { // Max ~8s wachten zodat reverse-proxy niet timeout; rest via client-poll const raced = await Promise.race([ xmlPromise.then((m) => ({ ready: true as const, m })), new Promise<{ ready: false }>((resolve) => setTimeout(() => resolve({ ready: false }), 8_000)), ]); if (raced.ready) { applyXmltvMap(raced.m); } else { void xmlPromise .then((m) => { const merged = mergeTitles( new Map(epgCategoryCache.get(key)?.titles ?? titles), titlesFromXmltv(streams, m, tokenId) ); const ids = streams.map((s) => s.stream_id); epgCategoryCache.set(key, { at: Date.now(), titles: merged, complete: coverageRatio(merged, ids) >= 0.7, }); }) .catch((err) => console.warn("[iptv-epg] xmltv bg:", err)); } } else { void xmlPromise .then((m) => { const merged = mergeTitles( new Map(epgCategoryCache.get(key)?.titles ?? titles), titlesFromXmltv(streams, m, tokenId) ); const ids = streams.map((s) => s.stream_id); epgCategoryCache.set(key, { at: Date.now(), titles: merged, complete: coverageRatio(merged, ids) >= 0.7, }); }) .catch(() => undefined); } const afterXmltvCov = coverageRatio(titles, streamIds); epgCategoryCache.set(key, { at: Date.now(), titles, complete: afterXmltvCov >= 0.7, }); // Ontbrekende zenders (geen epg_channel_id / niet in XMLTV) via short_epg bijwerken const missing = streams .filter((s) => { const t = titles.get(s.stream_id); return !t?.now && !t?.next; }) .map((s) => s.stream_id); if (missing.length > 0) { void fillMissingShortEpg(tokenId, categoryId, missing, client); } return titles; } function refreshCategoryEpgInBackground( tokenId: string, categoryId: string | undefined, streams: XtreamLiveStream[], client: XtreamClient ) { const key = epgCacheKey(tokenId, categoryId); if (epgCategoryRefreshing.has(key)) return; epgCategoryRefreshing.add(key); void (async () => { try { await resolveCategoryEpgTitles(tokenId, categoryId, streams, client, { waitForXmltv: true, }); } catch { // stale blijft } finally { epgCategoryRefreshing.delete(key); } })(); } function fillMissingShortEpg( tokenId: string, categoryId: string | undefined, streamIds: number[], client: XtreamClient ) { const key = epgCacheKey(tokenId, categoryId); const jobKey = `${key}:missing`; if (epgCategoryRefreshing.has(jobKey)) return; epgCategoryRefreshing.add(jobKey); void (async () => { try { // Beperk burst: max 120 missing per ronde, concurrency 8 const batch = streamIds.slice(0, 120); const filled = await fetchNowNextTitlesParallel(client, batch, 8, tokenId); const cached = epgCategoryCache.get(key); const merged = mergeTitles(new Map(cached?.titles ?? []), filled); const allIds = [...merged.keys()]; // coverage t.o.v. deze fill-batch + bestaande epgCategoryCache.set(key, { at: Date.now(), titles: merged, complete: cached?.complete === true || coverageRatio(merged, allIds) >= 0.6, }); } catch { // ignore } finally { epgCategoryRefreshing.delete(jobKey); } })(); } async function fetchNowNextTitlesParallel( client: XtreamClient, streamIds: number[], concurrency: number, tokenId?: string ): Promise> { const result = new Map(); if (streamIds.length === 0) return result; const queue = [...streamIds]; const workers = Array.from({ length: Math.min(concurrency, queue.length) }, async () => { while (queue.length > 0) { const streamId = queue.shift(); if (streamId == null) return; try { const epg = await client.getShortEpg(streamId, 8); 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 } } }); await Promise.all(workers); return result; } /** * Warm XMLTV + categorie-caches op zodat TV-lijsten meteen Nu/Straks hebben. */ export function startIptvEpgWarmup(config: Config) { const INTERVAL_MS = 15 * 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) { if (xmltvLineRefreshing.has(line.addonTokenId)) continue; xmltvLineRefreshing.add(line.addonTokenId); try { const xt = await loadXtreamForToken(line.addonTokenId, config); if (!xt?.line.enableLive) continue; // Primair: hele XMLTV warmen (ttl=0 forceert refresh) try { await xt.client.getXmltvNowNextByChannelId(0); } catch (err) { console.warn("[iptv-epg-warmup] xmltv", err); } const cats = await xt.client.getLiveCategories().catch(() => []); const sorted = [...cats].sort((a, b) => friendlyCategoryName(a.category_name).localeCompare( friendlyCategoryName(b.category_name), "nl" ) ); // Top categorieën + alle zenders via XMLTV-mapping for (const cat of sorted.slice(0, 20)) { const catId = String(cat.category_id); let streams = await xt.client.getLiveStreams(catId).catch(() => []); streams = sortChannels(streams); if (streams.length === 0) continue; await resolveCategoryEpgTitles(line.addonTokenId, catId, streams, xt.client, { waitForXmltv: true, }); } const all = sortChannels(await xt.client.getLiveStreams().catch(() => [])); if (all.length > 0) { await resolveCategoryEpgTitles(line.addonTokenId, "__all__", all, xt.client, { waitForXmltv: true, }); } } finally { xmltvLineRefreshing.delete(line.addonTokenId); } } } catch (err) { console.warn("[iptv-epg-warmup]", err); } }; setTimeout(() => void tick(), 12_000); setInterval(() => void tick(), INTERVAL_MS); }