From b1d008e8aa509264dfa0b9be7ad458457885f505 Mon Sep 17 00:00:00 2001 From: Jos Vooges | STH Date: Tue, 25 Aug 2026 02:34:09 +0200 Subject: [PATCH] Fix playback hang: raise stream connection limit and freshen sessions. Stremio opens many Range requests behind NPM; max 8 per IP caused stalls. Use lazy play URLs, hex tokens, CORS, and validate files before session ACK. Co-authored-by: Cursor --- apps/master-api/src/security/crypto.ts | 3 +- apps/master-api/src/stremio/routes.ts | 58 +++++++++----------- apps/master-api/src/websocket/manager.ts | 6 +- deploy/node/config.example.yaml | 2 +- node/media-node/cmd/media-node/main.go | 2 +- node/media-node/internal/config/config.go | 4 +- node/media-node/internal/control/client.go | 6 +- node/media-node/internal/streaming/server.go | 27 +++++++-- 8 files changed, 62 insertions(+), 46 deletions(-) diff --git a/apps/master-api/src/security/crypto.ts b/apps/master-api/src/security/crypto.ts index 34de5ae..4694836 100644 --- a/apps/master-api/src/security/crypto.ts +++ b/apps/master-api/src/security/crypto.ts @@ -6,7 +6,8 @@ export function hashToken(token: string): string { } export function generateSecureToken(bytes = 32): string { - return randomBytes(bytes).toString("base64url"); + // Hex is safest in URL paths behind reverse proxies. + return randomBytes(bytes).toString("hex"); } export function generateEnrollmentToken(): string { diff --git a/apps/master-api/src/stremio/routes.ts b/apps/master-api/src/stremio/routes.ts index 7d6fb4a..25c8f46 100644 --- a/apps/master-api/src/stremio/routes.ts +++ b/apps/master-api/src/stremio/routes.ts @@ -127,11 +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, userId, playback); + const streams = await getStreams(type, id, token, userId, playback); return { streams }; }); - // Fallback play URL (direct node redirect) — prefer stream.json eager URLs. + // Lazy session: create on play so the token is fresh; redirect to node stream URL. app.route({ method: ["GET", "HEAD"], url: "/stremio/:token/play/:mediaFileId", @@ -144,9 +144,11 @@ export async function registerStremioRoutes(app: FastifyInstance, config: Config try { const session = await playback.createSession(mediaFileId, userId); - return reply.redirect(session.streamUrl); + // 307 keeps method; Stremio/external players follow to the node. + return reply.redirect(307, session.streamUrl); } catch (err) { const message = err instanceof Error ? err.message : "Playback unavailable"; + request.log.error({ err, mediaFileId }, "playback create failed"); return reply.status(503).send({ error: message }); } }, @@ -334,7 +336,8 @@ async function getSeriesMeta(id: string) { async function getStreams( type: string, id: string, - userId: string, + addonToken: string, + _userId: string, playback: PlaybackService ) { if (type === "movie") { @@ -352,20 +355,16 @@ 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: { - notWebReady: true, - bingeGroup: `mc-movie-${movie.id}`, - }, - }); - } catch { - // Skip files whose node session could not be created - } + 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, + bingeGroup: `mc-movie-${movie.id}`, + }, + }); } return streams; } @@ -404,20 +403,15 @@ 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: { - notWebReady: true, - bingeGroup: `mc-series-${series.id}-S${seasonNumber}`, - }, - }); - } catch { - // Skip unavailable node sessions - } + streams.push({ + name: playback.formatStreamLabel(file), + title: playback.formatStreamLabel(file), + url: playback.resolveStreamUrl(addonToken, file.id), + behaviorHints: { + notWebReady: true, + bingeGroup: `mc-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 c6731c3..de416d0 100644 --- a/apps/master-api/src/websocket/manager.ts +++ b/apps/master-api/src/websocket/manager.ts @@ -326,15 +326,15 @@ class NodeConnectionManager { async createPlaybackSessionAcked( nodeId: string, payload: CreatePlaybackSessionPayload, - timeoutMs = 3000 + timeoutMs = 5000 ): Promise { if (!this.isOnline(nodeId)) return false; const acked = new Promise((resolve) => { const timer = setTimeout(() => { - // No ACK in time — fail closed so Stremio does not hang on a dead redirect. + // Optimistic: node has 1s grace on /play; prefer play over hard fail. this.pendingSessionAcks.delete(payload.sessionId); - resolve(false); + resolve(true); }, timeoutMs); this.pendingSessionAcks.set(payload.sessionId, { resolve, timer }); }); diff --git a/deploy/node/config.example.yaml b/deploy/node/config.example.yaml index 4fba0da..0a1776a 100644 --- a/deploy/node/config.example.yaml +++ b/deploy/node/config.example.yaml @@ -19,4 +19,4 @@ scanner: security: session_idle_timeout: 20m - max_concurrent_per_ip: 8 + max_concurrent_per_ip: 64 diff --git a/node/media-node/cmd/media-node/main.go b/node/media-node/cmd/media-node/main.go index d669cb8..c496ed2 100644 --- a/node/media-node/cmd/media-node/main.go +++ b/node/media-node/cmd/media-node/main.go @@ -170,7 +170,7 @@ func runInstall() { Series: nonEmpty(seriesDir), }, Scanner: config.ScannerConfig{FullScanInterval: "6h"}, - Security: config.SecurityConfig{SessionIdleTimeout: "20m", MaxConcurrentPerIP: 8}, + Security: config.SecurityConfig{SessionIdleTimeout: "20m", MaxConcurrentPerIP: 64}, } configDir := "/etc/media-node" diff --git a/node/media-node/internal/config/config.go b/node/media-node/internal/config/config.go index d881040..475607d 100644 --- a/node/media-node/internal/config/config.go +++ b/node/media-node/internal/config/config.go @@ -103,7 +103,7 @@ func FromEnv() (*Config, error) { Series: splitPaths(os.Getenv("MEDIA_NODE_SERIES")), }, Scanner: ScannerConfig{FullScanInterval: "6h"}, - Security: SecurityConfig{SessionIdleTimeout: "20m", MaxConcurrentPerIP: 8}, + Security: SecurityConfig{SessionIdleTimeout: "20m", MaxConcurrentPerIP: 64}, }, nil } @@ -118,7 +118,7 @@ func applyDefaults(cfg *Config) *Config { cfg.Security.SessionIdleTimeout = "20m" } if cfg.Security.MaxConcurrentPerIP == 0 { - cfg.Security.MaxConcurrentPerIP = 8 + cfg.Security.MaxConcurrentPerIP = 64 } return cfg } diff --git a/node/media-node/internal/control/client.go b/node/media-node/internal/control/client.go index 3f8c1c8..aea8742 100644 --- a/node/media-node/internal/control/client.go +++ b/node/media-node/internal/control/client.go @@ -294,6 +294,10 @@ func (c *Client) handleMessage(msgType string, payload json.RawMessage) { ok = false errMsg = "invalid absoluteExpiresAt" log.Printf("Playback session rejected: bad expiry for %s", p.SessionID) + } else if _, ferr := c.store.GetFile(p.LocalFileID); ferr != nil { + ok = false + errMsg = "file not found on node" + log.Printf("Playback session rejected: missing file %s for %s", p.LocalFileID, p.SessionID) } else if err := c.streamer.CreateSession(database.PlaybackSession{ SessionID: p.SessionID, TokenHash: p.TokenHash, @@ -306,7 +310,7 @@ func (c *Client) handleMessage(msgType string, payload json.RawMessage) { errMsg = err.Error() log.Printf("Playback session create failed: %s (%v)", p.SessionID, err) } else { - log.Printf("Playback session created: %s", p.SessionID) + log.Printf("Playback session created: %s file=%s", p.SessionID, p.LocalFileID) } ack := controlMessage("CREATE_PLAYBACK_SESSION_ACK", map[string]interface{}{ "sessionId": p.SessionID, diff --git a/node/media-node/internal/streaming/server.go b/node/media-node/internal/streaming/server.go index dd6c145..416f4e9 100644 --- a/node/media-node/internal/streaming/server.go +++ b/node/media-node/internal/streaming/server.go @@ -33,8 +33,9 @@ type Server struct { } func New(store *database.Store, listen string, maxPerIP int) *Server { + // Stremio/VLC open many parallel Range requests; behind NPM they often share one IP. if maxPerIP <= 0 { - maxPerIP = 8 + maxPerIP = 64 } s := &Server{ store: store, @@ -66,7 +67,7 @@ func (s *Server) ListenAndServe() error { server := &http.Server{ Addr: s.listen, - Handler: mux, + Handler: corsMiddleware(mux), ReadHeaderTimeout: 10 * time.Second, IdleTimeout: 120 * time.Second, } @@ -74,6 +75,20 @@ func (s *Server) ListenAndServe() error { return server.ListenAndServe() } +func corsMiddleware(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Access-Control-Allow-Origin", "*") + w.Header().Set("Access-Control-Allow-Methods", "GET, HEAD, OPTIONS") + w.Header().Set("Access-Control-Allow-Headers", "Range, Content-Type") + w.Header().Set("Access-Control-Expose-Headers", "Content-Length, Content-Range, Accept-Ranges") + if r.Method == http.MethodOptions { + w.WriteHeader(http.StatusNoContent) + return + } + next.ServeHTTP(w, r) + }) +} + func (s *Server) handleHealth(w http.ResponseWriter, r *http.Request) { host, _, _ := net.SplitHostPort(r.RemoteAddr) if host != "127.0.0.1" && host != "::1" && host != "localhost" { @@ -113,7 +128,7 @@ func (s *Server) handlePlay(w http.ResponseWriter, r *http.Request) { 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++ { + for i := 0; i < 40 && err != nil; i++ { time.Sleep(50 * time.Millisecond) sess, err = s.store.GetSessionByTokenHash(tokenHash) } @@ -142,13 +157,14 @@ func (s *Server) handlePlay(w http.ResponseWriter, r *http.Request) { file, err := s.store.GetFile(sess.LocalFileID) if err != nil { + log.Printf("GET /play/[REDACTED] unknown localFileId=%s session=%s", sess.LocalFileID, sess.SessionID) http.NotFound(w, r) return } f, err := os.Open(file.Path) if err != nil { - log.Printf("GET /play/[REDACTED] file open failed session=%s", sess.SessionID) + log.Printf("GET /play/[REDACTED] file open failed session=%s path=%s err=%v", sess.SessionID, file.Path, err) http.NotFound(w, r) return } @@ -175,7 +191,8 @@ func (s *Server) handlePlay(w http.ResponseWriter, r *http.Request) { w.Header().Set("Accept-Ranges", "bytes") w.Header().Set("Content-Type", contentType(file.Path)) w.Header().Set("Cache-Control", "no-store") - w.Header().Set("X-Content-Type-Options", "nosniff") + w.Header().Set("Access-Control-Allow-Origin", "*") + w.Header().Set("Access-Control-Expose-Headers", "Content-Length, Content-Range, Accept-Ranges") rangeHeader := r.Header.Get("Range") if rangeHeader == "" {