Fix Live TV EPG by using bulk XMLTV instead of partial short_epg.
Category now/next titles were incomplete because warmup capped streams and per-channel fetches timed out; load xmltv.php once, fix timezone parsing, and fill remaining gaps in the background. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
16f3f3f93c
commit
388962e6d1
3 changed files with 510 additions and 74 deletions
|
|
@ -69,6 +69,13 @@ export type XtreamEpgListing = {
|
|||
channel_id?: string;
|
||||
start_timestamp?: number;
|
||||
stop_timestamp?: number;
|
||||
/** Sommige panels markeren het huidige programma expliciet. */
|
||||
now_playing?: number | boolean | string;
|
||||
};
|
||||
|
||||
export type XmltvNowNext = {
|
||||
now: { title: string; start: string; end: string | null } | null;
|
||||
next: { title: string; start: string; end: string | null } | null;
|
||||
};
|
||||
|
||||
type CacheEntry<T> = { at: number; data: T };
|
||||
|
|
@ -254,6 +261,51 @@ export class XtreamClient {
|
|||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Bulk EPG via xmltv.php — één download voor alle zenders (veel betrouwbaarder
|
||||
* dan honderden get_short_epg-calls).
|
||||
*/
|
||||
async getXmltvNowNextByChannelId(ttlMs = 20 * 60_000): Promise<Map<string, XmltvNowNext>> {
|
||||
const cacheKey = "xmltv:now-next";
|
||||
const cached = this.cacheGet<Map<string, XmltvNowNext>>(cacheKey, ttlMs);
|
||||
if (cached) return cached;
|
||||
|
||||
const user = encodeURIComponent(this.username);
|
||||
const pass = encodeURIComponent(this.password);
|
||||
const urls = [
|
||||
`${this.baseUrl}/xmltv.php?username=${user}&password=${pass}`,
|
||||
`${this.baseUrl}/xmltv.php?username=${user}&password=${pass}&type=xml`,
|
||||
];
|
||||
|
||||
let xml = "";
|
||||
let lastErr: unknown;
|
||||
for (const url of urls) {
|
||||
try {
|
||||
const res = await fetch(url, {
|
||||
headers: {
|
||||
"User-Agent": "MediaCluster/1.0",
|
||||
Accept: "application/xml,text/xml,*/*",
|
||||
"Accept-Encoding": "gzip, deflate",
|
||||
},
|
||||
signal: AbortSignal.timeout(120_000),
|
||||
});
|
||||
if (!res.ok) throw new Error(`XMLTV HTTP ${res.status}`);
|
||||
xml = await res.text();
|
||||
if (xml.includes("<programme") || xml.includes("<tv")) break;
|
||||
xml = "";
|
||||
} catch (err) {
|
||||
lastErr = err;
|
||||
}
|
||||
}
|
||||
if (!xml) {
|
||||
throw lastErr instanceof Error ? lastErr : new Error("XMLTV niet beschikbaar");
|
||||
}
|
||||
|
||||
const map = parseXmltvNowNext(xml);
|
||||
this.cacheSet(cacheKey, map);
|
||||
return map;
|
||||
}
|
||||
|
||||
/**
|
||||
* XUI playlists: output=hls | ts | llhls (LLOD low-latency, XUI 1.5+).
|
||||
* On-demand + LLOD v3 levert HLS-segmenten — geen raw TS continuous stream.
|
||||
|
|
@ -394,33 +446,188 @@ export function parseXuiPlaylistPlayUrls(m3u: string): Map<number, string> {
|
|||
|
||||
export function decodeEpgText(raw?: string): string {
|
||||
if (!raw) return "";
|
||||
const unescaped = raw
|
||||
.replace(/&/g, "&")
|
||||
.replace(/</g, "<")
|
||||
.replace(/>/g, ">")
|
||||
.replace(/"/g, '"')
|
||||
.replace(/'/g, "'");
|
||||
try {
|
||||
// Xtream often base64-encodes title/description
|
||||
if (/^[A-Za-z0-9+/=]+$/.test(raw) && raw.length % 4 === 0) {
|
||||
const decoded = Buffer.from(raw, "base64").toString("utf8");
|
||||
if (decoded && !decoded.includes("\uFFFD")) return decoded;
|
||||
if (/^[A-Za-z0-9+/=]+$/.test(unescaped) && unescaped.length % 4 === 0 && unescaped.length >= 8) {
|
||||
const decoded = Buffer.from(unescaped, "base64").toString("utf8");
|
||||
if (decoded && !decoded.includes("\uFFFD") && /[\p{L}\p{N}]/u.test(decoded)) {
|
||||
return decoded.trim();
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
// fall through
|
||||
}
|
||||
return raw;
|
||||
return unescaped.trim();
|
||||
}
|
||||
|
||||
/** Parse naive wall-clock in Europe/Amsterdam to epoch ms. */
|
||||
export function parseAmsterdamLocalMs(raw: string): number | null {
|
||||
const m = raw
|
||||
.trim()
|
||||
.match(/^(\d{4})-(\d{2})-(\d{2})[ T](\d{2}):(\d{2})(?::(\d{2}))?/);
|
||||
if (!m) return null;
|
||||
const y = Number(m[1]);
|
||||
const mo = Number(m[2]);
|
||||
const d = Number(m[3]);
|
||||
const h = Number(m[4]);
|
||||
const mi = Number(m[5]);
|
||||
const s = Number(m[6] ?? "0");
|
||||
const desiredAsUtc = Date.UTC(y, mo - 1, d, h, mi, s);
|
||||
const dtf = new Intl.DateTimeFormat("en-US", {
|
||||
timeZone: "Europe/Amsterdam",
|
||||
year: "numeric",
|
||||
month: "2-digit",
|
||||
day: "2-digit",
|
||||
hour: "2-digit",
|
||||
minute: "2-digit",
|
||||
second: "2-digit",
|
||||
hour12: false,
|
||||
});
|
||||
const asLocalUtc = (ms: number) => {
|
||||
const parts = dtf.formatToParts(new Date(ms));
|
||||
const get = (type: Intl.DateTimeFormatPartTypes) =>
|
||||
parts.find((p) => p.type === type)?.value ?? "0";
|
||||
let hour = Number(get("hour"));
|
||||
if (hour === 24) hour = 0;
|
||||
return Date.UTC(
|
||||
Number(get("year")),
|
||||
Number(get("month")) - 1,
|
||||
Number(get("day")),
|
||||
hour,
|
||||
Number(get("minute")),
|
||||
Number(get("second"))
|
||||
);
|
||||
};
|
||||
return desiredAsUtc + (desiredAsUtc - asLocalUtc(desiredAsUtc));
|
||||
}
|
||||
|
||||
export function parseXmltvTimeMs(raw: string): number | null {
|
||||
const t = raw.trim();
|
||||
// 20260905160000 +0200 / 20260905160000+0200 / 20260905160000
|
||||
const m = t.match(/^(\d{4})(\d{2})(\d{2})(\d{2})(\d{2})(\d{2})\s*([+-]\d{2})(\d{2})?/);
|
||||
if (m) {
|
||||
const y = Number(m[1]);
|
||||
const mo = Number(m[2]);
|
||||
const d = Number(m[3]);
|
||||
const h = Number(m[4]);
|
||||
const mi = Number(m[5]);
|
||||
const s = Number(m[6]);
|
||||
const sign = m[7].startsWith("-") ? -1 : 1;
|
||||
const offH = Number(m[7].replace(/[+-]/, ""));
|
||||
const offM = Number(m[8] ?? "0");
|
||||
const utc = Date.UTC(y, mo - 1, d, h, mi, s) - sign * (offH * 3600_000 + offM * 60_000);
|
||||
return utc;
|
||||
}
|
||||
const compact = t.match(/^(\d{4})(\d{2})(\d{2})(\d{2})(\d{2})(\d{2})$/);
|
||||
if (compact) {
|
||||
// Geen offset → Amsterdam lokale tijd
|
||||
return parseAmsterdamLocalMs(
|
||||
`${compact[1]}-${compact[2]}-${compact[3]} ${compact[4]}:${compact[5]}:${compact[6]}`
|
||||
);
|
||||
}
|
||||
return parseAmsterdamLocalMs(t) ?? (Number.isFinite(Date.parse(t)) ? Date.parse(t) : null);
|
||||
}
|
||||
|
||||
export function epgToIso(listing: XtreamEpgListing, which: "start" | "end"): string | undefined {
|
||||
const ts = which === "start" ? listing.start_timestamp : listing.stop_timestamp;
|
||||
if (ts && Number(ts) > 0) {
|
||||
return new Date(Number(ts) * 1000).toISOString();
|
||||
// Panels leveren soms ms i.p.v. seconden
|
||||
const n = Number(ts);
|
||||
const ms = n > 1e12 ? n : n * 1000;
|
||||
return new Date(ms).toISOString();
|
||||
}
|
||||
const raw = which === "start" ? listing.start : listing.end;
|
||||
if (!raw) return undefined;
|
||||
// Formats: "2026-08-25 18:00:00" or "20260825180000 +0000"
|
||||
const compact = raw.match(/^(\d{4})(\d{2})(\d{2})(\d{2})(\d{2})(\d{2})/);
|
||||
if (compact) {
|
||||
const iso = `${compact[1]}-${compact[2]}-${compact[3]}T${compact[4]}:${compact[5]}:${compact[6]}Z`;
|
||||
const d = new Date(iso);
|
||||
if (!Number.isNaN(d.getTime())) return d.toISOString();
|
||||
|
||||
const withOffset = String(raw).match(
|
||||
/^(\d{4})(\d{2})(\d{2})(\d{2})(\d{2})(\d{2})\s*([+-]\d{2})(\d{2})?/
|
||||
);
|
||||
if (withOffset) {
|
||||
const ms = parseXmltvTimeMs(String(raw));
|
||||
if (ms != null) return new Date(ms).toISOString();
|
||||
}
|
||||
const d = new Date(raw.replace(" ", "T") + (raw.includes("Z") || raw.includes("+") ? "" : "Z"));
|
||||
|
||||
const compact = String(raw).match(/^(\d{4})(\d{2})(\d{2})(\d{2})(\d{2})(\d{2})$/);
|
||||
if (compact) {
|
||||
const ms = parseAmsterdamLocalMs(
|
||||
`${compact[1]}-${compact[2]}-${compact[3]} ${compact[4]}:${compact[5]}:${compact[6]}`
|
||||
);
|
||||
if (ms != null) return new Date(ms).toISOString();
|
||||
}
|
||||
|
||||
const amsterdam = parseAmsterdamLocalMs(String(raw));
|
||||
if (amsterdam != null) return new Date(amsterdam).toISOString();
|
||||
|
||||
const d = new Date(String(raw));
|
||||
if (!Number.isNaN(d.getTime())) return d.toISOString();
|
||||
return undefined;
|
||||
}
|
||||
|
||||
/**
|
||||
* Lightweight XMLTV parser → now/next per channel id.
|
||||
* Vermijdt zware XML-deps; werkt op typische Xtream xmltv.php output.
|
||||
*/
|
||||
export function parseXmltvNowNext(xml: string): Map<string, XmltvNowNext> {
|
||||
const now = Date.now();
|
||||
type Prog = { start: number; end: number; title: string };
|
||||
const byChannel = new Map<string, Prog[]>();
|
||||
|
||||
const progRe =
|
||||
/<programme\b([^>]*)>([\s\S]*?)<\/programme>/gi;
|
||||
let match: RegExpExecArray | null;
|
||||
while ((match = progRe.exec(xml)) !== null) {
|
||||
const attrs = match[1];
|
||||
const body = match[2];
|
||||
const channel = attrs.match(/\bchannel="([^"]+)"/i)?.[1];
|
||||
const startRaw = attrs.match(/\bstart="([^"]+)"/i)?.[1];
|
||||
const stopRaw = attrs.match(/\bstop="([^"]+)"/i)?.[1];
|
||||
if (!channel || !startRaw) continue;
|
||||
const start = parseXmltvTimeMs(startRaw);
|
||||
if (start == null) continue;
|
||||
const end = stopRaw ? parseXmltvTimeMs(stopRaw) : start + 30 * 60_000;
|
||||
if (end == null) continue;
|
||||
// Alleen programma's rond nu bewaren (geheugen)
|
||||
if (end < now - 2 * 60 * 60_000 || start > now + 6 * 60 * 60_000) continue;
|
||||
|
||||
const titleTag =
|
||||
body.match(/<title\b[^>]*>([\s\S]*?)<\/title>/i)?.[1] ??
|
||||
body.match(/<title>([\s\S]*?)<\/title>/i)?.[1] ??
|
||||
"";
|
||||
const title = decodeEpgText(titleTag.replace(/<!\[CDATA\[([\s\S]*?)\]\]>/g, "$1")) || "Programma";
|
||||
const list = byChannel.get(channel) ?? [];
|
||||
list.push({ start, end, title });
|
||||
byChannel.set(channel, list);
|
||||
}
|
||||
|
||||
const result = new Map<string, XmltvNowNext>();
|
||||
for (const [channel, progs] of byChannel) {
|
||||
progs.sort((a, b) => a.start - b.start);
|
||||
let nowProg: XmltvNowNext["now"] = null;
|
||||
let nextProg: XmltvNowNext["next"] = null;
|
||||
for (const p of progs) {
|
||||
if (p.start <= now && now < p.end) {
|
||||
nowProg = {
|
||||
title: p.title,
|
||||
start: new Date(p.start).toISOString(),
|
||||
end: new Date(p.end).toISOString(),
|
||||
};
|
||||
} else if (p.start > now && !nextProg) {
|
||||
nextProg = {
|
||||
title: p.title,
|
||||
start: new Date(p.start).toISOString(),
|
||||
end: new Date(p.end).toISOString(),
|
||||
};
|
||||
}
|
||||
}
|
||||
if (nowProg || nextProg) {
|
||||
result.set(channel, { now: nowProg, next: nextProg });
|
||||
}
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
|
|
|||
33
apps/master-api/src/viewer/epg-parse.test.ts
Normal file
33
apps/master-api/src/viewer/epg-parse.test.ts
Normal file
|
|
@ -0,0 +1,33 @@
|
|||
import { decodeEpgText, epgToIso, parseAmsterdamLocalMs, parseXmltvTimeMs } from "../iptv/xtream";
|
||||
import { parseNowNextPrograms } from "./iptv-client";
|
||||
|
||||
function assert(cond: unknown, msg: string) {
|
||||
if (!cond) throw new Error(msg);
|
||||
}
|
||||
|
||||
const startMs = parseXmltvTimeMs("20260905140000 +0200");
|
||||
assert(startMs === Date.UTC(2026, 8, 5, 12, 0, 0), `offset parse got ${startMs}`);
|
||||
|
||||
const local = parseAmsterdamLocalMs("2026-09-05 14:00:00");
|
||||
assert(local != null && Number.isFinite(local), "amsterdam local");
|
||||
|
||||
const iso = epgToIso(
|
||||
{ title: "X", start: "2026-09-05 14:00:00", end: "2026-09-05 15:00:00" },
|
||||
"start"
|
||||
);
|
||||
assert(iso != null, "epgToIso naive");
|
||||
|
||||
assert(decodeEpgText(Buffer.from("Hallo").toString("base64")).includes("Hallo"), "b64");
|
||||
|
||||
const listed = parseNowNextPrograms([
|
||||
{
|
||||
title: Buffer.from("Nu live").toString("base64"),
|
||||
start: "2000-01-01 00:00:00",
|
||||
end: "2100-01-01 00:00:00",
|
||||
start_timestamp: Math.floor(Date.now() / 1000) - 60,
|
||||
stop_timestamp: Math.floor(Date.now() / 1000) + 3600,
|
||||
},
|
||||
]);
|
||||
assert(listed.now?.title === "Nu live", `now title=${listed.now?.title}`);
|
||||
|
||||
console.log("epg-parse tests ok");
|
||||
|
|
@ -8,6 +8,7 @@ import {
|
|||
type XtreamClient,
|
||||
type XtreamEpgListing,
|
||||
type XtreamLiveStream,
|
||||
type XmltvNowNext,
|
||||
} from "../iptv/xtream";
|
||||
|
||||
/** Maak Xtream-categorienamen leesbaarder voor TV-gebruikers. */
|
||||
|
|
@ -125,10 +126,10 @@ export async function getViewerIptvChannels(
|
|||
: await xt.client.getLiveStreams(categoryId);
|
||||
streams = sortChannels(streams);
|
||||
|
||||
// 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 })
|
||||
? await resolveCategoryEpgTitles(tokenId, categoryId, streams, xt.client, {
|
||||
waitForXmltv: true,
|
||||
})
|
||||
: peekCategoryEpgTitles(tokenId, categoryId);
|
||||
|
||||
return {
|
||||
|
|
@ -164,7 +165,6 @@ export async function getViewerIptvPlayUrl(
|
|||
|
||||
const resolved = await xt.client.resolveLivePlayUrls(streamId);
|
||||
|
||||
// Smarters-gedrag: klassiek Xtream-pad eerst (/live/user/pass/id.ts), daarna playlist-fallbacks
|
||||
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;
|
||||
|
|
@ -204,7 +204,7 @@ export async function getViewerIptvEpg(viewerId: string, config: Config, streamI
|
|||
if (!ch) throw new AppError("NOT_FOUND", "Zender niet gevonden", 404);
|
||||
|
||||
const cached = getCachedStreamEpg(tokenId, streamId);
|
||||
if (cached) {
|
||||
if (cached && (cached.now || cached.next)) {
|
||||
return {
|
||||
streamId,
|
||||
channelName: ch.name,
|
||||
|
|
@ -214,7 +214,28 @@ export async function getViewerIptvEpg(viewerId: string, config: Config, streamI
|
|||
};
|
||||
}
|
||||
|
||||
const epg = await xt.client.getShortEpg(streamId, 8);
|
||||
// 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);
|
||||
|
||||
|
|
@ -227,14 +248,20 @@ export async function getViewerIptvEpg(viewerId: string, config: Config, streamI
|
|||
};
|
||||
}
|
||||
|
||||
/** 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;
|
||||
/** 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> };
|
||||
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;
|
||||
|
|
@ -267,21 +294,33 @@ function setCachedStreamEpg(
|
|||
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();
|
||||
let nowProgram: { title: string; start: string; end: string | null } | null = null;
|
||||
let nextProgram: { title: string; start: string; end: string | null } | null = null;
|
||||
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");
|
||||
if (!start) continue;
|
||||
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 };
|
||||
|
|
@ -290,6 +329,9 @@ export function parseNowNextPrograms(epg: XtreamEpgListing[]): {
|
|||
}
|
||||
}
|
||||
|
||||
// Tijd-match heeft voorrang; now_playing als fallback als tijd mist/scheef zit
|
||||
if (!nowProgram && flaggedNow) nowProgram = flaggedNow;
|
||||
|
||||
return { now: nowProgram, next: nextProgram };
|
||||
}
|
||||
|
||||
|
|
@ -304,36 +346,148 @@ function peekCategoryEpgTitles(
|
|||
return cached.titles;
|
||||
}
|
||||
|
||||
async function getCategoryEpgTitles(
|
||||
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: { waitIfCold: boolean }
|
||||
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;
|
||||
|
||||
if (cached && age < EPG_STALE_MS) {
|
||||
// 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;
|
||||
}
|
||||
|
||||
if (!opts.waitIfCold) {
|
||||
void refreshCategoryEpgInBackground(tokenId, categoryId, streams, client);
|
||||
return cached?.titles ?? new Map();
|
||||
// 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);
|
||||
}
|
||||
|
||||
const titles = await fetchNowNextTitlesParallel(
|
||||
client,
|
||||
streams.map((s) => s.stream_id),
|
||||
14,
|
||||
tokenId
|
||||
);
|
||||
epgCategoryCache.set(key, { at: Date.now(), titles });
|
||||
return titles;
|
||||
}
|
||||
|
||||
|
|
@ -348,21 +502,50 @@ function refreshCategoryEpgInBackground(
|
|||
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 });
|
||||
await resolveCategoryEpgTitles(tokenId, categoryId, streams, client, {
|
||||
waitForXmltv: true,
|
||||
});
|
||||
} catch {
|
||||
// stale cache blijft bruikbaar
|
||||
// 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[],
|
||||
|
|
@ -378,7 +561,7 @@ async function fetchNowNextTitlesParallel(
|
|||
const streamId = queue.shift();
|
||||
if (streamId == null) return;
|
||||
try {
|
||||
const epg = await client.getShortEpg(streamId, 4);
|
||||
const epg = await client.getShortEpg(streamId, 8);
|
||||
const { now, next } = parseNowNextPrograms(epg);
|
||||
result.set(streamId, {
|
||||
now: now?.title ?? null,
|
||||
|
|
@ -398,13 +581,10 @@ async function fetchNowNextTitlesParallel(
|
|||
}
|
||||
|
||||
/**
|
||||
* Warm populaire categorieën op zodat TV-lijsten meteen Nu/Straks hebben.
|
||||
* Draait periodiek op de master-api.
|
||||
* Warm XMLTV + categorie-caches op zodat TV-lijsten meteen Nu/Straks hebben.
|
||||
*/
|
||||
export function startIptvEpgWarmup(config: Config) {
|
||||
const WARM_CAT_LIMIT = 10;
|
||||
const WARM_STREAM_CAP = 80;
|
||||
const INTERVAL_MS = 12 * 60_000;
|
||||
const INTERVAL_MS = 15 * 60_000;
|
||||
|
||||
const tick = async () => {
|
||||
try {
|
||||
|
|
@ -416,30 +596,46 @@ export function startIptvEpgWarmup(config: Config) {
|
|||
});
|
||||
|
||||
for (const line of lines) {
|
||||
const xt = await loadXtreamForToken(line.addonTokenId, config);
|
||||
if (!xt?.line.enableLive) continue;
|
||||
if (xmltvLineRefreshing.has(line.addonTokenId)) continue;
|
||||
xmltvLineRefreshing.add(line.addonTokenId);
|
||||
try {
|
||||
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"
|
||||
)
|
||||
);
|
||||
// Primair: hele XMLTV warmen (ttl=0 forceert refresh)
|
||||
try {
|
||||
await xt.client.getXmltvNowNextByChannelId(0);
|
||||
} catch (err) {
|
||||
console.warn("[iptv-epg-warmup] xmltv", err);
|
||||
}
|
||||
|
||||
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;
|
||||
const cats = await xt.client.getLiveCategories().catch(() => []);
|
||||
const sorted = [...cats].sort((a, b) =>
|
||||
friendlyCategoryName(a.category_name).localeCompare(
|
||||
friendlyCategoryName(b.category_name),
|
||||
"nl"
|
||||
)
|
||||
);
|
||||
|
||||
let streams = await xt.client.getLiveStreams(catId).catch(() => []);
|
||||
streams = sortChannels(streams).slice(0, WARM_STREAM_CAP);
|
||||
if (streams.length === 0) continue;
|
||||
// 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,
|
||||
});
|
||||
}
|
||||
|
||||
await getCategoryEpgTitles(line.addonTokenId, catId, streams, xt.client, {
|
||||
waitIfCold: 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) {
|
||||
|
|
@ -447,6 +643,6 @@ export function startIptvEpgWarmup(config: Config) {
|
|||
}
|
||||
};
|
||||
|
||||
setTimeout(() => void tick(), 20_000);
|
||||
setTimeout(() => void tick(), 12_000);
|
||||
setInterval(() => void tick(), INTERVAL_MS);
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue