From 6f78f13af5caebaf45cb08ecdbc9d3269d1f8e2c Mon Sep 17 00:00:00 2001 From: Jos Vooges | STH Date: Tue, 25 Aug 2026 02:40:50 +0200 Subject: [PATCH] Serve direct node stream URLs; ACK library sync without blocking WS. Stremio was hanging on Master 307 redirects. Sessions are created in stream.json with the node play URL. Full sync ACKs immediately so playback stays responsive. Co-authored-by: Cursor --- apps/master-api/src/stremio/routes.ts | 61 ++++++++++++--------- apps/master-api/src/websocket/manager.ts | 67 +++++++++++++----------- 2 files changed, 74 insertions(+), 54 deletions(-) diff --git a/apps/master-api/src/stremio/routes.ts b/apps/master-api/src/stremio/routes.ts index b3f5ff0..783f8bd 100644 --- a/apps/master-api/src/stremio/routes.ts +++ b/apps/master-api/src/stremio/routes.ts @@ -127,12 +127,11 @@ export async function registerStremioRoutes(app: FastifyInstance, config: Config return reply.status(401).send({ error: "Invalid addon token" }); } - const streams = await getStreams(type, id, token, userId, playback); + const streams = await getStreams(type, id, userId, playback, request.log); return { streams }; }); - // Lazy session: create on GET play so the token is fresh; redirect to node stream URL. - // HEAD is treated as a probe — do not create a session (Stremio binge/prefetch uses HEAD/GET). + // Fallback play URL (direct node redirect) — stream.json prefers eager node URLs. app.route({ method: ["GET", "HEAD"], url: "/stremio/:token/play/:mediaFileId", @@ -160,7 +159,11 @@ export async function registerStremioRoutes(app: FastifyInstance, config: Config try { const session = await playback.createSession(mediaFileId, userId); - return reply.code(307).redirect(session.streamUrl); + request.log.info( + { mediaFileId, streamUrl: session.streamUrl }, + "playback redirect to node" + ); + return reply.code(302).redirect(session.streamUrl); } catch (err) { const message = err instanceof Error ? err.message : "Playback unavailable"; request.log.error({ err, mediaFileId }, "playback create failed"); @@ -351,9 +354,9 @@ async function getSeriesMeta(id: string) { async function getStreams( type: string, id: string, - addonToken: string, - _userId: string, - playback: PlaybackService + userId: string, + playback: PlaybackService, + log?: { warn: (obj: unknown, msg?: string) => void } ) { if (type === "movie") { const movie = await prisma.movie.findFirst({ @@ -370,15 +373,20 @@ async function getStreams( const streams = []; for (const file of movie.mediaFiles) { if (file.node.status !== "ONLINE" || file.node.revoked) continue; - streams.push({ - name: playback.formatStreamLabel(file), - title: playback.formatStreamLabel(file), - // Lazy master URL — session is created when Stremio actually opens it. - url: playback.resolveStreamUrl(addonToken, file.id), - behaviorHints: { - notWebReady: true, - }, - }); + try { + // Direct node URL — Stremio does not reliably follow Master 302/307 for video. + const session = await playback.createSession(file.id, userId); + streams.push({ + name: playback.formatStreamLabel(file), + title: playback.formatStreamLabel(file), + url: session.streamUrl, + behaviorHints: { + notWebReady: true, + }, + }); + } catch (err) { + log?.warn({ err, mediaFileId: file.id }, "skip stream: session create failed"); + } } return streams; } @@ -417,14 +425,19 @@ async function getStreams( const streams = []; for (const file of episode.mediaFiles) { if (file.node.status !== "ONLINE" || file.node.revoked) continue; - streams.push({ - name: playback.formatStreamLabel(file), - title: playback.formatStreamLabel(file), - url: playback.resolveStreamUrl(addonToken, file.id), - behaviorHints: { - notWebReady: true, - }, - }); + try { + const session = await playback.createSession(file.id, userId); + streams.push({ + name: playback.formatStreamLabel(file), + title: playback.formatStreamLabel(file), + url: session.streamUrl, + behaviorHints: { + notWebReady: true, + }, + }); + } catch (err) { + log?.warn({ err, mediaFileId: file.id }, "skip stream: session create failed"); + } } return streams; } diff --git a/apps/master-api/src/websocket/manager.ts b/apps/master-api/src/websocket/manager.ts index de416d0..92036fe 100644 --- a/apps/master-api/src/websocket/manager.ts +++ b/apps/master-api/src/websocket/manager.ts @@ -39,6 +39,7 @@ class NodeConnectionManager { private offlineCheckInterval: ReturnType | null = null; private pendingSessionAcks = new Map(); private syncState = new Map(); + private syncQueues = new Map>(); init(config: Config): void { this.config = config; @@ -257,39 +258,45 @@ class NodeConnectionManager { this.syncState.set(conn.nodeId, state); } - try { - await this.librarySync.processFullSync( - conn.nodeId, - payload.nodeRevision, - payload.files, - { isLast, fileIdsSoFar: state.fileIds } - ); + // ACK immediately so the node WS stays free for playback / next batches. + // Metadata matching can take seconds per batch and must not block the socket. + this.send( + socket, + createMessage("FULL_LIBRARY_SYNC_ACK", { + accepted: true, + syncId: payload.syncId, + batchIndex: payload.batchIndex, + }) + ); - this.send( - socket, - createMessage("FULL_LIBRARY_SYNC_ACK", { - accepted: true, - syncId: payload.syncId, - batchIndex: payload.batchIndex, - }) - ); - - if (isLast) { - console.log( - `Full library sync complete for ${conn.nodeId}: ${state.fileIds.size} files` + - (payload.batchCount != null ? ` (${payload.batchCount} batches)` : "") + const prev = this.syncQueues.get(conn.nodeId) ?? Promise.resolve(); + const job = prev + .then(async () => { + if (!this.librarySync) return; + await this.librarySync.processFullSync( + conn.nodeId, + payload.nodeRevision, + payload.files, + { isLast, fileIdsSoFar: state!.fileIds } ); + + if (isLast) { + console.log( + `Full library sync complete for ${conn.nodeId}: ${state!.fileIds.size} files` + + (payload.batchCount != null ? ` (${payload.batchCount} batches)` : "") + ); + this.syncState.delete(conn.nodeId); + } else if (batchIndex % 10 === 0) { + console.log( + `Full library sync progress for ${conn.nodeId}: batch ${(batchIndex ?? 0) + 1}/${payload.batchCount ?? "?"}` + ); + } + }) + .catch((err) => { + console.error(`Full sync failed for ${conn.nodeId}:`, err); this.syncState.delete(conn.nodeId); - } else if (batchIndex % 10 === 0) { - console.log( - `Full library sync progress for ${conn.nodeId}: batch ${(batchIndex ?? 0) + 1}/${payload.batchCount ?? "?"}` - ); - } - } catch (err) { - console.error(`Full sync failed for ${conn.nodeId}:`, err); - this.syncState.delete(conn.nodeId); - this.send(socket, createMessage("ERROR", { message: "Full sync failed" })); - } + }); + this.syncQueues.set(conn.nodeId, job); } private handlePlaybackSessionAck(payload: CreatePlaybackSessionAckPayload): void {