Improve IPTV EPG caching with SWR, warmup, and spotlight logos.
Speeds up Live TV now/next titles via longer cache, background refresh, and category warmup; also expose logoUrl on spotlight items. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
826e95aaab
commit
16f3f3f93c
3 changed files with 201 additions and 33 deletions
|
|
@ -18,6 +18,7 @@ import { startGooglePlayReconciliationPoller } from "./google-play/service";
|
||||||
import { nodeConnectionManager } from "./websocket/manager";
|
import { nodeConnectionManager } from "./websocket/manager";
|
||||||
import { MetadataService } from "./metadata/service";
|
import { MetadataService } from "./metadata/service";
|
||||||
import { toErrorResponse } from "./security/errors";
|
import { toErrorResponse } from "./security/errors";
|
||||||
|
import { startIptvEpgWarmup } from "./viewer/iptv-client";
|
||||||
|
|
||||||
async function main() {
|
async function main() {
|
||||||
const config = loadConfig();
|
const config = loadConfig();
|
||||||
|
|
@ -81,6 +82,7 @@ async function main() {
|
||||||
const googlePlay = registerGooglePlayRoutes(app, config);
|
const googlePlay = registerGooglePlayRoutes(app, config);
|
||||||
registerSettingsRoutes(app, config, googlePlay);
|
registerSettingsRoutes(app, config, googlePlay);
|
||||||
startGooglePlayReconciliationPoller(googlePlay, config.GOOGLE_PLAY_SYNC_INTERVAL_MS);
|
startGooglePlayReconciliationPoller(googlePlay, config.GOOGLE_PLAY_SYNC_INTERVAL_MS);
|
||||||
|
startIptvEpgWarmup(config);
|
||||||
|
|
||||||
app.get("/api/v1/node/connect", { websocket: true }, (socket) => {
|
app.get("/api/v1/node/connect", { websocket: true }, (socket) => {
|
||||||
void nodeConnectionManager.handleConnection(socket);
|
void nodeConnectionManager.handleConnection(socket);
|
||||||
|
|
|
||||||
|
|
@ -105,6 +105,8 @@ function sortChannels(streams: XtreamLiveStream[]): XtreamLiveStream[] {
|
||||||
return [...streams].sort((a, b) => (a.num ?? a.stream_id) - (b.num ?? b.stream_id));
|
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(
|
export async function getViewerIptvChannels(
|
||||||
viewerId: string,
|
viewerId: string,
|
||||||
config: Config,
|
config: Config,
|
||||||
|
|
@ -123,19 +125,25 @@ export async function getViewerIptvChannels(
|
||||||
: await xt.client.getLiveStreams(categoryId);
|
: await xt.client.getLiveStreams(categoryId);
|
||||||
streams = sortChannels(streams);
|
streams = sortChannels(streams);
|
||||||
|
|
||||||
const nowTitles = includeEpg
|
// Bij includeEpg=false toch warm cache meeleveren (instant Nu-titels).
|
||||||
? await getCachedCategoryNowTitles(tokenId, categoryId, streams, xt.client)
|
// Bij includeEpg=true: SWR — direct cache, achtergrond-refresh; bij cold miss wachten.
|
||||||
: null;
|
const epgTitles = includeEpg
|
||||||
|
? await getCategoryEpgTitles(tokenId, categoryId, streams, xt.client, { waitIfCold: true })
|
||||||
|
: peekCategoryEpgTitles(tokenId, categoryId);
|
||||||
|
|
||||||
return {
|
return {
|
||||||
channels: streams.map((s) => ({
|
channels: streams.map((s) => {
|
||||||
streamId: s.stream_id,
|
const pair = epgTitles?.get(s.stream_id);
|
||||||
name: s.name.trim(),
|
return {
|
||||||
logoUrl: s.stream_icon || null,
|
streamId: s.stream_id,
|
||||||
number: s.num ?? null,
|
name: s.name.trim(),
|
||||||
categoryId: s.category_id ? String(s.category_id) : null,
|
logoUrl: s.stream_icon || null,
|
||||||
nowTitle: nowTitles?.get(s.stream_id) ?? 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);
|
const ch = streams.find((s) => s.stream_id === streamId);
|
||||||
if (!ch) throw new AppError("NOT_FOUND", "Zender niet gevonden", 404);
|
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 epg = await xt.client.getShortEpg(streamId, 8);
|
||||||
const { now: nowProgram, next: nextProgram } = parseNowNextPrograms(epg);
|
const { now: nowProgram, next: nextProgram } = parseNowNextPrograms(epg);
|
||||||
|
setCachedStreamEpg(tokenId, streamId, nowProgram, nextProgram);
|
||||||
|
|
||||||
return {
|
return {
|
||||||
streamId,
|
streamId,
|
||||||
|
|
@ -207,13 +227,46 @@ export async function getViewerIptvEpg(viewerId: string, config: Config, streamI
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
const EPG_NOW_CACHE_TTL_MS = 90_000;
|
/** Fris tot 5 min; stale tot 20 min (SWR). */
|
||||||
const epgNowCache = new Map<string, { at: number; titles: Map<number, string> }>();
|
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<number, EpgTitlePair> };
|
||||||
|
const epgCategoryCache = new Map<string, CategoryEpgCacheEntry>();
|
||||||
|
const epgCategoryRefreshing = 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 {
|
function epgCacheKey(tokenId: string, categoryId?: string): string {
|
||||||
return `${tokenId}:${categoryId?.trim() || "__all__"}`;
|
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[]): {
|
export function parseNowNextPrograms(epg: XtreamEpgListing[]): {
|
||||||
now: { title: string; start: string; end: string | null } | null;
|
now: { title: string; start: string; end: string | null } | null;
|
||||||
next: { 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 };
|
return { now: nowProgram, next: nextProgram };
|
||||||
}
|
}
|
||||||
|
|
||||||
async function getCachedCategoryNowTitles(
|
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;
|
||||||
|
}
|
||||||
|
|
||||||
|
async function getCategoryEpgTitles(
|
||||||
|
tokenId: string,
|
||||||
|
categoryId: string | undefined,
|
||||||
|
streams: XtreamLiveStream[],
|
||||||
|
client: XtreamClient,
|
||||||
|
opts: { waitIfCold: boolean }
|
||||||
|
): Promise<Map<number, EpgTitlePair>> {
|
||||||
|
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,
|
tokenId: string,
|
||||||
categoryId: string | undefined,
|
categoryId: string | undefined,
|
||||||
streams: XtreamLiveStream[],
|
streams: XtreamLiveStream[],
|
||||||
client: XtreamClient
|
client: XtreamClient
|
||||||
): Promise<Map<number, string>> {
|
) {
|
||||||
const key = epgCacheKey(tokenId, categoryId);
|
const key = epgCacheKey(tokenId, categoryId);
|
||||||
const cached = epgNowCache.get(key);
|
if (epgCategoryRefreshing.has(key)) return;
|
||||||
if (cached && Date.now() - cached.at < EPG_NOW_CACHE_TTL_MS) {
|
epgCategoryRefreshing.add(key);
|
||||||
return cached.titles;
|
void (async () => {
|
||||||
}
|
try {
|
||||||
|
const titles = await fetchNowNextTitlesParallel(
|
||||||
const titles = await fetchNowTitlesParallel(
|
client,
|
||||||
client,
|
streams.map((s) => s.stream_id),
|
||||||
streams.map((s) => s.stream_id),
|
14,
|
||||||
14
|
tokenId
|
||||||
);
|
);
|
||||||
epgNowCache.set(key, { at: Date.now(), titles });
|
epgCategoryCache.set(key, { at: Date.now(), titles });
|
||||||
return titles;
|
} catch {
|
||||||
|
// stale cache blijft bruikbaar
|
||||||
|
} finally {
|
||||||
|
epgCategoryRefreshing.delete(key);
|
||||||
|
}
|
||||||
|
})();
|
||||||
}
|
}
|
||||||
|
|
||||||
async function fetchNowTitlesParallel(
|
async function fetchNowNextTitlesParallel(
|
||||||
client: XtreamClient,
|
client: XtreamClient,
|
||||||
streamIds: number[],
|
streamIds: number[],
|
||||||
concurrency: number
|
concurrency: number,
|
||||||
): Promise<Map<number, string>> {
|
tokenId?: string
|
||||||
const result = new Map<number, string>();
|
): Promise<Map<number, EpgTitlePair>> {
|
||||||
|
const result = new Map<number, EpgTitlePair>();
|
||||||
if (streamIds.length === 0) return result;
|
if (streamIds.length === 0) return result;
|
||||||
|
|
||||||
const queue = [...streamIds];
|
const queue = [...streamIds];
|
||||||
|
|
@ -276,8 +379,14 @@ async function fetchNowTitlesParallel(
|
||||||
if (streamId == null) return;
|
if (streamId == null) return;
|
||||||
try {
|
try {
|
||||||
const epg = await client.getShortEpg(streamId, 4);
|
const epg = await client.getShortEpg(streamId, 4);
|
||||||
const { now } = parseNowNextPrograms(epg);
|
const { now, next } = parseNowNextPrograms(epg);
|
||||||
if (now?.title) result.set(streamId, now.title);
|
result.set(streamId, {
|
||||||
|
now: now?.title ?? null,
|
||||||
|
next: next?.title ?? null,
|
||||||
|
});
|
||||||
|
if (tokenId) {
|
||||||
|
setCachedStreamEpg(tokenId, streamId, now, next);
|
||||||
|
}
|
||||||
} catch {
|
} catch {
|
||||||
// zender zonder EPG overslaan
|
// zender zonder EPG overslaan
|
||||||
}
|
}
|
||||||
|
|
@ -287,3 +396,57 @@ async function fetchNowTitlesParallel(
|
||||||
await Promise.all(workers);
|
await Promise.all(workers);
|
||||||
return result;
|
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);
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -449,6 +449,7 @@ export class ViewerService {
|
||||||
year: number | null;
|
year: number | null;
|
||||||
posterUrl: string | null;
|
posterUrl: string | null;
|
||||||
backdropUrl: string | null;
|
backdropUrl: string | null;
|
||||||
|
logoUrl: string | null;
|
||||||
mediaFileId: string | null;
|
mediaFileId: string | null;
|
||||||
addedAt: number;
|
addedAt: number;
|
||||||
badge: string;
|
badge: string;
|
||||||
|
|
@ -478,6 +479,7 @@ export class ViewerService {
|
||||||
year: m.year,
|
year: m.year,
|
||||||
posterUrl: m.posterUrl,
|
posterUrl: m.posterUrl,
|
||||||
backdropUrl: m.backdropUrl,
|
backdropUrl: m.backdropUrl,
|
||||||
|
logoUrl: m.logoUrl,
|
||||||
mediaFileId: m.mediaFiles[0]?.id ?? null,
|
mediaFileId: m.mediaFiles[0]?.id ?? null,
|
||||||
addedAt: added.get(m.id) ?? 0,
|
addedAt: added.get(m.id) ?? 0,
|
||||||
badge: "Nieuw",
|
badge: "Nieuw",
|
||||||
|
|
@ -501,6 +503,7 @@ export class ViewerService {
|
||||||
year: s.year,
|
year: s.year,
|
||||||
posterUrl: s.posterUrl,
|
posterUrl: s.posterUrl,
|
||||||
backdropUrl: s.backdropUrl,
|
backdropUrl: s.backdropUrl,
|
||||||
|
logoUrl: s.logoUrl,
|
||||||
mediaFileId: null,
|
mediaFileId: null,
|
||||||
addedAt: added.get(s.id) ?? 0,
|
addedAt: added.get(s.id) ?? 0,
|
||||||
badge: "Nieuwe aflevering",
|
badge: "Nieuwe aflevering",
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue