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 <cursoragent@cursor.com>
This commit is contained in:
parent
eb8ae1e85a
commit
b1d008e8aa
8 changed files with 62 additions and 46 deletions
|
|
@ -6,7 +6,8 @@ export function hashToken(token: string): string {
|
||||||
}
|
}
|
||||||
|
|
||||||
export function generateSecureToken(bytes = 32): 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 {
|
export function generateEnrollmentToken(): string {
|
||||||
|
|
|
||||||
|
|
@ -127,11 +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, userId, playback);
|
const streams = await getStreams(type, id, token, userId, playback);
|
||||||
return { streams };
|
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({
|
app.route({
|
||||||
method: ["GET", "HEAD"],
|
method: ["GET", "HEAD"],
|
||||||
url: "/stremio/:token/play/:mediaFileId",
|
url: "/stremio/:token/play/:mediaFileId",
|
||||||
|
|
@ -144,9 +144,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.redirect(session.streamUrl);
|
// 307 keeps method; Stremio/external players follow to the node.
|
||||||
|
return reply.redirect(307, 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");
|
||||||
return reply.status(503).send({ error: message });
|
return reply.status(503).send({ error: message });
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
|
@ -334,7 +336,8 @@ async function getSeriesMeta(id: string) {
|
||||||
async function getStreams(
|
async function getStreams(
|
||||||
type: string,
|
type: string,
|
||||||
id: string,
|
id: string,
|
||||||
userId: string,
|
addonToken: string,
|
||||||
|
_userId: string,
|
||||||
playback: PlaybackService
|
playback: PlaybackService
|
||||||
) {
|
) {
|
||||||
if (type === "movie") {
|
if (type === "movie") {
|
||||||
|
|
@ -352,20 +355,16 @@ 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;
|
||||||
try {
|
streams.push({
|
||||||
const session = await playback.createSession(file.id, userId);
|
name: playback.formatStreamLabel(file),
|
||||||
streams.push({
|
title: playback.formatStreamLabel(file),
|
||||||
name: playback.formatStreamLabel(file),
|
// Lazy master URL — session is created when Stremio actually opens it.
|
||||||
title: playback.formatStreamLabel(file),
|
url: playback.resolveStreamUrl(addonToken, file.id),
|
||||||
url: session.streamUrl,
|
behaviorHints: {
|
||||||
behaviorHints: {
|
notWebReady: true,
|
||||||
notWebReady: true,
|
bingeGroup: `mc-movie-${movie.id}`,
|
||||||
bingeGroup: `mc-movie-${movie.id}`,
|
},
|
||||||
},
|
});
|
||||||
});
|
|
||||||
} catch {
|
|
||||||
// Skip files whose node session could not be created
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
return streams;
|
return streams;
|
||||||
}
|
}
|
||||||
|
|
@ -404,20 +403,15 @@ 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;
|
||||||
try {
|
streams.push({
|
||||||
const session = await playback.createSession(file.id, userId);
|
name: playback.formatStreamLabel(file),
|
||||||
streams.push({
|
title: playback.formatStreamLabel(file),
|
||||||
name: playback.formatStreamLabel(file),
|
url: playback.resolveStreamUrl(addonToken, file.id),
|
||||||
title: playback.formatStreamLabel(file),
|
behaviorHints: {
|
||||||
url: session.streamUrl,
|
notWebReady: true,
|
||||||
behaviorHints: {
|
bingeGroup: `mc-series-${series.id}-S${seasonNumber}`,
|
||||||
notWebReady: true,
|
},
|
||||||
bingeGroup: `mc-series-${series.id}-S${seasonNumber}`,
|
});
|
||||||
},
|
|
||||||
});
|
|
||||||
} catch {
|
|
||||||
// Skip unavailable node sessions
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
return streams;
|
return streams;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -326,15 +326,15 @@ class NodeConnectionManager {
|
||||||
async createPlaybackSessionAcked(
|
async createPlaybackSessionAcked(
|
||||||
nodeId: string,
|
nodeId: string,
|
||||||
payload: CreatePlaybackSessionPayload,
|
payload: CreatePlaybackSessionPayload,
|
||||||
timeoutMs = 3000
|
timeoutMs = 5000
|
||||||
): Promise<boolean> {
|
): Promise<boolean> {
|
||||||
if (!this.isOnline(nodeId)) return false;
|
if (!this.isOnline(nodeId)) return false;
|
||||||
|
|
||||||
const acked = new Promise<boolean>((resolve) => {
|
const acked = new Promise<boolean>((resolve) => {
|
||||||
const timer = setTimeout(() => {
|
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);
|
this.pendingSessionAcks.delete(payload.sessionId);
|
||||||
resolve(false);
|
resolve(true);
|
||||||
}, timeoutMs);
|
}, timeoutMs);
|
||||||
this.pendingSessionAcks.set(payload.sessionId, { resolve, timer });
|
this.pendingSessionAcks.set(payload.sessionId, { resolve, timer });
|
||||||
});
|
});
|
||||||
|
|
|
||||||
|
|
@ -19,4 +19,4 @@ scanner:
|
||||||
|
|
||||||
security:
|
security:
|
||||||
session_idle_timeout: 20m
|
session_idle_timeout: 20m
|
||||||
max_concurrent_per_ip: 8
|
max_concurrent_per_ip: 64
|
||||||
|
|
|
||||||
|
|
@ -170,7 +170,7 @@ func runInstall() {
|
||||||
Series: nonEmpty(seriesDir),
|
Series: nonEmpty(seriesDir),
|
||||||
},
|
},
|
||||||
Scanner: config.ScannerConfig{FullScanInterval: "6h"},
|
Scanner: config.ScannerConfig{FullScanInterval: "6h"},
|
||||||
Security: config.SecurityConfig{SessionIdleTimeout: "20m", MaxConcurrentPerIP: 8},
|
Security: config.SecurityConfig{SessionIdleTimeout: "20m", MaxConcurrentPerIP: 64},
|
||||||
}
|
}
|
||||||
|
|
||||||
configDir := "/etc/media-node"
|
configDir := "/etc/media-node"
|
||||||
|
|
|
||||||
|
|
@ -103,7 +103,7 @@ func FromEnv() (*Config, error) {
|
||||||
Series: splitPaths(os.Getenv("MEDIA_NODE_SERIES")),
|
Series: splitPaths(os.Getenv("MEDIA_NODE_SERIES")),
|
||||||
},
|
},
|
||||||
Scanner: ScannerConfig{FullScanInterval: "6h"},
|
Scanner: ScannerConfig{FullScanInterval: "6h"},
|
||||||
Security: SecurityConfig{SessionIdleTimeout: "20m", MaxConcurrentPerIP: 8},
|
Security: SecurityConfig{SessionIdleTimeout: "20m", MaxConcurrentPerIP: 64},
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -118,7 +118,7 @@ func applyDefaults(cfg *Config) *Config {
|
||||||
cfg.Security.SessionIdleTimeout = "20m"
|
cfg.Security.SessionIdleTimeout = "20m"
|
||||||
}
|
}
|
||||||
if cfg.Security.MaxConcurrentPerIP == 0 {
|
if cfg.Security.MaxConcurrentPerIP == 0 {
|
||||||
cfg.Security.MaxConcurrentPerIP = 8
|
cfg.Security.MaxConcurrentPerIP = 64
|
||||||
}
|
}
|
||||||
return cfg
|
return cfg
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -294,6 +294,10 @@ func (c *Client) handleMessage(msgType string, payload json.RawMessage) {
|
||||||
ok = false
|
ok = false
|
||||||
errMsg = "invalid absoluteExpiresAt"
|
errMsg = "invalid absoluteExpiresAt"
|
||||||
log.Printf("Playback session rejected: bad expiry for %s", p.SessionID)
|
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{
|
} else if err := c.streamer.CreateSession(database.PlaybackSession{
|
||||||
SessionID: p.SessionID,
|
SessionID: p.SessionID,
|
||||||
TokenHash: p.TokenHash,
|
TokenHash: p.TokenHash,
|
||||||
|
|
@ -306,7 +310,7 @@ func (c *Client) handleMessage(msgType string, payload json.RawMessage) {
|
||||||
errMsg = err.Error()
|
errMsg = err.Error()
|
||||||
log.Printf("Playback session create failed: %s (%v)", p.SessionID, err)
|
log.Printf("Playback session create failed: %s (%v)", p.SessionID, err)
|
||||||
} else {
|
} 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{}{
|
ack := controlMessage("CREATE_PLAYBACK_SESSION_ACK", map[string]interface{}{
|
||||||
"sessionId": p.SessionID,
|
"sessionId": p.SessionID,
|
||||||
|
|
|
||||||
|
|
@ -33,8 +33,9 @@ type Server struct {
|
||||||
}
|
}
|
||||||
|
|
||||||
func New(store *database.Store, listen string, maxPerIP int) *Server {
|
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 {
|
if maxPerIP <= 0 {
|
||||||
maxPerIP = 8
|
maxPerIP = 64
|
||||||
}
|
}
|
||||||
s := &Server{
|
s := &Server{
|
||||||
store: store,
|
store: store,
|
||||||
|
|
@ -66,7 +67,7 @@ func (s *Server) ListenAndServe() error {
|
||||||
|
|
||||||
server := &http.Server{
|
server := &http.Server{
|
||||||
Addr: s.listen,
|
Addr: s.listen,
|
||||||
Handler: mux,
|
Handler: corsMiddleware(mux),
|
||||||
ReadHeaderTimeout: 10 * time.Second,
|
ReadHeaderTimeout: 10 * time.Second,
|
||||||
IdleTimeout: 120 * time.Second,
|
IdleTimeout: 120 * time.Second,
|
||||||
}
|
}
|
||||||
|
|
@ -74,6 +75,20 @@ func (s *Server) ListenAndServe() error {
|
||||||
return server.ListenAndServe()
|
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) {
|
func (s *Server) handleHealth(w http.ResponseWriter, r *http.Request) {
|
||||||
host, _, _ := net.SplitHostPort(r.RemoteAddr)
|
host, _, _ := net.SplitHostPort(r.RemoteAddr)
|
||||||
if host != "127.0.0.1" && host != "::1" && host != "localhost" {
|
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)
|
sess, err := s.store.GetSessionByTokenHash(tokenHash)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// Brief grace: Master may have redirected before WS session landed.
|
// 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)
|
time.Sleep(50 * time.Millisecond)
|
||||||
sess, err = s.store.GetSessionByTokenHash(tokenHash)
|
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)
|
file, err := s.store.GetFile(sess.LocalFileID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
log.Printf("GET /play/[REDACTED] unknown localFileId=%s session=%s", sess.LocalFileID, sess.SessionID)
|
||||||
http.NotFound(w, r)
|
http.NotFound(w, r)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
f, err := os.Open(file.Path)
|
f, err := os.Open(file.Path)
|
||||||
if err != nil {
|
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)
|
http.NotFound(w, r)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -175,7 +191,8 @@ func (s *Server) handlePlay(w http.ResponseWriter, r *http.Request) {
|
||||||
w.Header().Set("Accept-Ranges", "bytes")
|
w.Header().Set("Accept-Ranges", "bytes")
|
||||||
w.Header().Set("Content-Type", contentType(file.Path))
|
w.Header().Set("Content-Type", contentType(file.Path))
|
||||||
w.Header().Set("Cache-Control", "no-store")
|
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")
|
rangeHeader := r.Header.Get("Range")
|
||||||
if rangeHeader == "" {
|
if rangeHeader == "" {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue