diff --git a/apps/admin-ui/src/app/streams/page.tsx b/apps/admin-ui/src/app/streams/page.tsx index f42e17e..78afead 100644 --- a/apps/admin-ui/src/app/streams/page.tsx +++ b/apps/admin-ui/src/app/streams/page.tsx @@ -84,8 +84,8 @@ export default function StreamsPage() {

- Alleen echte kijkers (bytes stromen). Stopt iemand → verdwijnt binnen ~3 minuten uit deze lijst. - Voortgang is een schatting op basis van range-positie in het bestand. + Eén rij per kijker + titel. Alleen echte streams (>512 KB). Stopt iemand → verdwijnt + binnen ~3 minuten. Voortgang is een schatting op range-positie.

diff --git a/apps/master-api/src/admin/routes.ts b/apps/master-api/src/admin/routes.ts index 954bf45..2c9ec92 100644 --- a/apps/master-api/src/admin/routes.ts +++ b/apps/master-api/src/admin/routes.ts @@ -43,7 +43,7 @@ export async function registerAdminRoutes(app: FastifyInstance, config: Config) revoked: false, absoluteExpiresAt: { gt: new Date() }, lastActivity: { gt: new Date(Date.now() - 3 * 60 * 1000) }, - bytesSent: { gt: 0 }, + bytesSent: { gt: 512n * 1024n }, }, }), prisma.node.findMany({ @@ -253,15 +253,16 @@ export async function registerAdminRoutes(app: FastifyInstance, config: Config) data: { status: "EXPIRED" }, }); - // Only recently active + actually streamed bytes (no prefetch ghosts). + // Alleen echte kijkers: >512KB gestreamd + recent actief. const liveSince = new Date(now.getTime() - 3 * 60 * 1000); + const minBytes = 512n * 1024n; const sessions = await prisma.playbackSession.findMany({ where: { status: "ACTIVE", revoked: false, absoluteExpiresAt: { gt: now }, lastActivity: { gt: liveSince }, - bytesSent: { gt: 0 }, + bytesSent: { gt: minBytes }, }, include: { node: { select: { id: true, name: true } }, @@ -279,8 +280,23 @@ export async function registerAdminRoutes(app: FastifyInstance, config: Config) take: 100, }); + // Eén rij per kijker+titel (hoogste bytes wint) — niet per mediabestand. + const best = new Map(); + for (const s of sessions) { + const titleKey = + s.mediaFile.movieId ?? + (s.mediaFile.episode + ? `ep:${s.mediaFile.episode.seriesId ?? ""}:${s.mediaFile.episode.seasonNumber}:${s.mediaFile.episode.episodeNumber}` + : s.mediaFileId); + const key = `${s.addonTokenId ?? s.id}:${titleKey}`; + const prev = best.get(key); + if (!prev || s.bytesSent > prev.bytesSent) best.set(key, s); + } + return { - streams: sessions.map((s) => { + streams: [...best.values()] + .sort((a, b) => b.lastActivity.getTime() - a.lastActivity.getTime()) + .map((s) => { const fileSize = Number(s.fileSizeBytes || s.mediaFile.sizeBytes || 0); const offset = Number(s.lastByteOffset || 0); const progressPercent = diff --git a/apps/master-api/src/playback/service.ts b/apps/master-api/src/playback/service.ts index 775f317..795f514 100644 --- a/apps/master-api/src/playback/service.ts +++ b/apps/master-api/src/playback/service.ts @@ -8,10 +8,42 @@ import { AppError } from "../security/errors"; import { nodeConnectionManager } from "../websocket/manager"; import type { Config } from "../config"; +type SessionResult = { + sessionId: string; + streamUrl: string; + token: string; + nodeId: string; + mediaFileId: string; +}; + +/** In-memory play tokens so concurrent Stremio probes reuse one session. */ +const playTokens = new Map(); +const createLocks = new Map>(); + +function sessionKey(addonTokenId: string | undefined, mediaFileId: string): string { + return `${addonTokenId || "anon"}:${mediaFileId}`; +} + export class PlaybackService { constructor(private readonly config: Config) {} async createSession(mediaFileId: string, userId?: string, addonTokenId?: string) { + const key = sessionKey(addonTokenId, mediaFileId); + const existingLock = createLocks.get(key); + if (existingLock) return existingLock; + + const work = this.createSessionUnlocked(mediaFileId, userId, addonTokenId).finally(() => { + createLocks.delete(key); + }); + createLocks.set(key, work); + return work; + } + + private async createSessionUnlocked( + mediaFileId: string, + userId?: string, + addonTokenId?: string + ): Promise { const mediaFile = await prisma.mediaFile.findUnique({ where: { id: mediaFileId }, include: { node: true }, @@ -29,10 +61,65 @@ export class PlaybackService { throw new AppError("NODE_OFFLINE", "Selected media node is unavailable", 503); } + const key = sessionKey(addonTokenId, mediaFileId); + const now = new Date(); + + // Reuse one live session per kijker + bestand (Stremio opent vaak meerdere Range/HEAD hits). + if (addonTokenId) { + const active = await prisma.playbackSession.findFirst({ + where: { + addonTokenId, + mediaFileId: mediaFile.id, + status: "ACTIVE", + revoked: false, + absoluteExpiresAt: { gt: now }, + }, + orderBy: { lastActivity: "desc" }, + }); + const cached = playTokens.get(key); + if ( + active && + cached && + cached.sessionId === active.id && + hashToken(cached.token) === active.tokenHash + ) { + await prisma.playbackSession.update({ + where: { id: active.id }, + data: { lastActivity: now }, + }); + const streamUrl = `${mediaFile.node.publicStreamUrl.replace(/\/$/, "")}/play/${cached.token}`; + return { + sessionId: active.id, + streamUrl, + token: cached.token, + nodeId: mediaFile.nodeId, + mediaFileId: mediaFile.id, + }; + } + + // Oude parallelle sessies opruimen voordat we een nieuwe maken. + const stale = await prisma.playbackSession.findMany({ + where: { + addonTokenId, + mediaFileId: mediaFile.id, + status: "ACTIVE", + revoked: false, + }, + select: { id: true, nodeId: true }, + }); + for (const s of stale) { + await prisma.playbackSession.update({ + where: { id: s.id }, + data: { status: "COMPLETED" }, + }); + nodeConnectionManager.revokePlaybackSession(s.nodeId, s.id); + } + playTokens.delete(key); + } + const token = generateSecureToken(32); const sessionId = generateSessionId(); const tokenHash = hashToken(token); - const now = new Date(); const absoluteExpiresAt = new Date( now.getTime() + this.config.PLAYBACK_ABSOLUTE_LIFETIME_SECONDS * 1000 ); @@ -69,6 +156,7 @@ export class PlaybackService { throw new AppError("NODE_OFFLINE", "Failed to register session on node", 503); } + playTokens.set(key, { token, sessionId }); const streamUrl = `${mediaFile.node.publicStreamUrl.replace(/\/$/, "")}/play/${token}`; return { @@ -97,6 +185,10 @@ export class PlaybackService { data: { status: "REVOKED", revoked: true }, }); + for (const [key, val] of playTokens) { + if (val.sessionId === sessionId) playTokens.delete(key); + } + nodeConnectionManager.revokePlaybackSession(session.nodeId, sessionId); } diff --git a/deploy/docker/Dockerfile.master-api b/deploy/docker/Dockerfile.master-api index 397ade4..dbf6142 100644 --- a/deploy/docker/Dockerfile.master-api +++ b/deploy/docker/Dockerfile.master-api @@ -7,8 +7,8 @@ WORKDIR /src RUN apk add --no-cache git ca-certificates COPY node/media-node/ ./ RUN go mod tidy \ - && CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -ldflags="-s -w -X main.version=1.2.8" -o /media-node-linux-amd64 ./cmd/media-node \ - && CGO_ENABLED=0 GOOS=linux GOARCH=arm64 go build -ldflags="-s -w -X main.version=1.2.8" -o /media-node-linux-arm64 ./cmd/media-node + && CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -ldflags="-s -w -X main.version=1.2.9" -o /media-node-linux-amd64 ./cmd/media-node \ + && CGO_ENABLED=0 GOOS=linux GOARCH=arm64 go build -ldflags="-s -w -X main.version=1.2.9" -o /media-node-linux-arm64 ./cmd/media-node FROM base AS builder COPY package.json pnpm-workspace.yaml ./ diff --git a/node/media-node/internal/streaming/server.go b/node/media-node/internal/streaming/server.go index df0c810..46ae367 100644 --- a/node/media-node/internal/streaming/server.go +++ b/node/media-node/internal/streaming/server.go @@ -295,17 +295,17 @@ func (s *Server) noteProgress(sessionID string, offset, fileSize, deltaBytes int } if deltaBytes > 0 { st.bytesSent += deltaBytes - } else if st.bytesSent == 0 && force { - // Mark as live even before first chunk (Range open). - st.bytesSent = 1 - } - if offset >= st.lastOffset { - st.lastOffset = offset + // Alleen offset bijwerken bij echte data — voorkomt 100% door 1-byte end-of-file probes. + if offset > st.lastOffset { + st.lastOffset = offset + } } if fileSize > 0 { st.fileSize = fileSize } - shouldReport := force || time.Since(st.lastReport) >= 3*time.Second + shouldReport := (deltaBytes > 0 || force) && + st.bytesSent > 0 && + (force || time.Since(st.lastReport) >= 3*time.Second) bytesSent := st.bytesSent lastOffset := st.lastOffset size := st.fileSize