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 {