import { prisma } from "../database/client"; import type { Config } from "../config"; import { AppError } from "../security/errors"; import { loadXtreamByLineId, loadXtreamForViewer } from "../iptv/service"; import { decodeEpgText, epgToIso, setExternalEpgXmlUrl, type XtreamClient, type XtreamEpgListing, type XtreamLiveStream, type XmltvChannelGuide, type XmltvNowNext, type XmltvProgram, } 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\u2605\u2606\u25CF\u25AA\u2022\-\u2013\u2014]+/u, ""); 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 getViewerIptvStatus(viewerId: string, config: Config) { const xt = await loadXtreamForViewer(viewerId, 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 xt = await loadXtreamForViewer(viewerId, config); if (!xt) throw new AppError("FORBIDDEN", "Geen TV-toegang", 403); 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(() => [] as Awaited>), Promise.race([ xt.client.getLiveStreams().catch(() => [] as XtreamLiveStream[]), new Promise((resolve) => setTimeout(() => resolve([]), 6_000)), ]), ]); const cats = Array.isArray(catsRaw) ? catsRaw : []; const streams: XtreamLiveStream[] = 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\u2605\u2606\u25CF\u25AA\u2022\-\u2013\u2014]+/u, "")); 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 xt = await loadXtreamForViewer(viewerId, config); if (!xt) throw new AppError("FORBIDDEN", "Geen TV-toegang", 403); if (!xt?.line.enableLive) throw new AppError("FORBIDDEN", "Live TV niet beschikbaar", 403); setExternalEpgXmlUrl(config.IPTV_EPG_XML_URL); 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(xt.cacheKey, categoryId, streams, xt.client, { waitForXmltv: true, }) : peekCategoryEpgTitles(xt.cacheKey, categoryId); return { channels: streams.map((s) => { const pair = epgTitles?.get(s.stream_id); const cached = getCachedStreamEpg(xt.cacheKey, 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, nowImageUrl: cached?.now?.imageUrl ?? null, }; }), }; } export async function getViewerIptvPlayUrl( viewerId: string, config: Config, streamId: number ) { const xt = await loadXtreamForViewer(viewerId, config); if (!xt) throw new AppError("FORBIDDEN", "Geen TV-toegang", 403); 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"; void import("../iptv/active-watches") .then(({ touchIptvWatchForViewer }) => touchIptvWatchForViewer(viewerId, "live", streamId, ch.name) ) .catch(() => undefined); 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 xt = await loadXtreamForViewer(viewerId, config); if (!xt) throw new AppError("FORBIDDEN", "Geen TV-toegang", 403); if (!xt?.line.enableLive) throw new AppError("FORBIDDEN", "Live TV niet beschikbaar", 403); setExternalEpgXmlUrl(config.IPTV_EPG_XML_URL); 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); void import("../iptv/active-watches") .then(({ refreshIptvWatch }) => refreshIptvWatch({ iptvLineId: xt.cacheKey, streamId, title: ch.name }) ) .catch(() => undefined); const cached = getCachedStreamEpg(xt.cacheKey, streamId); if (cached && (cached.now || cached.next || cached.schedule.length > 0)) { return { streamId, channelName: ch.name, logoUrl: ch.stream_icon || null, now: cached.now, next: cached.next, schedule: cached.schedule, }; } // Probeer eerst rijke XMLTV (extern Odido + XUI) via epg_channel_id const epgId = ch.epg_channel_id?.trim(); if (epgId) { try { const guides = await xt.client.getXmltvGuidesByChannelId(); const hit = guides.get(epgId) ?? guides.get(epgId.toLowerCase()) ?? guides.get(epgId.toUpperCase()); if (hit && (hit.now || hit.next || hit.schedule.length > 0)) { setCachedStreamEpg(xt.cacheKey, streamId, hit.now, hit.next, hit.schedule); return { streamId, channelName: ch.name, logoUrl: ch.stream_icon || null, now: hit.now, next: hit.next, schedule: hit.schedule, }; } } catch { // fall through naar short_epg } } const epg = await xt.client.getShortEpg(streamId, 48); const { now: nowProgram, next: nextProgram, schedule } = parseRichShortEpg(epg); setCachedStreamEpg(xt.cacheKey, streamId, nowProgram, nextProgram, schedule); return { streamId, channelName: ch.name, logoUrl: ch.stream_icon || null, now: nowProgram, next: nextProgram, schedule, }; } /** 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: XmltvProgram | null; next: XmltvProgram | null; schedule: XmltvProgram[]; }; 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: XmltvProgram | null, next: XmltvProgram | null, schedule: XmltvProgram[] = [] ) { epgStreamCache.set(streamEpgKey(tokenId, streamId), { at: Date.now(), now, next, schedule, }); } function isNowPlayingFlag(v: XtreamEpgListing["now_playing"]): boolean { return v === 1 || v === true || v === "1"; } function listingToProgram(listing: XtreamEpgListing): XmltvProgram | null { const start = epgToIso(listing, "start"); if (!start) return null; const end = epgToIso(listing, "end") ?? null; const title = decodeEpgText(listing.title) || "Programma"; const description = decodeEpgText(listing.description) || null; return { title, subTitle: null, description: description || null, imageUrl: null, category: null, episode: null, start, end, }; } export function parseNowNextPrograms(epg: XtreamEpgListing[]): { now: XmltvProgram | null; next: XmltvProgram | null; } { const { now, next } = parseRichShortEpg(epg); return { now, next }; } export function parseRichShortEpg(epg: XtreamEpgListing[]): { now: XmltvProgram | null; next: XmltvProgram | null; schedule: XmltvProgram[]; } { const nowMs = Date.now(); const windowEnd = nowMs + 24 * 60 * 60_000; let nowProgram: XmltvProgram | null = null; let nextProgram: XmltvProgram | null = null; let flaggedNow: XmltvProgram | null = null; const schedule: XmltvProgram[] = []; for (const listing of epg) { const prog = listingToProgram(listing); if (!prog) continue; const startMs = Date.parse(prog.start); const endMs = prog.end ? Date.parse(prog.end) : startMs + 30 * 60_000; if (!Number.isFinite(startMs)) continue; if (isNowPlayingFlag(listing.now_playing)) { flaggedNow = prog; } if (endMs >= nowMs - 30 * 60_000 && startMs <= windowEnd) { schedule.push(prog); } if (startMs <= nowMs && nowMs < endMs) { nowProgram = prog; } else if (startMs > nowMs && !nextProgram) { nextProgram = prog; } } if (!nowProgram && flaggedNow) nowProgram = flaggedNow; schedule.sort((a, b) => Date.parse(a.start) - Date.parse(b.start)); return { now: nowProgram, next: nextProgram, schedule }; } 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 || ("schedule" in hit && Array.isArray(hit.schedule) && hit.schedule.length > 0)) { const schedule: XmltvProgram[] = "schedule" in hit && Array.isArray(hit.schedule) ? (hit.schedule as XmltvProgram[]) : []; setCachedStreamEpg(tokenId, s.stream_id, hit.now, hit.next, schedule); } } 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.getXmltvGuidesByChannelId(); 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; setExternalEpgXmlUrl(config.IPTV_EPG_XML_URL); const tick = async () => { try { setExternalEpgXmlUrl(config.IPTV_EPG_XML_URL); const lines = await prisma.iptvLine.findMany({ where: { enableLive: true, lastError: null, OR: [ { viewerUserId: { not: null } }, { addonToken: { is: { revoked: false } } }, ], }, select: { id: true }, take: 8, orderBy: { updatedAt: "desc" }, }); for (const line of lines) { if (xmltvLineRefreshing.has(line.id)) continue; xmltvLineRefreshing.add(line.id); try { const xt = await loadXtreamByLineId(line.id, config); const cacheKey = xt?.cacheKey ?? line.id; if (!xt?.line.enableLive) continue; // Primair: hele rijke XMLTV warmen (ttl=0 forceert refresh) try { await xt.client.getXmltvGuidesByChannelId(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(cacheKey, catId, streams, xt.client, { waitForXmltv: true, }); } const all = sortChannels(await xt.client.getLiveStreams().catch(() => [])); if (all.length > 0) { await resolveCategoryEpgTitles(cacheKey, "__all__", all, xt.client, { waitForXmltv: true, }); } } finally { xmltvLineRefreshing.delete(line.id); } } } catch (err) { console.warn("[iptv-epg-warmup]", err); } }; setTimeout(() => void tick(), 12_000); setInterval(() => void tick(), INTERVAL_MS); }