admin streams presence for events/live-lists with client heartbeat.

This commit is contained in:
Jos Vooges | STH 2026-09-19 01:04:04 +02:00
parent 3c356001b2
commit 79f15719a3
12 changed files with 378 additions and 24 deletions

View file

@ -5,7 +5,8 @@ import { Nav, useAuth } from "@/components/Nav";
interface Stream { interface Stream {
id: string; id: string;
source?: "library" | "iptv"; source?: "library" | "iptv" | "client";
kind?: string;
node: string; node: string;
status: string; status: string;
viewer: string; viewer: string;
@ -44,6 +45,18 @@ function relativeTime(iso: string): string {
return new Date(iso).toLocaleTimeString(); return new Date(iso).toLocaleTimeString();
} }
function isLivePresence(s: Stream): boolean {
return s.source === "iptv" || s.source === "client" || s.node === "IPTV";
}
function presenceHint(s: Stream): string {
if (s.source === "client" || s.kind === "event" || s.kind === "live-list") {
return s.kind === "event" ? "Live event · heartbeat" : "Custom live · heartbeat";
}
if (s.source === "iptv" || s.node === "IPTV") return "Direct naar Xtream";
return `${formatBytes(s.bytesSent)} gestreamd`;
}
export default function StreamsPage() { export default function StreamsPage() {
useAuth(); useAuth();
const [streams, setStreams] = useState<Stream[]>([]); const [streams, setStreams] = useState<Stream[]>([]);
@ -61,7 +74,11 @@ export default function StreamsPage() {
}, []); }, []);
async function revoke(id: string) { async function revoke(id: string) {
if (!confirm("Playback-sessie intrekken?")) return; const isClient = id.startsWith("cw-");
const msg = isClient
? "Presence wissen? De stream op het apparaat stopt niet (direct play)."
: "Playback-sessie intrekken?";
if (!confirm(msg)) return;
await fetch(`/api/v1/admin/streams/${id}/revoke`, { await fetch(`/api/v1/admin/streams/${id}/revoke`, {
method: "POST", method: "POST",
credentials: "include", credentials: "include",
@ -82,8 +99,9 @@ export default function StreamsPage() {
</div> </div>
<p className="muted page-lead"> <p className="muted page-lead">
Bibliotheek én Live TV / IPTV. Bibliotheek verdwijnt binnen ~3 min na stoppen; IPTV Bibliotheek, Live TV / IPTV, Live Events en custom live-lijsten. Bibliotheek verdwijnt
blijft zolang Xtream een actieve connectie ziet. binnen ~3 min na stoppen; IPTV blijft zolang Xtream een actieve connectie ziet; events
en custom live verdwijnen ~75s zonder heartbeat.
</p> </p>
<div className="desktop-table card"> <div className="desktop-table card">
@ -101,7 +119,8 @@ export default function StreamsPage() {
</thead> </thead>
<tbody> <tbody>
{streams.map((s) => { {streams.map((s) => {
const isIptv = s.source === "iptv" || s.node === "IPTV"; const live = isLivePresence(s);
const canRevoke = s.source === "library" || s.source === "client" || s.id.startsWith("cw-");
const pct = s.progressPercent ?? 0; const pct = s.progressPercent ?? 0;
return ( return (
<tr key={s.id}> <tr key={s.id}>
@ -110,12 +129,10 @@ export default function StreamsPage() {
</td> </td>
<td> <td>
<div>{s.title}</div> <div>{s.title}</div>
<div className="muted tiny"> <div className="muted tiny">{presenceHint(s)}</div>
{isIptv ? "Direct naar Xtream" : `${formatBytes(s.bytesSent)} gestreamd`}
</div>
</td> </td>
<td style={{ minWidth: 140 }}> <td style={{ minWidth: 140 }}>
{isIptv ? ( {live ? (
<div className="muted mono tiny">live</div> <div className="muted mono tiny">live</div>
) : ( ) : (
<> <>
@ -132,9 +149,9 @@ export default function StreamsPage() {
<td>{s.resolution ?? "-"}</td> <td>{s.resolution ?? "-"}</td>
<td className="mono tiny">{relativeTime(s.lastActivity)}</td> <td className="mono tiny">{relativeTime(s.lastActivity)}</td>
<td> <td>
{!isIptv && ( {canRevoke && !(s.source === "iptv" || s.node === "IPTV") && (
<button type="button" className="danger" onClick={() => revoke(s.id)}> <button type="button" className="danger" onClick={() => revoke(s.id)}>
Intrekken {s.source === "client" || s.id.startsWith("cw-") ? "Wissen" : "Intrekken"}
</button> </button>
)} )}
</td> </td>
@ -154,7 +171,8 @@ export default function StreamsPage() {
<div className="list-cards"> <div className="list-cards">
{streams.map((s) => { {streams.map((s) => {
const isIptv = s.source === "iptv" || s.node === "IPTV"; const live = isLivePresence(s);
const canRevoke = s.source === "library" || s.source === "client" || s.id.startsWith("cw-");
const pct = s.progressPercent ?? 0; const pct = s.progressPercent ?? 0;
return ( return (
<article key={s.id} className="list-card stream-card"> <article key={s.id} className="list-card stream-card">
@ -163,25 +181,25 @@ export default function StreamsPage() {
<span className="mono tiny">{relativeTime(s.lastActivity)}</span> <span className="mono tiny">{relativeTime(s.lastActivity)}</span>
</div> </div>
<div className="stream-title">{s.title}</div> <div className="stream-title">{s.title}</div>
{!isIptv && ( {!live && (
<div className="progress-track" title={`${pct.toFixed(0)}%`}> <div className="progress-track" title={`${pct.toFixed(0)}%`}>
<div className="progress-fill" style={{ width: `${pct}%` }} /> <div className="progress-fill" style={{ width: `${pct}%` }} />
</div> </div>
)} )}
<div className="list-card-meta"> <div className="list-card-meta">
<span className="mono"> <span className="mono">
{isIptv {live
? "live · Xtream" ? presenceHint(s)
: `${s.progressPercent != null ? `${pct.toFixed(0)}%` : "—"} · ${formatBytes(s.bytesSent)}`} : `${s.progressPercent != null ? `${pct.toFixed(0)}%` : "—"} · ${formatBytes(s.bytesSent)}`}
</span> </span>
<span> <span>
{s.resolution ?? "?"} · {s.node} {s.resolution ?? "?"} · {s.node}
</span> </span>
</div> </div>
{!isIptv && ( {canRevoke && !(s.source === "iptv" || s.node === "IPTV") && (
<div className="action-row"> <div className="action-row">
<button type="button" className="danger" onClick={() => revoke(s.id)}> <button type="button" className="danger" onClick={() => revoke(s.id)}>
Intrekken {s.source === "client" || s.id.startsWith("cw-") ? "Wissen" : "Intrekken"}
</button> </button>
</div> </div>
)} )}

View file

@ -556,10 +556,30 @@ class ApiClient(
fallbackStreamUrl = o.optString("fallbackStreamUrl").ifBlank { null }, fallbackStreamUrl = o.optString("fallbackStreamUrl").ifBlank { null },
fallbackFormat = o.optString("fallbackFormat").ifBlank { null }, fallbackFormat = o.optString("fallbackFormat").ifBlank { null },
drm = drm, drm = drm,
watchSessionId = o.optString("watchSessionId").ifBlank { null },
) )
} }
} }
/**
* Presence-heartbeat voor custom live-lists (best-effort).
*/
suspend fun watchingHeartbeat(
watchSessionId: String? = null,
kind: String? = null,
id: String? = null,
) = withContext(Dispatchers.IO) {
val body = JSONObject()
if (!watchSessionId.isNullOrBlank()) body.put("watchSessionId", watchSessionId)
if (!kind.isNullOrBlank()) body.put("kind", kind)
if (!id.isNullOrBlank()) body.put("id", id)
val req = authed("$baseUrl/api/v1/client/watching/heartbeat")
.post(body.toString().toRequestBody(jsonMedia))
.build()
runCatching { execute(req) }
Unit
}
suspend fun iptvEpg(streamId: String): IptvEpg = withContext(Dispatchers.IO) { suspend fun iptvEpg(streamId: String): IptvEpg = withContext(Dispatchers.IO) {
val enc = android.net.Uri.encode(streamId) val enc = android.net.Uri.encode(streamId)
val req = authed("$baseUrl/api/v1/client/iptv/channels/$enc/epg").get().build() val req = authed("$baseUrl/api/v1/client/iptv/channels/$enc/epg").get().build()
@ -932,6 +952,7 @@ data class IptvPlayInfo(
val fallbackStreamUrl: String? = null, val fallbackStreamUrl: String? = null,
val fallbackFormat: String? = null, val fallbackFormat: String? = null,
val drm: DrmInfo? = null, val drm: DrmInfo? = null,
val watchSessionId: String? = null,
) )
data class EpgProgram( data class EpgProgram(

View file

@ -119,6 +119,7 @@ fun LiveTvPlayerScreen(
var fallbackUrl by remember { mutableStateOf<String?>(null) } var fallbackUrl by remember { mutableStateOf<String?>(null) }
var fallbackFormat by remember { mutableStateOf<String?>(null) } var fallbackFormat by remember { mutableStateOf<String?>(null) }
var triedFallback by remember { mutableStateOf(false) } var triedFallback by remember { mutableStateOf(false) }
var watchSessionId by remember { mutableStateOf<String?>(null) }
var menuChannels by remember { mutableStateOf<List<IptvChannel>>(emptyList()) } var menuChannels by remember { mutableStateOf<List<IptvChannel>>(emptyList()) }
var menuCategories by remember { mutableStateOf<List<IptvCategory>>(emptyList()) } var menuCategories by remember { mutableStateOf<List<IptvCategory>>(emptyList()) }
@ -142,11 +143,13 @@ fun LiveTvPlayerScreen(
if (isSwitch) switching = true else loading = true if (isSwitch) switching = true else loading = true
error = null error = null
triedFallback = false triedFallback = false
watchSessionId = null
try { try {
val play = api.iptvPlayUrl(streamIdToPlay) val play = api.iptvPlayUrl(streamIdToPlay)
channelName = play.name channelName = play.name
logoUrl = play.logoUrl logoUrl = play.logoUrl
streamUrl = play.streamUrl streamUrl = play.streamUrl
watchSessionId = play.watchSessionId
val isDash = play.format.equals("dash", ignoreCase = true) val isDash = play.format.equals("dash", ignoreCase = true)
val tsUrl = when { val tsUrl = when {
isDash -> null isDash -> null
@ -178,6 +181,7 @@ fun LiveTvPlayerScreen(
onStreamChange(streamIdToPlay) onStreamChange(streamIdToPlay)
} catch (e: Exception) { } catch (e: Exception) {
error = e.message ?: "Kon zender niet starten" error = e.message ?: "Kon zender niet starten"
watchSessionId = null
} finally { } finally {
loading = false loading = false
switching = false switching = false
@ -197,6 +201,22 @@ fun LiveTvPlayerScreen(
loadAndPlay(currentStreamId, isSwitch = currentStreamId != streamId) loadAndPlay(currentStreamId, isSwitch = currentStreamId != streamId)
} }
// Presence heartbeat voor custom live-lists (watchSessionId uit play).
LaunchedEffect(watchSessionId, currentStreamId, error, loading) {
val session = watchSessionId ?: return@LaunchedEffect
if (error != null || loading) return@LaunchedEffect
while (true) {
runCatching {
api.watchingHeartbeat(
watchSessionId = session,
kind = "live-list",
id = currentStreamId,
)
}
delay(25_000)
}
}
LaunchedEffect(currentStreamId) { LaunchedEffect(currentStreamId) {
delay(2_500) delay(2_500)
runCatching { epg = api.iptvEpg(currentStreamId) } runCatching { epg = api.iptvEpg(currentStreamId) }

View file

@ -20,8 +20,8 @@ android {
applicationId = "nl.vonas.mediacluster.tv" applicationId = "nl.vonas.mediacluster.tv"
minSdk = 24 minSdk = 24
targetSdk = 36 targetSdk = 36
versionCode = 118 versionCode = 119
versionName = "0.14.55" versionName = "0.14.56"
buildConfigField("String", "DEFAULT_API_BASE", "\"https://master.vonas.nl\"") buildConfigField("String", "DEFAULT_API_BASE", "\"https://master.vonas.nl\"")
ndk { ndk {
abiFilters += listOf("arm64-v8a", "armeabi-v7a") abiFilters += listOf("arm64-v8a", "armeabi-v7a")

View file

@ -570,6 +570,7 @@ class ApiClient(
fallbackStreamUrl = o.optString("fallbackStreamUrl").ifBlank { null }, fallbackStreamUrl = o.optString("fallbackStreamUrl").ifBlank { null },
fallbackFormat = o.optString("fallbackFormat").ifBlank { null }, fallbackFormat = o.optString("fallbackFormat").ifBlank { null },
drm = drm, drm = drm,
watchSessionId = o.optString("watchSessionId").ifBlank { null },
) )
} }
} }
@ -661,10 +662,31 @@ class ApiClient(
fallbackStreamUrl = o.optString("fallbackStreamUrl").ifBlank { null }, fallbackStreamUrl = o.optString("fallbackStreamUrl").ifBlank { null },
fallbackFormat = o.optString("fallbackFormat").ifBlank { null }, fallbackFormat = o.optString("fallbackFormat").ifBlank { null },
drm = drm, drm = drm,
watchSessionId = o.optString("watchSessionId").ifBlank { null },
) )
} }
} }
/**
* Presence-heartbeat voor events / custom live-lists.
* Best-effort: never throws to callers that ignore failures.
*/
suspend fun watchingHeartbeat(
watchSessionId: String? = null,
kind: String? = null,
id: String? = null,
) = withContext(Dispatchers.IO) {
val body = JSONObject()
if (!watchSessionId.isNullOrBlank()) body.put("watchSessionId", watchSessionId)
if (!kind.isNullOrBlank()) body.put("kind", kind)
if (!id.isNullOrBlank()) body.put("id", id)
val req = authed("$baseUrl/api/v1/client/watching/heartbeat")
.post(body.toString().toRequestBody(jsonMedia))
.build()
runCatching { execute(req) }
Unit
}
private fun JSONObject.toScheduleEvent(): ScheduleEvent? { private fun JSONObject.toScheduleEvent(): ScheduleEvent? {
val id = optClean("id") ?: return null val id = optClean("id") ?: return null
val name = optClean("name") ?: return null val name = optClean("name") ?: return null
@ -1128,6 +1150,7 @@ data class IptvPlayInfo(
val fallbackStreamUrl: String? = null, val fallbackStreamUrl: String? = null,
val fallbackFormat: String? = null, val fallbackFormat: String? = null,
val drm: DrmInfo? = null, val drm: DrmInfo? = null,
val watchSessionId: String? = null,
) )
data class EpgProgram( data class EpgProgram(

View file

@ -70,6 +70,7 @@ fun EventPlayerScreen(
var loading by remember { mutableStateOf(true) } var loading by remember { mutableStateOf(true) }
var error by remember { mutableStateOf<String?>(null) } var error by remember { mutableStateOf<String?>(null) }
var showChrome by remember { mutableStateOf(true) } var showChrome by remember { mutableStateOf(true) }
var watchSessionId by remember { mutableStateOf<String?>(null) }
fun safeStopAndClear() { fun safeStopAndClear() {
runCatching { runCatching {
@ -82,12 +83,14 @@ fun EventPlayerScreen(
LaunchedEffect(eventId) { LaunchedEffect(eventId) {
loading = true loading = true
error = null error = null
watchSessionId = null
reachedReady.set(false) reachedReady.set(false)
showChrome = true showChrome = true
try { try {
val play = api.scheduleEventPlay(eventId) val play = api.scheduleEventPlay(eventId)
if (!isActive) return@LaunchedEffect if (!isActive) return@LaunchedEffect
title = play.name title = play.name
watchSessionId = play.watchSessionId
val url = play.streamUrl.trim() val url = play.streamUrl.trim()
if (url.isEmpty()) { if (url.isEmpty()) {
error = "Geen stream-URL beschikbaar" error = "Geen stream-URL beschikbaar"
@ -122,6 +125,22 @@ fun EventPlayerScreen(
} }
} }
// Presence heartbeat — fire-and-forget, geen UI/focus-impact.
LaunchedEffect(watchSessionId, error, loading) {
val session = watchSessionId ?: return@LaunchedEffect
if (error != null || loading) return@LaunchedEffect
while (isActive) {
runCatching {
api.watchingHeartbeat(
watchSessionId = session,
kind = "event",
id = eventId,
)
}
delay(25_000)
}
}
// Geen eindeloze zwart-scherm hang: toon fout als READY uitblijft. // Geen eindeloze zwart-scherm hang: toon fout als READY uitblijft.
LaunchedEffect(eventId, error, loading) { LaunchedEffect(eventId, error, loading) {
if (error != null || loading) return@LaunchedEffect if (error != null || loading) return@LaunchedEffect

View file

@ -114,6 +114,7 @@ fun LiveTvPlayerScreen(
var fallbackUrl by remember { mutableStateOf<String?>(null) } var fallbackUrl by remember { mutableStateOf<String?>(null) }
var fallbackFormat by remember { mutableStateOf<String?>(null) } var fallbackFormat by remember { mutableStateOf<String?>(null) }
var triedFallback by remember { mutableStateOf(false) } var triedFallback by remember { mutableStateOf(false) }
var watchSessionId by remember { mutableStateOf<String?>(null) }
var menuChannels by remember { mutableStateOf<List<IptvChannel>>(emptyList()) } var menuChannels by remember { mutableStateOf<List<IptvChannel>>(emptyList()) }
var menuCategories by remember { mutableStateOf<List<IptvCategory>>(emptyList()) } var menuCategories by remember { mutableStateOf<List<IptvCategory>>(emptyList()) }
@ -148,10 +149,12 @@ fun LiveTvPlayerScreen(
if (isSwitch) switching = true else loading = true if (isSwitch) switching = true else loading = true
error = null error = null
triedFallback = false triedFallback = false
watchSessionId = null
try { try {
val play = api.iptvPlayUrl(streamIdToPlay) val play = api.iptvPlayUrl(streamIdToPlay)
channelName = play.name channelName = play.name
logoUrl = play.logoUrl logoUrl = play.logoUrl
watchSessionId = play.watchSessionId
val isDash = play.format.equals("dash", ignoreCase = true) val isDash = play.format.equals("dash", ignoreCase = true)
val tsUrl = when { val tsUrl = when {
isDash -> null isDash -> null
@ -187,6 +190,7 @@ fun LiveTvPlayerScreen(
onStreamChange(streamIdToPlay) onStreamChange(streamIdToPlay)
} catch (e: Exception) { } catch (e: Exception) {
error = e.message ?: "Kon zender niet starten" error = e.message ?: "Kon zender niet starten"
watchSessionId = null
} finally { } finally {
loading = false loading = false
switching = false switching = false
@ -206,6 +210,22 @@ fun LiveTvPlayerScreen(
loadAndPlay(currentStreamId, isSwitch = currentStreamId != streamId) loadAndPlay(currentStreamId, isSwitch = currentStreamId != streamId)
} }
// Presence heartbeat voor custom live-lists (watchSessionId uit play).
LaunchedEffect(watchSessionId, currentStreamId, error, loading) {
val session = watchSessionId ?: return@LaunchedEffect
if (error != null || loading) return@LaunchedEffect
while (true) {
runCatching {
api.watchingHeartbeat(
watchSessionId = session,
kind = "live-list",
id = currentStreamId,
)
}
delay(25_000)
}
}
LaunchedEffect(currentStreamId) { LaunchedEffect(currentStreamId) {
delay(2_500) delay(2_500)
runCatching { epg = api.iptvEpg(currentStreamId) } runCatching { epg = api.iptvEpg(currentStreamId) }

View file

@ -11,6 +11,11 @@ import { nodeConnectionManager } from "../websocket/manager";
import { createMetadataService } from "../settings/metadata-factory"; import { createMetadataService } from "../settings/metadata-factory";
import type { Config } from "../config"; import type { Config } from "../config";
import { listActiveIptvWatches } from "../iptv/active-watches"; import { listActiveIptvWatches } from "../iptv/active-watches";
import {
isClientWatchId,
listActiveClientWatches,
revokeClientWatch,
} from "../viewer/client-watches";
import { import {
parseMovieFromPath, parseMovieFromPath,
parseSeriesFromPath, parseSeriesFromPath,
@ -56,7 +61,8 @@ export async function registerAdminRoutes(app: FastifyInstance, config: Config)
}, },
}); });
const iptv = await listActiveIptvWatches(config); const iptv = await listActiveIptvWatches(config);
return library + iptv.length; const client = listActiveClientWatches();
return library + iptv.length + client.length;
})(), })(),
prisma.node.findMany({ prisma.node.findMany({
where: { revoked: false, status: "ONLINE" }, where: { revoked: false, status: "ONLINE" },
@ -496,12 +502,14 @@ export async function registerAdminRoutes(app: FastifyInstance, config: Config)
}); });
const iptvWatches = await listActiveIptvWatches(config); const iptvWatches = await listActiveIptvWatches(config);
const clientWatches = listActiveClientWatches();
// App-clients zetten ViewerUser.id in userId (zonder addonToken). // App-clients zetten ViewerUser.id in userId (zonder addonToken).
const viewerIds = [ const viewerIds = [
...new Set([ ...new Set([
...sessions.map((s) => s.userId).filter((id): id is string => !!id), ...sessions.map((s) => s.userId).filter((id): id is string => !!id),
...iptvWatches.map((w) => w.viewerUserId).filter((id): id is string => !!id), ...iptvWatches.map((w) => w.viewerUserId).filter((id): id is string => !!id),
...clientWatches.map((w) => w.viewerUserId),
]), ]),
]; ];
const viewerUsers = const viewerUsers =
@ -559,6 +567,12 @@ export async function registerAdminRoutes(app: FastifyInstance, config: Config)
return "Onbekend"; return "Onbekend";
} }
function clientViewerLabel(viewerUserId: string): string {
const v = viewerById.get(viewerUserId);
if (v) return v.name?.trim() || v.email;
return "Onbekend";
}
// Eén rij per kijker+titel (hoogste bytes wint) — niet per mediabestand. // Eén rij per kijker+titel (hoogste bytes wint) — niet per mediabestand.
const best = new Map<string, (typeof sessions)[number]>(); const best = new Map<string, (typeof sessions)[number]>();
for (const s of sessions) { for (const s of sessions) {
@ -590,6 +604,7 @@ export async function registerAdminRoutes(app: FastifyInstance, config: Config)
return { return {
id: s.id, id: s.id,
source: "library" as const, source: "library" as const,
kind: "library" as const,
node: s.node.name, node: s.node.name,
status: s.status, status: s.status,
viewer: streamViewerLabel(s), viewer: streamViewerLabel(s),
@ -611,6 +626,7 @@ export async function registerAdminRoutes(app: FastifyInstance, config: Config)
return { return {
id: w.id, id: w.id,
source: "iptv" as const, source: "iptv" as const,
kind: w.kind,
node: "IPTV", node: "IPTV",
status: "ACTIVE", status: "ACTIVE",
viewer: iptvViewerLabel(w), viewer: iptvViewerLabel(w),
@ -627,8 +643,32 @@ export async function registerAdminRoutes(app: FastifyInstance, config: Config)
}; };
}); });
const clientStreams = clientWatches.map((w) => {
const isEvent = w.kind === "event";
const titlePrefix = isEvent ? "Event" : "Live";
const node = w.provider?.trim() || (isEvent ? "Event" : "Live");
return {
id: w.id,
source: "client" as const,
kind: w.kind,
node,
status: "ACTIVE",
viewer: clientViewerLabel(w.viewerUserId),
addonTokenId: null as string | null,
createdAt: w.startedAt,
lastActivity: w.lastActivity,
absoluteExpiresAt: new Date(w.lastActivity.getTime() + 90_000),
title: `${titlePrefix} · ${w.title}`,
resolution: null as string | null,
bytesSent: "0",
lastByteOffset: "0",
fileSizeBytes: "0",
progressPercent: null as number | null,
};
});
return { return {
streams: [...libraryStreams, ...iptvStreams].sort( streams: [...libraryStreams, ...iptvStreams, ...clientStreams].sort(
(a, b) => new Date(b.lastActivity).getTime() - new Date(a.lastActivity).getTime() (a, b) => new Date(b.lastActivity).getTime() - new Date(a.lastActivity).getTime()
), ),
}; };
@ -636,6 +676,17 @@ export async function registerAdminRoutes(app: FastifyInstance, config: Config)
app.post("/api/v1/admin/streams/:id/revoke", { preHandler: requireAdmin }, async (request) => { app.post("/api/v1/admin/streams/:id/revoke", { preHandler: requireAdmin }, async (request) => {
const { id } = request.params as { id: string }; const { id } = request.params as { id: string };
if (isClientWatchId(id)) {
// Direct play (geen proxy): presence wissen stopt de stream op het apparaat niet.
const cleared = revokeClientWatch(id);
if (!cleared) throw new AppError("NOT_FOUND", "Session not found", 404);
return {
ok: true,
presenceCleared: true,
streamStopped: false,
note: "Presence gewist; de stream op het apparaat stopt niet (direct play).",
};
}
if (id.startsWith("iptv-")) { if (id.startsWith("iptv-")) {
throw new AppError( throw new AppError(
"INVALID_REQUEST", "INVALID_REQUEST",

View file

@ -0,0 +1,117 @@
import { randomUUID } from "crypto";
/** Direct-play kinds zonder Xtream active_cons (events + custom live-lists). */
export type ClientWatchKind = "event" | "live-list";
export type ClientActiveWatch = {
id: string;
viewerUserId: string;
kind: ClientWatchKind;
resourceId: string;
title: string;
provider: string | null;
startedAt: Date;
lastActivity: Date;
};
/** Zonder heartbeat verdwijnt de entry (~75s, midden in 60–90). */
const TTL_MS = 75_000;
/** Key: viewer + kind + resource — één actieve kijk per zender/event per viewer. */
const byViewerKey = new Map<string, ClientActiveWatch>();
const bySessionId = new Map<string, ClientActiveWatch>();
function viewerKey(viewerUserId: string, kind: ClientWatchKind, resourceId: string): string {
return `${viewerUserId}:${kind}:${resourceId}`;
}
function purgeExpired(now = Date.now()): void {
for (const [key, w] of [...byViewerKey.entries()]) {
if (now - w.lastActivity.getTime() > TTL_MS) {
byViewerKey.delete(key);
bySessionId.delete(w.id);
}
}
}
export function touchClientWatch(input: {
viewerUserId: string;
kind: ClientWatchKind;
resourceId: string;
title: string;
provider?: string | null;
}): ClientActiveWatch {
purgeExpired();
const now = new Date();
const key = viewerKey(input.viewerUserId, input.kind, input.resourceId);
const existing = byViewerKey.get(key);
const watch: ClientActiveWatch = {
id: existing?.id ?? `cw-${randomUUID()}`,
viewerUserId: input.viewerUserId,
kind: input.kind,
resourceId: input.resourceId,
title: input.title.trim() || existing?.title || input.resourceId,
provider: input.provider?.trim() || existing?.provider || null,
startedAt: existing?.startedAt ?? now,
lastActivity: now,
};
byViewerKey.set(key, watch);
bySessionId.set(watch.id, watch);
return watch;
}
/**
* Heartbeat: verlengt TTL.
* Accepteert watchSessionId óf kind+id (resourceId) voor dezelfde viewer.
*/
export function heartbeatClientWatch(opts: {
viewerUserId: string;
watchSessionId?: string | null;
kind?: ClientWatchKind | null;
resourceId?: string | null;
}): ClientActiveWatch | null {
purgeExpired();
const now = new Date();
let watch: ClientActiveWatch | undefined;
const sessionId = opts.watchSessionId?.trim();
if (sessionId) {
watch = bySessionId.get(sessionId);
if (watch && watch.viewerUserId !== opts.viewerUserId) return null;
}
if (!watch && opts.kind && opts.resourceId?.trim()) {
watch = byViewerKey.get(viewerKey(opts.viewerUserId, opts.kind, opts.resourceId.trim()));
}
if (!watch) return null;
watch.lastActivity = now;
return watch;
}
export function listActiveClientWatches(): ClientActiveWatch[] {
purgeExpired();
return [...byViewerKey.values()].sort(
(a, b) => b.lastActivity.getTime() - a.lastActivity.getTime()
);
}
export function countActiveClientWatches(): number {
purgeExpired();
return byViewerKey.size;
}
/**
* Presence clear only — stopt de player-stream niet (direct play, geen proxy).
* Returns true if an entry was removed.
*/
export function revokeClientWatch(sessionId: string): boolean {
const w = bySessionId.get(sessionId);
if (!w) return false;
bySessionId.delete(sessionId);
byViewerKey.delete(viewerKey(w.viewerUserId, w.kind, w.resourceId));
return true;
}
export function isClientWatchId(id: string): boolean {
return id.startsWith("cw-");
}

View file

@ -416,6 +416,7 @@ export async function getCustomPlayUrl(
viewers: { some: { viewerUserId: viewerId } }, viewers: { some: { viewerUserId: viewerId } },
}, },
}, },
include: { list: { select: { name: true } } },
}); });
if (!channel) throw new AppError("FORBIDDEN", "Geen toegang tot deze zender", 403); if (!channel) throw new AppError("FORBIDDEN", "Geen toegang tot deze zender", 403);
@ -432,8 +433,24 @@ export async function getCustomPlayUrl(
} }
: null; : null;
const sid = customStreamId(channel.id);
let watchSessionId: string | null = null;
try {
const { touchClientWatch } = await import("./client-watches");
const watch = touchClientWatch({
viewerUserId: viewerId,
kind: "live-list",
resourceId: sid,
title: channel.name,
provider: channel.list.name,
});
watchSessionId = watch.id;
} catch {
// Presence is best-effort; play must not fail.
}
return { return {
streamId: customStreamId(channel.id), streamId: sid,
name: channel.name, name: channel.name,
logoUrl: channel.logoUrl, logoUrl: channel.logoUrl,
streamUrl: channel.mpdUrl, streamUrl: channel.mpdUrl,
@ -441,6 +458,7 @@ export async function getCustomPlayUrl(
fallbackStreamUrl: null, fallbackStreamUrl: null,
fallbackFormat: null, fallbackFormat: null,
drm, drm,
watchSessionId,
}; };
} }

View file

@ -871,7 +871,35 @@ export function registerViewerRoutes(
if (!id?.trim()) throw new AppError("INVALID_REQUEST", "Ongeldig event-id", 400); if (!id?.trim()) throw new AppError("INVALID_REQUEST", "Ongeldig event-id", 400);
const { getScheduleEventPlay, resolveViewerLiveEventsAccess } = await import("./schedule-events"); const { getScheduleEventPlay, resolveViewerLiveEventsAccess } = await import("./schedule-events");
const access = await resolveViewerLiveEventsAccess(auth.viewerId, config); const access = await resolveViewerLiveEventsAccess(auth.viewerId, config);
return getScheduleEventPlay(config, decodeURIComponent(id), access); return getScheduleEventPlay(config, decodeURIComponent(id), access, auth.viewerId);
});
/**
* Lichte presence-heartbeat voor Live Events / custom live-lists.
* Fire-and-forget vanaf de client; failures mogen genegeerd worden.
*/
app.post("/api/v1/client/watching/heartbeat", async (request) => {
const auth = await viewers.authFromBearer(request.headers.authorization);
const body = (request.body ?? {}) as {
watchSessionId?: string;
kind?: string;
id?: string;
};
const { heartbeatClientWatch } = await import("./client-watches");
const kindRaw = body.kind?.trim().toLowerCase();
const kind =
kindRaw === "event"
? ("event" as const)
: kindRaw === "live-list" || kindRaw === "live"
? ("live-list" as const)
: null;
const watch = heartbeatClientWatch({
viewerUserId: auth.viewerId,
watchSessionId: body.watchSessionId,
kind,
resourceId: body.id,
});
return { ok: !!watch };
}); });
app.get("/api/v1/client/subtitles/search", async (request) => { app.get("/api/v1/client/subtitles/search", async (request) => {

View file

@ -575,7 +575,8 @@ export async function getScheduleEvent(
export async function getScheduleEventPlay( export async function getScheduleEventPlay(
config: Config, config: Config,
eventId: string, eventId: string,
access?: LiveEventsAccessFilter | null access?: LiveEventsAccessFilter | null,
viewerUserId?: string | null
) { ) {
assertAccessEnabled(access); assertAccessEnabled(access);
const id = str(eventId); const id = str(eventId);
@ -609,6 +610,23 @@ export async function getScheduleEventPlay(
); );
} }
let watchSessionId: string | null = null;
if (viewerUserId) {
try {
const { touchClientWatch } = await import("./client-watches");
const watch = touchClientWatch({
viewerUserId,
kind: "event",
resourceId: pub.id,
title: pub.name,
provider: pub.provider,
});
watchSessionId = watch.id;
} catch {
// Presence is best-effort; play must not fail.
}
}
return { return {
eventId: pub.id, eventId: pub.id,
name: pub.name, name: pub.name,
@ -621,5 +639,6 @@ export async function getScheduleEventPlay(
type: "clearkey" as const, type: "clearkey" as const,
keys, keys,
}, },
watchSessionId,
}; };
} }