stremio/apps/master-api/src/viewer/iptv-client.ts
Jos Vooges | STH 318f613b40 Restore IPTV categories; only drop provider aggregate folders.
Keep one Alle zenders entry, backfill per-category counts when the bulk stream list is incomplete, and do not block the linked IPTV line on lastError.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-07 19:20:31 +02:00

721 lines
22 KiB
TypeScript

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";
/** Maak Xtream-categorienamen leesbaarder voor TV-gebruikers. */
export function friendlyCategoryName(raw: string): string {
let s = raw.trim();
s = s.replace(/^[A-Z]{2,3}\s*[|·:]\s*/i, "");
s = s.replace(/^[\s★☆●▪•\-–—]+/, "");
return s.trim() || raw.trim() || "Overig";
}
export async function resolveViewerIptvToken(viewerId: string): Promise<string | null> {
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);
const catsRaw = await xt.client.getLiveCategories().catch(() => null);
const streamsRaw = await xt.client.getLiveStreams().catch(() => null);
const cats = Array.isArray(catsRaw) ? catsRaw : [];
const streams = Array.isArray(streamsRaw) ? streamsRaw : [];
const counts = new Map<string, number>();
for (const s of streams) {
const key = normalizeIptvCategoryId(s.category_id);
if (!key) continue;
counts.set(key, (counts.get(key) ?? 0) + 1);
}
// Als de bulk-lijst incompleet is: per categorie bijtellen (gecached in XtreamClient)
const idsNeedingCount = cats
.map((c) => normalizeIptvCategoryId(c.category_id))
.filter((id) => id && (counts.get(id) ?? 0) === 0);
if (idsNeedingCount.length > 0) {
await mapPool(idsNeedingCount, 8, async (id) => {
const list = await xt.client.getLiveStreams(id).catch(() => null);
if (Array.isArray(list) && list.length > 0) {
counts.set(id, list.length);
}
});
}
const categories = cats
.map((c) => {
const id = normalizeIptvCategoryId(c.category_id);
const rawName = c.category_name ?? "";
return {
id,
name: friendlyCategoryName(rawName),
rawName,
channelCount: counts.get(id) ?? 0,
};
})
// Alleen lege ids + provider-verzamelcategorieën weg (niet echte genre-planken)
.filter((c) => c.id.length > 0 && !isAggregateCategoryName(c.rawName))
.sort((a, b) => a.name.localeCompare(b.name, "nl"));
// Één verzamelknop van ons — niet van de IPTV-provider
const totalChannels =
streams.length > 0
? streams.length
: [...counts.values()].reduce((a, b) => a + b, 0);
categories.unshift({
id: "__all__",
name: "Alle zenders",
rawName: "Alle zenders",
channelCount: totalChannels,
});
return { categories };
}
/** Alleen echte verzamelbakken (Alles / All / Alle zenders), geen Sport/Films/etc. */
export function isAggregateCategoryName(raw: string): boolean {
const n = friendlyCategoryName(raw).toLowerCase().replace(/\s+/g, " ").trim();
return (
n === "alles" ||
n === "all" ||
n === "alle" ||
n === "alle zenders" ||
n === "all channels" ||
n === "all channel" ||
n === "all categories" ||
n === "all cats" ||
n === "everything"
);
}
/** 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;
}
async function mapPool<T>(items: T[], concurrency: number, fn: (item: T) => Promise<void>) {
const queue = [...items];
const workers = Array.from({ length: Math.min(concurrency, queue.length) }, async () => {
while (queue.length > 0) {
const item = queue.shift();
if (item === undefined) return;
await fn(item);
}
});
await Promise.all(workers);
}
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<number, EpgTitlePair>;
/** true als XMLTV + short_epg-pass voor missing klaar is / voldoende dekking */
complete: boolean;
};
const epgCategoryCache = new Map<string, CategoryEpgCacheEntry>();
const epgCategoryRefreshing = new Set<string>();
const xmltvLineRefreshing = new Set<string>();
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<string, StreamEpgCacheEntry>();
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<number, EpgTitlePair> | 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<number, EpgTitlePair>, 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<number, EpgTitlePair>,
from: Map<number, EpgTitlePair>
): Map<number, EpgTitlePair> {
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<string, XmltvNowNext>,
tokenId: string
): Map<number, EpgTitlePair> {
const titles = new Map<number, EpgTitlePair>();
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<Map<number, EpgTitlePair>> {
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<number, EpgTitlePair>(cached?.titles ?? []);
const applyXmltvMap = (xmlMap: Map<string, XmltvNowNext>) => {
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<Map<number, EpgTitlePair>> {
const result = new Map<number, EpgTitlePair>();
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);
}