From 7741500e73931d8bcaf61c61a074f294a48d8fb7 Mon Sep 17 00:00:00 2001 From: Jos Vooges | STH Date: Tue, 25 Aug 2026 01:47:54 +0200 Subject: [PATCH] Fix Stremio playback: lazy sessions, remove notWebReady, node ACK. Co-authored-by: Cursor --- apps/admin-ui/src/app/settings/page.tsx | 4 +- apps/master-api/src/playback/service.ts | 10 ++- apps/master-api/src/stremio/routes.ts | 84 ++++++++++++-------- apps/master-api/src/websocket/manager.ts | 48 +++++++++++ node/media-node/internal/control/client.go | 28 ++++++- node/media-node/internal/streaming/server.go | 7 ++ packages/protocol/src/index.ts | 7 ++ 7 files changed, 148 insertions(+), 40 deletions(-) diff --git a/apps/admin-ui/src/app/settings/page.tsx b/apps/admin-ui/src/app/settings/page.tsx index e067f62..f2a4043 100644 --- a/apps/admin-ui/src/app/settings/page.tsx +++ b/apps/admin-ui/src/app/settings/page.tsx @@ -18,8 +18,8 @@ export default function SettingsPage() { }); const data = await res.json(); setAddonToken(data.token); - const base = window.location.origin.replace(":3001", ":3000"); - setInstallUrl(`${base}/stremio/${data.token}/manifest.json`); + // Addon draait op Master, niet op Admin + setInstallUrl(`https://master.vonas.nl/stremio/${data.token}/manifest.json`); } return ( diff --git a/apps/master-api/src/playback/service.ts b/apps/master-api/src/playback/service.ts index d5c7bc3..8447174 100644 --- a/apps/master-api/src/playback/service.ts +++ b/apps/master-api/src/playback/service.ts @@ -51,7 +51,7 @@ export class PlaybackService { }, }); - const pushed = nodeConnectionManager.createPlaybackSession(mediaFile.nodeId, { + const acked = await nodeConnectionManager.createPlaybackSessionAcked(mediaFile.nodeId, { sessionId, tokenHash, localFileId: mediaFile.localFileId, @@ -59,7 +59,7 @@ export class PlaybackService { absoluteExpiresAt: absoluteExpiresAt.toISOString(), }); - if (!pushed) { + if (!acked) { await prisma.playbackSession.update({ where: { id: sessionId }, data: { status: "REVOKED", revoked: true }, @@ -78,6 +78,12 @@ export class PlaybackService { }; } + /** Public stream URL for Stremio — session is created only when this URL is opened. */ + resolveStreamUrl(addonToken: string, mediaFileId: string): string { + const base = this.config.PUBLIC_URL.replace(/\/$/, ""); + return `${base}/stremio/${addonToken}/play/${mediaFileId}`; + } + async revokeSession(sessionId: string): Promise { const session = await prisma.playbackSession.findUnique({ where: { id: sessionId }, diff --git a/apps/master-api/src/stremio/routes.ts b/apps/master-api/src/stremio/routes.ts index e8eed25..e5a3213 100644 --- a/apps/master-api/src/stremio/routes.ts +++ b/apps/master-api/src/stremio/routes.ts @@ -110,9 +110,30 @@ export async function registerStremioRoutes(app: FastifyInstance, config: Config return reply.status(401).send({ error: "Invalid addon token" }); } - const streams = await getStreams(type, id, userId, playback); + const streams = await getStreams(type, id, token, playback); return { streams }; }); + + // Lazy session: Stremio opens this URL only when actually playing (or prefetching media). + app.route({ + method: ["GET", "HEAD"], + url: "/stremio/:token/play/:mediaFileId", + handler: async (request, reply) => { + const { token, mediaFileId } = request.params as { token: string; mediaFileId: string }; + const userId = await validateAddonToken(token); + if (!userId) { + return reply.status(401).send({ error: "Invalid addon token" }); + } + + try { + const session = await playback.createSession(mediaFileId, userId); + return reply.redirect(session.streamUrl); + } catch (err) { + const message = err instanceof Error ? err.message : "Playback unavailable"; + return reply.status(503).send({ error: message }); + } + }, + }); } function parseExtraSearch(extra: string): string | undefined { @@ -247,7 +268,12 @@ async function getSeriesMeta(id: string) { }, include: { seasons: { - include: { episodes: { orderBy: [{ seasonNumber: "asc" }, { episodeNumber: "asc" }] } }, + include: { + episodes: { + orderBy: [{ seasonNumber: "asc" }, { episodeNumber: "asc" }], + include: { mediaFiles: { where: { available: true }, take: 1 } }, + }, + }, orderBy: { seasonNumber: "asc" }, }, }, @@ -255,14 +281,16 @@ async function getSeriesMeta(id: string) { if (!series) return null; const videos = series.seasons.flatMap((season) => - season.episodes.map((ep) => ({ - id: `${series.imdbId ?? `sth:series:${series.id}`}:${ep.seasonNumber}:${ep.episodeNumber}`, - title: ep.title ?? `Episode ${ep.episodeNumber}`, - season: ep.seasonNumber, - episode: ep.episodeNumber, - overview: ep.overview ?? undefined, - released: ep.airDate?.toISOString().slice(0, 10), - })) + season.episodes + .filter((ep) => ep.mediaFiles.length > 0) + .map((ep) => ({ + id: `${series.imdbId ?? `sth:series:${series.id}`}:${ep.seasonNumber}:${ep.episodeNumber}`, + title: ep.title ?? `Episode ${ep.episodeNumber}`, + season: ep.seasonNumber, + episode: ep.episodeNumber, + overview: ep.overview ?? undefined, + released: ep.airDate?.toISOString().slice(0, 10), + })) ); return { @@ -281,7 +309,7 @@ async function getSeriesMeta(id: string) { async function getStreams( type: string, id: string, - userId: string, + addonToken: string, playback: PlaybackService ) { if (type === "movie") { @@ -299,17 +327,12 @@ async function getStreams( const streams = []; for (const file of movie.mediaFiles) { if (file.node.status !== "ONLINE" || file.node.revoked) continue; - try { - const session = await playback.createSession(file.id, userId); - streams.push({ - name: playback.formatStreamLabel(file), - title: playback.formatStreamLabel(file), - url: session.streamUrl, - behaviorHints: { bingeGroup: movie.id, notWebReady: true }, - }); - } catch { - // skip unavailable nodes - } + streams.push({ + name: playback.formatStreamLabel(file), + title: playback.formatStreamLabel(file), + url: playback.resolveStreamUrl(addonToken, file.id), + behaviorHints: { bingeGroup: `sth-movie-${movie.id}` }, + }); } return streams; } @@ -348,17 +371,12 @@ async function getStreams( const streams = []; for (const file of episode.mediaFiles) { if (file.node.status !== "ONLINE" || file.node.revoked) continue; - try { - const session = await playback.createSession(file.id, userId); - streams.push({ - name: playback.formatStreamLabel(file), - title: playback.formatStreamLabel(file), - url: session.streamUrl, - behaviorHints: { bingeGroup: `${series.id}-S${seasonNumber}`, notWebReady: true }, - }); - } catch { - // skip - } + streams.push({ + name: playback.formatStreamLabel(file), + title: playback.formatStreamLabel(file), + url: playback.resolveStreamUrl(addonToken, file.id), + behaviorHints: { bingeGroup: `sth-series-${series.id}-S${seasonNumber}` }, + }); } return streams; } diff --git a/apps/master-api/src/websocket/manager.ts b/apps/master-api/src/websocket/manager.ts index b00fdf8..bf304d1 100644 --- a/apps/master-api/src/websocket/manager.ts +++ b/apps/master-api/src/websocket/manager.ts @@ -13,6 +13,7 @@ import type { HelloPayload, LibraryEventPayload, PlaybackSessionEndedPayload, + CreatePlaybackSessionAckPayload, } from "@media-cluster/protocol"; interface NodeConnection { @@ -21,11 +22,17 @@ interface NodeConnection { lastHeartbeat: Date; } +interface PendingAck { + resolve: (ok: boolean) => void; + timer: ReturnType; +} + class NodeConnectionManager { private connections = new Map(); private librarySync: LibrarySyncService | null = null; private config: Config | null = null; private offlineCheckInterval: ReturnType | null = null; + private pendingSessionAcks = new Map(); init(config: Config): void { this.config = config; @@ -112,6 +119,9 @@ class NodeConnectionManager { case "PLAYBACK_SESSION_ENDED": await this.handlePlaybackEnded(message.payload as PlaybackSessionEndedPayload); break; + case "CREATE_PLAYBACK_SESSION_ACK": + this.handlePlaybackSessionAck(message.payload as CreatePlaybackSessionAckPayload); + break; default: this.send(socket, createMessage("ERROR", { message: `Unknown type: ${message.type}` })); } @@ -230,6 +240,14 @@ class NodeConnectionManager { this.send(socket, createMessage("FULL_LIBRARY_SYNC_ACK", { accepted: true })); } + private handlePlaybackSessionAck(payload: CreatePlaybackSessionAckPayload): void { + const pending = this.pendingSessionAcks.get(payload.sessionId); + if (!pending) return; + clearTimeout(pending.timer); + this.pendingSessionAcks.delete(payload.sessionId); + pending.resolve(payload.ok); + } + private async handlePlaybackEnded(payload: PlaybackSessionEndedPayload): Promise { await prisma.playbackSession.updateMany({ where: { id: payload.sessionId }, @@ -252,6 +270,36 @@ class NodeConnectionManager { return this.sendToNode(nodeId, "CREATE_PLAYBACK_SESSION", payload); } + /** Push session to node and wait for ACK (or timeout). */ + async createPlaybackSessionAcked( + nodeId: string, + payload: CreatePlaybackSessionPayload, + timeoutMs = 800 + ): Promise { + if (!this.isOnline(nodeId)) return false; + + const acked = new Promise((resolve) => { + const timer = setTimeout(() => { + // Older nodes may not send ACK — proceed optimistically; node has grace retry. + this.pendingSessionAcks.delete(payload.sessionId); + resolve(true); + }, timeoutMs); + this.pendingSessionAcks.set(payload.sessionId, { resolve, timer }); + }); + + const pushed = this.createPlaybackSession(nodeId, payload); + if (!pushed) { + const pending = this.pendingSessionAcks.get(payload.sessionId); + if (pending) { + clearTimeout(pending.timer); + this.pendingSessionAcks.delete(payload.sessionId); + } + return false; + } + + return acked; + } + revokePlaybackSession(nodeId: string, sessionId: string): boolean { return this.sendToNode(nodeId, "REVOKE_PLAYBACK_SESSION", { sessionId }); } diff --git a/node/media-node/internal/control/client.go b/node/media-node/internal/control/client.go index 82a89df..064f60d 100644 --- a/node/media-node/internal/control/client.go +++ b/node/media-node/internal/control/client.go @@ -229,16 +229,38 @@ func (c *Client) handleMessage(msgType string, payload json.RawMessage) { AbsoluteExpiresAt string `json:"absoluteExpiresAt"` } if json.Unmarshal(payload, &p) == nil { - expires, _ := time.Parse(time.RFC3339, p.AbsoluteExpiresAt) - _ = c.streamer.CreateSession(database.PlaybackSession{ + expires, err := time.Parse(time.RFC3339Nano, p.AbsoluteExpiresAt) + if err != nil { + expires, err = time.Parse(time.RFC3339, p.AbsoluteExpiresAt) + } + ok := true + errMsg := "" + if err != nil { + ok = false + errMsg = "invalid absoluteExpiresAt" + log.Printf("Playback session rejected: bad expiry for %s", p.SessionID) + } else if err := c.streamer.CreateSession(database.PlaybackSession{ SessionID: p.SessionID, TokenHash: p.TokenHash, LocalFileID: p.LocalFileID, IdleTimeoutSeconds: p.IdleTimeoutSeconds, AbsoluteExpiresAt: expires, LastActivity: time.Now().UTC(), + }); err != nil { + ok = false + errMsg = err.Error() + log.Printf("Playback session create failed: %s (%v)", p.SessionID, err) + } else { + log.Printf("Playback session created: %s", p.SessionID) + } + ack := controlMessage("CREATE_PLAYBACK_SESSION_ACK", map[string]interface{}{ + "sessionId": p.SessionID, + "ok": ok, + "error": errMsg, }) - log.Printf("Playback session created: %s", p.SessionID) + if werr := c.write(ack); werr != nil { + log.Printf("failed to send session ACK: %v", werr) + } } case "REVOKE_PLAYBACK_SESSION": var p struct { diff --git a/node/media-node/internal/streaming/server.go b/node/media-node/internal/streaming/server.go index a35cb18..dd6c145 100644 --- a/node/media-node/internal/streaming/server.go +++ b/node/media-node/internal/streaming/server.go @@ -111,6 +111,13 @@ func (s *Server) handlePlay(w http.ResponseWriter, r *http.Request) { tokenHash := hashToken(token) sess, err := s.store.GetSessionByTokenHash(tokenHash) + if err != nil { + // Brief grace: Master may have redirected before WS session landed. + for i := 0; i < 20 && err != nil; i++ { + time.Sleep(50 * time.Millisecond) + sess, err = s.store.GetSessionByTokenHash(tokenHash) + } + } if err != nil { log.Printf("GET /play/[REDACTED] invalid session from %s", clientIP) http.NotFound(w, r) diff --git a/packages/protocol/src/index.ts b/packages/protocol/src/index.ts index 13b32f0..3772909 100644 --- a/packages/protocol/src/index.ts +++ b/packages/protocol/src/index.ts @@ -18,6 +18,7 @@ export type ControlMessageType = | "FULL_LIBRARY_SYNC" | "FULL_LIBRARY_SYNC_ACK" | "CREATE_PLAYBACK_SESSION" + | "CREATE_PLAYBACK_SESSION_ACK" | "REVOKE_PLAYBACK_SESSION" | "RESCAN" | "CONFIG_UPDATE" @@ -68,6 +69,12 @@ export interface PlaybackSessionEndedPayload { reason: string; } +export interface CreatePlaybackSessionAckPayload { + sessionId: string; + ok: boolean; + error?: string; +} + export function createMessage( type: ControlMessageType, payload: T,