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 <cursoragent@cursor.com>
This commit is contained in:
Jos Vooges | STH 2026-08-25 02:40:50 +02:00
parent 181ab0b7f6
commit 6f78f13af5
2 changed files with 74 additions and 54 deletions

View file

@ -127,12 +127,11 @@ export async function registerStremioRoutes(app: FastifyInstance, config: Config
return reply.status(401).send({ error: "Invalid addon token" }); 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 }; return { streams };
}); });
// Lazy session: create on GET play so the token is fresh; redirect to node stream URL. // Fallback play URL (direct node redirect) — stream.json prefers eager node URLs.
// HEAD is treated as a probe — do not create a session (Stremio binge/prefetch uses HEAD/GET).
app.route({ app.route({
method: ["GET", "HEAD"], method: ["GET", "HEAD"],
url: "/stremio/:token/play/:mediaFileId", url: "/stremio/:token/play/:mediaFileId",
@ -160,7 +159,11 @@ export async function registerStremioRoutes(app: FastifyInstance, config: Config
try { try {
const session = await playback.createSession(mediaFileId, userId); 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) { } catch (err) {
const message = err instanceof Error ? err.message : "Playback unavailable"; const message = err instanceof Error ? err.message : "Playback unavailable";
request.log.error({ err, mediaFileId }, "playback create failed"); request.log.error({ err, mediaFileId }, "playback create failed");
@ -351,9 +354,9 @@ async function getSeriesMeta(id: string) {
async function getStreams( async function getStreams(
type: string, type: string,
id: string, id: string,
addonToken: string, userId: string,
_userId: string, playback: PlaybackService,
playback: PlaybackService log?: { warn: (obj: unknown, msg?: string) => void }
) { ) {
if (type === "movie") { if (type === "movie") {
const movie = await prisma.movie.findFirst({ const movie = await prisma.movie.findFirst({
@ -370,15 +373,20 @@ async function getStreams(
const streams = []; const streams = [];
for (const file of movie.mediaFiles) { for (const file of movie.mediaFiles) {
if (file.node.status !== "ONLINE" || file.node.revoked) continue; if (file.node.status !== "ONLINE" || file.node.revoked) continue;
streams.push({ try {
name: playback.formatStreamLabel(file), // Direct node URL — Stremio does not reliably follow Master 302/307 for video.
title: playback.formatStreamLabel(file), const session = await playback.createSession(file.id, userId);
// Lazy master URL — session is created when Stremio actually opens it. streams.push({
url: playback.resolveStreamUrl(addonToken, file.id), name: playback.formatStreamLabel(file),
behaviorHints: { title: playback.formatStreamLabel(file),
notWebReady: true, url: session.streamUrl,
}, behaviorHints: {
}); notWebReady: true,
},
});
} catch (err) {
log?.warn({ err, mediaFileId: file.id }, "skip stream: session create failed");
}
} }
return streams; return streams;
} }
@ -417,14 +425,19 @@ async function getStreams(
const streams = []; const streams = [];
for (const file of episode.mediaFiles) { for (const file of episode.mediaFiles) {
if (file.node.status !== "ONLINE" || file.node.revoked) continue; if (file.node.status !== "ONLINE" || file.node.revoked) continue;
streams.push({ try {
name: playback.formatStreamLabel(file), const session = await playback.createSession(file.id, userId);
title: playback.formatStreamLabel(file), streams.push({
url: playback.resolveStreamUrl(addonToken, file.id), name: playback.formatStreamLabel(file),
behaviorHints: { title: playback.formatStreamLabel(file),
notWebReady: true, url: session.streamUrl,
}, behaviorHints: {
}); notWebReady: true,
},
});
} catch (err) {
log?.warn({ err, mediaFileId: file.id }, "skip stream: session create failed");
}
} }
return streams; return streams;
} }

View file

@ -39,6 +39,7 @@ class NodeConnectionManager {
private offlineCheckInterval: ReturnType<typeof setInterval> | null = null; private offlineCheckInterval: ReturnType<typeof setInterval> | null = null;
private pendingSessionAcks = new Map<string, PendingAck>(); private pendingSessionAcks = new Map<string, PendingAck>();
private syncState = new Map<string, SyncState>(); private syncState = new Map<string, SyncState>();
private syncQueues = new Map<string, Promise<void>>();
init(config: Config): void { init(config: Config): void {
this.config = config; this.config = config;
@ -257,39 +258,45 @@ class NodeConnectionManager {
this.syncState.set(conn.nodeId, state); this.syncState.set(conn.nodeId, state);
} }
try { // ACK immediately so the node WS stays free for playback / next batches.
await this.librarySync.processFullSync( // Metadata matching can take seconds per batch and must not block the socket.
conn.nodeId, this.send(
payload.nodeRevision, socket,
payload.files, createMessage("FULL_LIBRARY_SYNC_ACK", {
{ isLast, fileIdsSoFar: state.fileIds } accepted: true,
); syncId: payload.syncId,
batchIndex: payload.batchIndex,
})
);
this.send( const prev = this.syncQueues.get(conn.nodeId) ?? Promise.resolve();
socket, const job = prev
createMessage("FULL_LIBRARY_SYNC_ACK", { .then(async () => {
accepted: true, if (!this.librarySync) return;
syncId: payload.syncId, await this.librarySync.processFullSync(
batchIndex: payload.batchIndex, 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)` : "")
); );
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); this.syncState.delete(conn.nodeId);
} else if (batchIndex % 10 === 0) { });
console.log( this.syncQueues.set(conn.nodeId, job);
`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" }));
}
} }
private handlePlaybackSessionAck(payload: CreatePlaybackSessionAckPayload): void { private handlePlaybackSessionAck(payload: CreatePlaybackSessionAckPayload): void {