fix(jellycompat): keep still-playing direct-play sessions in the activity view (#262)

* test(jellycompat): implement ListProgressSince on userstore test fake

PR #258 added ListProgressSince to the userstore.UserStore interface but
missed the progressCountingStore stub, leaving the jellycompat test build
broken on main. Stub it with the same panic("unused") the sibling methods
use.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(jellycompat): keep still-playing direct-play sessions in the activity view

An Infuse Static=true direct play skips PlaybackInfo, so its playback
reports carry a client-generated PlaySessionId. handlePlaybackReport only
resolved reports by exact PlaySessionId, so every report silently no-op'd:
the admin Activity View position froze, nothing refreshed the upstream
session's activity, and stale cleanup reaped the session ~45s in while the
client kept playing. Once reaped, ensureUpstreamPlayback early-returned the
dangling session id forever, so the stream card never came back.

Fix all three legs:
- handlePlaybackReport falls back to the same token-scoped route lookup the
  stream path uses (ItemId, then MediaSourceId) when the PlaySessionId is
  unknown, and revives a reaped upstream session when a progress report
  proves the client is still playing.
- HandleVideoStream marks an in-flight media transport around direct-play
  and remux serving, mirroring the native stream handler, so long-lived
  range transfers keep the session alive without progress reports.
- ensureUpstreamPlayback verifies the upstream session still exists before
  reusing it, recreating it under the same play session when it was reaped.

FindByRoute now compares RouteItemID with UUID normalization so report ids
match regardless of client dash/case formatting.

Closes #244

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(jellycompat): harden report fallback and upstream revive per adversarial review

Address three review findings on the issue #244 fix:

- Deterministic report binding: the stream path now records the client's
  own PlaySessionId (Static=true direct play never learns the server id)
  as an alias on the play session it binds, and playback reports resolve
  by that alias before falling back to the ambiguous item/source route
  scan. A Stopped report that only matched by bare route no longer tears
  the session down (it may not own it when the same item plays twice
  under one token); stale cleanup owns that path.
- Same-method revive closes any transcode still keyed to the reaped
  upstream id before recreating the session, so a second ffmpeg cannot
  start alongside an orphaned one.
- Attaching a new upstream session is now a compare-and-swap on the
  observed upstream id: the loser of a concurrent revive/range-request
  race stops its session and adopts the winner instead of leaving an
  orphan counted against the user's stream limits.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(jellycompat): guard alias ambiguity and method switches in upstream attach

Second-round adversarial findings on the issue #244 fix:

- A client PlaySessionId alias now only resolves when it identifies
  exactly one live session for the token, and the candidate must agree
  with the report's ItemId/MediaSourceId. A reused or stale client id
  degrades to the route scan (progress still lands on the right item)
  instead of binding an arbitrary session, and an ambiguous Stopped
  report tears nothing down.
- A CAS loser only adopts the concurrent winner when it serves the same
  play method; a concurrent method switch surfaces 409 Conflict instead
  of continuing with mismatched transcode bookkeeping.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(jellycompat): key audio-selection updates by the resolved play session id

Review follow-up (PR #262): after an alias or route fallback,
req.PlaySessionID is the client's own generated id and not a playback
store key, so the audio-track-change branch in handlePlaybackReport
silently dropped selection updates for Static direct-play sessions.
Use the resolved playSession.ID for the store mutation and logging.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Quick
2026-07-02 11:37:32 -04:00
committed by GitHub
co-authored by Claude Fable 5
parent b250dbb59b
commit 287ff3658f
6 changed files with 868 additions and 35 deletions
+45 -6
View File
@@ -18,11 +18,14 @@ import (
)
type testCompatSessionManager struct {
sessions map[string]*playback.Session
audioTrackCalls []compatAudioTrackCall
progressCalls int
stopCalls []string
startCalls int
sessions map[string]*playback.Session
audioTrackCalls []compatAudioTrackCall
progressCalls int
progressUpdates []compatProgressCall
stopCalls []string
startCalls int
beginTransportCalls []string
endTransportCalls []string
}
type compatAudioTrackCall struct {
@@ -31,6 +34,12 @@ type compatAudioTrackCall struct {
method playback.PlayMethod
}
type compatProgressCall struct {
sessionID string
position float64
isPaused bool
}
func (m *testCompatSessionManager) StartSession(userID int, profileID string, fileID int, method playback.PlayMethod, transcodeAudio bool) (*playback.Session, error) {
m.startCalls++
session := &playback.Session{
@@ -49,8 +58,38 @@ func (m *testCompatSessionManager) StartSession(userID int, profileID string, fi
return session, nil
}
func (m *testCompatSessionManager) UpdateProgress(string, float64, bool) error {
func (m *testCompatSessionManager) UpdateProgress(sessionID string, position float64, isPaused bool) error {
m.progressCalls++
m.progressUpdates = append(m.progressUpdates, compatProgressCall{
sessionID: sessionID,
position: position,
isPaused: isPaused,
})
if m.sessions != nil {
if _, ok := m.sessions[sessionID]; !ok {
return playback.ErrSessionNotFound
}
}
return nil
}
func (m *testCompatSessionManager) BeginTransport(sessionID string) error {
m.beginTransportCalls = append(m.beginTransportCalls, sessionID)
if m.sessions != nil {
if _, ok := m.sessions[sessionID]; !ok {
return playback.ErrSessionNotFound
}
}
return nil
}
func (m *testCompatSessionManager) EndTransport(sessionID string) error {
m.endTransportCalls = append(m.endTransportCalls, sessionID)
if m.sessions != nil {
if _, ok := m.sessions[sessionID]; !ok {
return playback.ErrSessionNotFound
}
}
return nil
}
+4
View File
@@ -2963,6 +2963,10 @@ func max(value, fallback int) int {
}
func writeCompatUpstreamError(w http.ResponseWriter, err error) {
if errors.Is(err, errUpstreamReplaced) {
writeError(w, http.StatusConflict, "Conflict", "Playback session changed concurrently; retry the request")
return
}
if errors.Is(err, playback.ErrTooManyStreams) {
writeError(w, http.StatusTooManyRequests, "TooManyStreams", "Too many concurrent streams")
return
@@ -94,6 +94,8 @@ type SessionManagerInterface interface {
StopSession(sessionID string) error
GetSession(sessionID string) (*playback.Session, error)
SetTranscodeNodeURL(sessionID, url string) error
BeginTransport(sessionID string) error
EndTransport(sessionID string) error
}
type sessionStarterContext interface {
@@ -0,0 +1,602 @@
package jellycompat
import (
"context"
"errors"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"github.com/Silo-Server/silo-server/internal/playback"
)
// The tests in this file cover issue #244: an Infuse Static=true direct play
// vanished from the admin Activity View mid-playback with a frozen position.
// Three defects compounded:
//
// 1. handlePlaybackReport resolved the play session strictly by PlaySessionId,
// but a Static direct play that skipped PlaybackInfo reports under the
// client's own generated id, so every progress report silently no-op'd and
// the position never advanced.
// 2. With progress reports dropped, nothing refreshed the upstream session's
// activity, so stale cleanup reaped it ~45s in while the client kept
// streaming. HandleVideoStream never marked an in-flight transport the way
// the native stream handler does.
// 3. ensureUpstreamPlayback early-returned whenever UpstreamSessionID was set
// without checking the session still existed, so a reaped session was
// never recreated and the stream card never came back.
// newReportLivenessHandler builds a PlaybackHandler with one stored play
// session whose upstream session id is registered in the fake session manager.
func newReportLivenessHandler(upstreamID string, registerUpstream bool) (*PlaybackHandler, *testCompatSessionManager, string, string) {
codec := NewResourceIDCodec()
version := testCompatVersion()
source := testCompatSource(codec, version)
encodedItemID := codec.EncodeStringID(EncodedIDItem, "movie-1")
playbackStore := NewPlaybackSessionStore(time.Hour, nil)
playbackStore.Put(PlaybackSession{
ID: "play-1",
CompatToken: "token-1",
ItemID: "movie-1",
RouteItemID: encodedItemID,
UserID: "user-1",
UpstreamSessionID: upstreamID,
UpstreamPlayMethod: "direct",
MediaSources: []PlaybackMediaSource{source},
})
sessions := map[string]*playback.Session{}
if registerUpstream {
sessions[upstreamID] = &playback.Session{
ID: upstreamID,
PlayMethod: playback.PlayDirect,
BasePlayMethod: playback.PlayDirect,
}
}
sessionMgr := &testCompatSessionManager{sessions: sessions}
handler := &PlaybackHandler{
codec: codec,
playbackStore: playbackStore,
sessionMgr: sessionMgr,
transcodes: make(map[string]*playback.TranscodeSession),
}
return handler, sessionMgr, encodedItemID, source.ID
}
func postProgressReport(handler *PlaybackHandler, body string) *httptest.ResponseRecorder {
req := httptest.NewRequest(http.MethodPost, "/Sessions/Playing/Progress", strings.NewReader(body))
req = req.WithContext(context.WithValue(req.Context(), compatSessionKey,
&Session{Token: "token-1", StreamAppUserID: 1, ProfileID: "profile-1"}))
rec := httptest.NewRecorder()
handler.HandleSessionPlayingProgress(rec, req)
return rec
}
// TestHandlePlaybackReport_ClientPlaySessionIDFallsBackToItemRoute proves a
// progress report whose PlaySessionId is unknown to the server (the client
// generated it because Static=true direct play skipped PlaybackInfo) still
// reaches the upstream session via the item-route fallback — the same
// route-scoped reuse the stream path performs in resolvePlaybackRoute.
// Without the fallback the report is a silent 204 no-op, the Activity View
// position freezes, and stale cleanup later drops the live session.
func TestHandlePlaybackReport_ClientPlaySessionIDFallsBackToItemRoute(t *testing.T) {
handler, mgr, encodedItemID, sourceID := newReportLivenessHandler("upstream-1", true)
rec := postProgressReport(handler,
`{"PlaySessionId":"infuse-client-psid","ItemId":"`+encodedItemID+`","MediaSourceId":"`+sourceID+`","PositionTicks":1234500000}`)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
if len(mgr.progressUpdates) != 1 {
t.Fatalf("UpdateProgress calls = %d, want 1 (client-generated PlaySessionId must fall back to the item route)", len(mgr.progressUpdates))
}
got := mgr.progressUpdates[0]
if got.sessionID != "upstream-1" {
t.Fatalf("UpdateProgress session = %q, want upstream-1", got.sessionID)
}
if got.position != 123.45 {
t.Fatalf("UpdateProgress position = %v, want 123.45", got.position)
}
}
// TestHandlePlaybackReport_MediaSourceIDFallbackWhenItemUnknown covers clients
// that omit or send an unmatchable ItemId: the MediaSourceId in the report
// still identifies the play session.
func TestHandlePlaybackReport_MediaSourceIDFallbackWhenItemUnknown(t *testing.T) {
handler, mgr, _, sourceID := newReportLivenessHandler("upstream-1", true)
rec := postProgressReport(handler,
`{"PlaySessionId":"infuse-client-psid","MediaSourceId":"`+sourceID+`","PositionTicks":600000000}`)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
if len(mgr.progressUpdates) != 1 {
t.Fatalf("UpdateProgress calls = %d, want 1 (MediaSourceId must resolve the play session)", len(mgr.progressUpdates))
}
if got := mgr.progressUpdates[0].sessionID; got != "upstream-1" {
t.Fatalf("UpdateProgress session = %q, want upstream-1", got)
}
}
// TestHandlePlaybackReport_ForeignTokenDoesNotMatch ensures the fallback stays
// scoped to the caller's own compat token: a report authenticated with a
// different token must not read or touch another user's play session.
func TestHandlePlaybackReport_ForeignTokenDoesNotMatch(t *testing.T) {
handler, mgr, encodedItemID, sourceID := newReportLivenessHandler("upstream-1", true)
req := httptest.NewRequest(http.MethodPost, "/Sessions/Playing/Progress", strings.NewReader(
`{"PlaySessionId":"infuse-client-psid","ItemId":"`+encodedItemID+`","MediaSourceId":"`+sourceID+`","PositionTicks":600000000}`))
req = req.WithContext(context.WithValue(req.Context(), compatSessionKey,
&Session{Token: "other-token", StreamAppUserID: 2, ProfileID: "profile-2"}))
rec := httptest.NewRecorder()
handler.HandleSessionPlayingProgress(rec, req)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
if len(mgr.progressUpdates) != 0 {
t.Fatalf("UpdateProgress calls = %d, want 0 (foreign token must not bind another caller's session)", len(mgr.progressUpdates))
}
}
// TestHandlePlaybackReport_RevivesReapedUpstreamSession proves a progress
// report for a play session whose upstream session was reaped by stale cleanup
// recreates the upstream session instead of silently dropping the report, so
// the still-playing client reappears in the admin Activity View.
func TestHandlePlaybackReport_RevivesReapedUpstreamSession(t *testing.T) {
handler, mgr, _, sourceID := newReportLivenessHandler("upstream-reaped", false)
rec := postProgressReport(handler,
`{"PlaySessionId":"play-1","MediaSourceId":"`+sourceID+`","PositionTicks":9000000000}`)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
if mgr.startCalls != 1 {
t.Fatalf("StartSession calls = %d, want 1 (reaped upstream session must be recreated)", mgr.startCalls)
}
updated, ok := handler.playbackStore.Get("play-1")
if !ok {
t.Fatal("expected play session to remain in store")
}
if updated.UpstreamSessionID != "upstream-started" {
t.Fatalf("UpstreamSessionID = %q, want upstream-started", updated.UpstreamSessionID)
}
last := mgr.progressUpdates[len(mgr.progressUpdates)-1]
if last.sessionID != "upstream-started" || last.position != 900 {
t.Fatalf("last progress update = %+v, want session upstream-started at 900s", last)
}
}
// TestHandlePlaybackReport_StoppedDoesNotReviveUpstream ensures the revival
// path does not resurrect a session for a Stopped report: stopping playback of
// an already-reaped session must stay a no-op teardown, not create a ghost.
func TestHandlePlaybackReport_StoppedDoesNotReviveUpstream(t *testing.T) {
handler, mgr, _, sourceID := newReportLivenessHandler("upstream-reaped", false)
req := httptest.NewRequest(http.MethodPost, "/Sessions/Playing/Stopped", strings.NewReader(
`{"PlaySessionId":"play-1","MediaSourceId":"`+sourceID+`","PositionTicks":9000000000}`))
req = req.WithContext(context.WithValue(req.Context(), compatSessionKey,
&Session{Token: "token-1", StreamAppUserID: 1, ProfileID: "profile-1"}))
rec := httptest.NewRecorder()
handler.HandleSessionPlayingStopped(rec, req)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
if mgr.startCalls != 0 {
t.Fatalf("StartSession calls = %d, want 0 (Stopped must not revive a reaped session)", mgr.startCalls)
}
if _, ok := handler.playbackStore.Get("play-1"); ok {
t.Fatal("expected Stopped report to tear down the play session")
}
}
// TestEnsureUpstreamPlayback_RecreatesReapedSession proves the stream path
// heals a reaped upstream session: the next range request must recreate it
// rather than early-returning a dangling UpstreamSessionID forever.
func TestEnsureUpstreamPlayback_RecreatesReapedSession(t *testing.T) {
handler, mgr, _, _ := newReportLivenessHandler("upstream-reaped", false)
playSession, _ := handler.playbackStore.Get("play-1")
source := playSession.MediaSources[0]
got, err := handler.ensureUpstreamPlayback(context.Background(),
&Session{Token: "token-1", StreamAppUserID: 1, ProfileID: "profile-1"},
"play-1", source, "direct")
if err != nil {
t.Fatalf("ensureUpstreamPlayback: %v", err)
}
if mgr.startCalls != 1 {
t.Fatalf("StartSession calls = %d, want 1 (reaped upstream must be recreated)", mgr.startCalls)
}
if got.UpstreamSessionID != "upstream-started" {
t.Fatalf("UpstreamSessionID = %q, want upstream-started", got.UpstreamSessionID)
}
}
// TestEnsureUpstreamPlayback_ReusesLiveSession guards the reuse path: a live
// upstream session must not be torn down or recreated by subsequent requests.
func TestEnsureUpstreamPlayback_ReusesLiveSession(t *testing.T) {
handler, mgr, _, _ := newReportLivenessHandler("upstream-1", true)
playSession, _ := handler.playbackStore.Get("play-1")
source := playSession.MediaSources[0]
got, err := handler.ensureUpstreamPlayback(context.Background(),
&Session{Token: "token-1", StreamAppUserID: 1, ProfileID: "profile-1"},
"play-1", source, "direct")
if err != nil {
t.Fatalf("ensureUpstreamPlayback: %v", err)
}
if mgr.startCalls != 0 {
t.Fatalf("StartSession calls = %d, want 0 (live upstream must be reused)", mgr.startCalls)
}
if got.UpstreamSessionID != "upstream-1" {
t.Fatalf("UpstreamSessionID = %q, want upstream-1", got.UpstreamSessionID)
}
}
// TestHandlePlaybackReport_AliasBindsDeterministicallyAmongDuplicates proves
// that when two play sessions for the same item live under one compat token,
// a report carrying the client's own PlaySessionId binds the session that
// recorded it as an alias — not whichever the route scan happens to hit first.
func TestHandlePlaybackReport_AliasBindsDeterministicallyAmongDuplicates(t *testing.T) {
handler, mgr, encodedItemID, sourceID := newReportLivenessHandler("upstream-1", true)
// A second live play of the same item under the same token, no alias.
other, _ := handler.playbackStore.Get("play-1")
sibling := *other
sibling.ID = "play-2"
sibling.UpstreamSessionID = "upstream-2"
handler.playbackStore.Put(sibling)
mgr.sessions["upstream-2"] = &playback.Session{ID: "upstream-2", PlayMethod: playback.PlayDirect}
// The stream path recorded the client's PlaySessionId on play-1.
if err := handler.playbackStore.Update("play-1", func(current *PlaybackSession) error {
current.ClientPlaySessionID = "infuse-client-psid"
return nil
}); err != nil {
t.Fatal(err)
}
rec := postProgressReport(handler,
`{"PlaySessionId":"infuse-client-psid","ItemId":"`+encodedItemID+`","MediaSourceId":"`+sourceID+`","PositionTicks":1234500000}`)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
if len(mgr.progressUpdates) != 1 {
t.Fatalf("UpdateProgress calls = %d, want 1", len(mgr.progressUpdates))
}
if got := mgr.progressUpdates[0].sessionID; got != "upstream-1" {
t.Fatalf("UpdateProgress session = %q, want upstream-1 (the aliased session)", got)
}
}
// TestHandleVideoStream_StaticRecordsClientPlaySessionAlias proves the stream
// path records the client-generated PlaySessionId so subsequent playback
// reports resolve the session without relying on ItemId/route matching.
func TestHandleVideoStream_StaticRecordsClientPlaySessionAlias(t *testing.T) {
handler, encodedID, _ := newStaticDirectPlayHandler(t)
rec := serveStaticStream(handler, encodedID, "Static=true&PlaySessionId=infuse-client-psid")
if rec.Code != 200 {
t.Fatalf("expected status 200; got %d, body=%s", rec.Code, rec.Body.String())
}
playSession, ok := handler.playbackStore.FindByClientPlaySessionID("token-1", "infuse-client-psid")
if !ok {
t.Fatal("expected the static play session to record the client PlaySessionId alias")
}
if playSession.UpstreamSessionID != "upstream-started" {
t.Fatalf("UpstreamSessionID = %q, want upstream-started", playSession.UpstreamSessionID)
}
// The full Infuse loop: a progress report carrying only the client id
// (no ItemId) must reach the upstream session via the alias.
mgr := handler.sessionMgr.(*testCompatSessionManager)
rep := postProgressReport(handler, `{"PlaySessionId":"infuse-client-psid","PositionTicks":600000000}`)
if rep.Code != http.StatusNoContent {
t.Fatalf("report status = %d", rep.Code)
}
if len(mgr.progressUpdates) != 1 || mgr.progressUpdates[0].sessionID != "upstream-started" {
t.Fatalf("progress updates = %+v, want one update on upstream-started", mgr.progressUpdates)
}
}
// TestHandlePlaybackReport_StopViaAliasTearsDown proves a Stopped report
// resolved through the recorded alias still tears the play session down —
// the alias is an exact, caller-owned match.
func TestHandlePlaybackReport_StopViaAliasTearsDown(t *testing.T) {
handler, mgr, _, sourceID := newReportLivenessHandler("upstream-1", true)
if err := handler.playbackStore.Update("play-1", func(current *PlaybackSession) error {
current.ClientPlaySessionID = "infuse-client-psid"
return nil
}); err != nil {
t.Fatal(err)
}
req := httptest.NewRequest(http.MethodPost, "/Sessions/Playing/Stopped", strings.NewReader(
`{"PlaySessionId":"infuse-client-psid","MediaSourceId":"`+sourceID+`","PositionTicks":9000000000}`))
req = req.WithContext(context.WithValue(req.Context(), compatSessionKey,
&Session{Token: "token-1", StreamAppUserID: 1, ProfileID: "profile-1"}))
rec := httptest.NewRecorder()
handler.HandleSessionPlayingStopped(rec, req)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
if len(mgr.stopCalls) != 1 || mgr.stopCalls[0] != "upstream-1" {
t.Fatalf("StopSession calls = %v, want exactly one for upstream-1", mgr.stopCalls)
}
if _, ok := handler.playbackStore.Get("play-1"); ok {
t.Fatal("expected alias-matched Stopped report to tear down the play session")
}
}
// TestHandlePlaybackReport_StopViaRouteMatchDoesNotTearDown proves a Stopped
// report that only matched by item/source route (an ambiguous match when the
// same item plays twice under one token) must not tear down the session it
// happened to hit; stale cleanup owns that session's end of life.
func TestHandlePlaybackReport_StopViaRouteMatchDoesNotTearDown(t *testing.T) {
handler, mgr, encodedItemID, sourceID := newReportLivenessHandler("upstream-1", true)
req := httptest.NewRequest(http.MethodPost, "/Sessions/Playing/Stopped", strings.NewReader(
`{"PlaySessionId":"never-seen-psid","ItemId":"`+encodedItemID+`","MediaSourceId":"`+sourceID+`","PositionTicks":9000000000}`))
req = req.WithContext(context.WithValue(req.Context(), compatSessionKey,
&Session{Token: "token-1", StreamAppUserID: 1, ProfileID: "profile-1"}))
rec := httptest.NewRecorder()
handler.HandleSessionPlayingStopped(rec, req)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
if len(mgr.stopCalls) != 0 {
t.Fatalf("StopSession calls = %v, want none for an ambiguous route-only match", mgr.stopCalls)
}
if _, ok := handler.playbackStore.Get("play-1"); !ok {
t.Fatal("expected route-only Stopped report to leave the play session in place")
}
// The final position still lands on the upstream session.
if len(mgr.progressUpdates) != 1 || mgr.progressUpdates[0].sessionID != "upstream-1" {
t.Fatalf("progress updates = %+v, want one update on upstream-1", mgr.progressUpdates)
}
}
// TestEnsureUpstreamPlayback_ReviveClosesStaleTranscode proves that when a
// reaped upstream session is recreated under the same play method, any
// transcode still keyed to the stale upstream id is closed first — otherwise
// a second ffmpeg would start alongside the orphaned one.
func TestEnsureUpstreamPlayback_ReviveClosesStaleTranscode(t *testing.T) {
handler, mgr, _, _ := newReportLivenessHandler("upstream-reaped", false)
handler.transcodes["upstream-reaped"] = nil // stale entry keyed by the reaped id
playSession, _ := handler.playbackStore.Get("play-1")
source := playSession.MediaSources[0]
if _, err := handler.ensureUpstreamPlayback(context.Background(),
&Session{Token: "token-1", StreamAppUserID: 1, ProfileID: "profile-1"},
"play-1", source, "direct"); err != nil {
t.Fatalf("ensureUpstreamPlayback: %v", err)
}
if mgr.startCalls != 1 {
t.Fatalf("StartSession calls = %d, want 1", mgr.startCalls)
}
if _, stale := handler.transcodes["upstream-reaped"]; stale {
t.Fatal("expected the stale transcode entry to be closed before recreating the upstream session")
}
}
// TestHandlePlaybackReport_AliasResolvedAudioSelectionApplies proves an
// audio-track change carried by an alias-resolved report is applied: store
// mutations must key on the resolved play session id, not the client's own
// PlaySessionId (which is not a store key for Static direct play).
func TestHandlePlaybackReport_AliasResolvedAudioSelectionApplies(t *testing.T) {
handler, mgr, _, sourceID := newReportLivenessHandler("upstream-1", true)
if err := handler.playbackStore.Update("play-1", func(current *PlaybackSession) error {
current.ClientPlaySessionID = "infuse-client-psid"
return nil
}); err != nil {
t.Fatal(err)
}
// testCompatSource pre-selects audio stream index 2; switch to index 1.
rec := postProgressReport(handler,
`{"PlaySessionId":"infuse-client-psid","MediaSourceId":"`+sourceID+`","AudioStreamIndex":1,"PositionTicks":600000000}`)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
updated, ok := handler.playbackStore.Get("play-1")
if !ok {
t.Fatal("expected play session in store")
}
if updated.MediaSources[0].SelectedAudioStreamIndex == nil || *updated.MediaSources[0].SelectedAudioStreamIndex != 1 {
t.Fatalf("SelectedAudioStreamIndex = %v, want 1 (audio change must apply to the resolved session)", updated.MediaSources[0].SelectedAudioStreamIndex)
}
if len(mgr.audioTrackCalls) != 1 {
t.Fatalf("upstream audio track calls = %d, want 1", len(mgr.audioTrackCalls))
}
}
// TestHandlePlaybackReport_DuplicateAliasFallsBackToRoute proves a client that
// reuses one PlaySessionId across different items cannot misbind reports: the
// ambiguous alias is skipped and the report resolves by ItemId route instead,
// and a Stopped report with only the ambiguous alias (no ItemId) tears nothing
// down.
func TestHandlePlaybackReport_DuplicateAliasFallsBackToRoute(t *testing.T) {
handler, mgr, encodedItemID, sourceID := newReportLivenessHandler("upstream-1", true)
// Same client alias on a second live session for a DIFFERENT item.
otherItemID := handler.codec.EncodeStringID(EncodedIDItem, "movie-2")
handler.playbackStore.Put(PlaybackSession{
ID: "play-2",
CompatToken: "token-1",
ItemID: "movie-2",
RouteItemID: otherItemID,
ClientPlaySessionID: "reused-psid",
UpstreamSessionID: "upstream-2",
UpstreamPlayMethod: "direct",
})
mgr.sessions["upstream-2"] = &playback.Session{ID: "upstream-2", PlayMethod: playback.PlayDirect}
if err := handler.playbackStore.Update("play-1", func(current *PlaybackSession) error {
current.ClientPlaySessionID = "reused-psid"
return nil
}); err != nil {
t.Fatal(err)
}
// Progress with the duplicate alias plus ItemId: must bind play-1 by route.
rec := postProgressReport(handler,
`{"PlaySessionId":"reused-psid","ItemId":"`+encodedItemID+`","MediaSourceId":"`+sourceID+`","PositionTicks":600000000}`)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
if len(mgr.progressUpdates) != 1 || mgr.progressUpdates[0].sessionID != "upstream-1" {
t.Fatalf("progress updates = %+v, want one update on upstream-1 via item route", mgr.progressUpdates)
}
// Stopped with only the duplicate alias: ambiguous, must tear nothing down.
req := httptest.NewRequest(http.MethodPost, "/Sessions/Playing/Stopped",
strings.NewReader(`{"PlaySessionId":"reused-psid","PositionTicks":600000000}`))
req = req.WithContext(context.WithValue(req.Context(), compatSessionKey,
&Session{Token: "token-1", StreamAppUserID: 1, ProfileID: "profile-1"}))
rr := httptest.NewRecorder()
handler.HandleSessionPlayingStopped(rr, req)
if len(mgr.stopCalls) != 0 {
t.Fatalf("StopSession calls = %v, want none for an ambiguous duplicate alias", mgr.stopCalls)
}
}
// TestHandlePlaybackReport_AliasContradictedByItemFallsBack proves an alias
// whose session disagrees with the report's ItemId is rejected (a stale or
// reused client id) and the report resolves by route instead.
func TestHandlePlaybackReport_AliasContradictedByItemFallsBack(t *testing.T) {
handler, mgr, _, _ := newReportLivenessHandler("upstream-1", true)
// Alias points at play-1 (movie-1), but the report is about movie-2.
if err := handler.playbackStore.Update("play-1", func(current *PlaybackSession) error {
current.ClientPlaySessionID = "stale-psid"
return nil
}); err != nil {
t.Fatal(err)
}
otherItemID := handler.codec.EncodeStringID(EncodedIDItem, "movie-2")
handler.playbackStore.Put(PlaybackSession{
ID: "play-2",
CompatToken: "token-1",
ItemID: "movie-2",
RouteItemID: otherItemID,
UpstreamSessionID: "upstream-2",
UpstreamPlayMethod: "direct",
})
mgr.sessions["upstream-2"] = &playback.Session{ID: "upstream-2", PlayMethod: playback.PlayDirect}
rec := postProgressReport(handler,
`{"PlaySessionId":"stale-psid","ItemId":"`+otherItemID+`","PositionTicks":600000000}`)
if rec.Code != http.StatusNoContent {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
if len(mgr.progressUpdates) != 1 || mgr.progressUpdates[0].sessionID != "upstream-2" {
t.Fatalf("progress updates = %+v, want one update on upstream-2 (the reported item)", mgr.progressUpdates)
}
}
// racingStartSessionManager simulates a concurrent request winning the
// upstream-attach race: while this caller is inside StartSession, the store's
// play session is re-pointed at a different upstream id/method.
type racingStartSessionManager struct {
testCompatSessionManager
store *PlaybackSessionStore
winnerMethod string
}
func (m *racingStartSessionManager) StartSession(userID int, profileID string, fileID int, method playback.PlayMethod, transcodeAudio bool) (*playback.Session, error) {
_ = m.store.Update("play-1", func(current *PlaybackSession) error {
current.UpstreamSessionID = "upstream-winner"
current.UpstreamPlayMethod = m.winnerMethod
return nil
})
return m.testCompatSessionManager.StartSession(userID, profileID, fileID, method, transcodeAudio)
}
// TestEnsureUpstreamPlayback_CASLoserAdoptsSameMethodWinner proves the loser
// of a concurrent attach race stops its own freshly created session and
// adopts the winner when the winner serves the same play method.
func TestEnsureUpstreamPlayback_CASLoserAdoptsSameMethodWinner(t *testing.T) {
handler, _, _, _ := newReportLivenessHandler("", false)
base := handler.sessionMgr.(*testCompatSessionManager)
racer := &racingStartSessionManager{
testCompatSessionManager: *base,
store: handler.playbackStore,
winnerMethod: "direct",
}
handler.sessionMgr = racer
playSession, _ := handler.playbackStore.Get("play-1")
source := playSession.MediaSources[0]
got, err := handler.ensureUpstreamPlayback(context.Background(),
&Session{Token: "token-1", StreamAppUserID: 1, ProfileID: "profile-1"},
"play-1", source, "direct")
if err != nil {
t.Fatalf("ensureUpstreamPlayback: %v", err)
}
if got.UpstreamSessionID != "upstream-winner" {
t.Fatalf("UpstreamSessionID = %q, want the concurrent winner", got.UpstreamSessionID)
}
if len(racer.stopCalls) != 1 || racer.stopCalls[0] != "upstream-started" {
t.Fatalf("StopSession calls = %v, want rollback of the loser's session", racer.stopCalls)
}
}
// TestEnsureUpstreamPlayback_CASLoserRejectsMethodSwitch proves the loser does
// NOT adopt a winner running a different play method — it rolls back its
// session and surfaces a conflict instead of continuing with mismatched
// transcode bookkeeping.
func TestEnsureUpstreamPlayback_CASLoserRejectsMethodSwitch(t *testing.T) {
handler, _, _, _ := newReportLivenessHandler("", false)
base := handler.sessionMgr.(*testCompatSessionManager)
racer := &racingStartSessionManager{
testCompatSessionManager: *base,
store: handler.playbackStore,
winnerMethod: "transcode",
}
handler.sessionMgr = racer
playSession, _ := handler.playbackStore.Get("play-1")
source := playSession.MediaSources[0]
_, err := handler.ensureUpstreamPlayback(context.Background(),
&Session{Token: "token-1", StreamAppUserID: 1, ProfileID: "profile-1"},
"play-1", source, "direct")
if !errors.Is(err, errUpstreamReplaced) {
t.Fatalf("error = %v, want errUpstreamReplaced", err)
}
if len(racer.stopCalls) != 1 || racer.stopCalls[0] != "upstream-started" {
t.Fatalf("StopSession calls = %v, want rollback of the loser's session", racer.stopCalls)
}
}
// TestHandleVideoStream_DirectPlayMarksTransport proves the compat stream
// handler marks an in-flight media transport on the upstream session while
// serving, mirroring the native stream handler. Without it, a long-lived
// direct-play range transfer emits no activity and stale cleanup reaps the
// session mid-transfer.
func TestHandleVideoStream_DirectPlayMarksTransport(t *testing.T) {
handler, encodedID, body := newStaticDirectPlayHandler(t)
mgr := handler.sessionMgr.(*testCompatSessionManager)
rec := serveStaticStream(handler, encodedID, "Static=true")
if rec.Code != 200 {
t.Fatalf("expected status 200; got %d, body=%s", rec.Code, rec.Body.String())
}
if got := rec.Body.String(); got != body {
t.Fatalf("expected file content %q; got %q", body, got)
}
if len(mgr.beginTransportCalls) != 1 || mgr.beginTransportCalls[0] != "upstream-started" {
t.Fatalf("BeginTransport calls = %v, want exactly one for upstream-started", mgr.beginTransportCalls)
}
if len(mgr.endTransportCalls) != 1 || mgr.endTransportCalls[0] != "upstream-started" {
t.Fatalf("EndTransport calls = %v, want exactly one for upstream-started", mgr.endTransportCalls)
}
}
+46 -6
View File
@@ -10,11 +10,16 @@ import (
// PlaybackSession stores compat-owned playback negotiation state before the
// native Silo playback session starts.
type PlaybackSession struct {
ID string
CompatToken string
ItemID string
RouteItemID string
UserID string
ID string
CompatToken string
ItemID string
RouteItemID string
// ClientPlaySessionID records the client's own generated PlaySessionId
// when it differs from ours (Static=true direct play skips PlaybackInfo,
// so the client never learns the server id). Playback reports carrying
// that id resolve to this session directly instead of by ambiguous route.
ClientPlaySessionID string
UserID string
InitialSeekSeconds float64
MediaSources []PlaybackMediaSource
UpstreamSessionID string
@@ -123,6 +128,38 @@ func (s *PlaybackSessionStore) Update(id string, fn func(*PlaybackSession) error
return nil
}
// FindByClientPlaySessionID resolves the client-generated PlaySessionId alias
// recorded for plays that skipped PlaybackInfo (Static=true direct play). The
// alias must identify exactly one live session: a client that reuses one
// PlaySessionId across plays makes the alias ambiguous, and the caller should
// fall back to route matching instead of binding an arbitrary session.
func (s *PlaybackSessionStore) FindByClientPlaySessionID(compatToken, clientPlaySessionID string) (*PlaybackSession, bool) {
if clientPlaySessionID == "" {
return nil, false
}
s.mu.RLock()
defer s.mu.RUnlock()
now := s.now()
var match *PlaybackSession
for _, session := range s.sessions {
if !session.ExpiresAt.After(now) {
continue
}
if session.CompatToken != compatToken {
continue
}
if session.ClientPlaySessionID == clientPlaySessionID {
if match != nil {
return nil, false
}
cp := session
match = &cp
}
}
return match, match != nil
}
// FindByRoute resolves a route item/media-source identifier to a compat playback session.
func (s *PlaybackSessionStore) FindByRoute(compatToken, routeID string) (*PlaybackSession, *PlaybackMediaSource, bool) {
s.mu.RLock()
@@ -136,7 +173,10 @@ func (s *PlaybackSessionStore) FindByRoute(compatToken, routeID string) (*Playba
if compatToken != "" && session.CompatToken != compatToken {
continue
}
if session.RouteItemID == routeID {
// UUID-normalized comparison: playback reports echo the item id in
// whatever casing/dash format the client model uses, which may differ
// from the raw route param captured at stream time.
if mediaSourceIDsEqual(session.RouteItemID, routeID) {
cp := session
return &cp, nil, true
}
+169 -23
View File
@@ -32,6 +32,10 @@ import (
// the near-head follow-up segments arrive quickly enough for browser playback.
const compatSegmentDuration = 2
// errUpstreamReplaced signals that a concurrent request attached a different
// upstream session to the play session while this one was being created.
var errUpstreamReplaced = errors.New("upstream session replaced concurrently")
type sessionReportRequest struct {
ItemID string `json:"ItemId"`
MediaSourceID string `json:"MediaSourceId"`
@@ -60,7 +64,8 @@ func (h *PlaybackHandler) HandleVideoStream(w http.ResponseWriter, r *http.Reque
// lookup must be case-insensitive: SenPlayer sends "static=true"
// (lowercase) and a case-sensitive Get("Static") would miss it, dropping
// the client to a 404 "Playback session not found" on every direct play.
playSession, source, err = h.createStaticPlaySession(r.Context(), session, routeID, mediaSourceID)
clientPlaySessionID := newCaseInsensitiveQuery(r.URL.Query()).Get("PlaySessionId")
playSession, source, err = h.createStaticPlaySession(r.Context(), session, routeID, mediaSourceID, clientPlaySessionID)
}
if err != nil {
writeError(w, http.StatusNotFound, "NotFound", "Playback session not found")
@@ -109,6 +114,18 @@ func (h *PlaybackHandler) HandleVideoStream(w http.ResponseWriter, r *http.Reque
}
}
// Mark an in-flight media transport, mirroring the native stream handler:
// a long-lived direct-play range transfer emits no progress reports, and
// without the transport marker stale cleanup reaps the session mid-stream.
if h.sessionMgr != nil && playSession.UpstreamSessionID != "" {
if err := h.sessionMgr.BeginTransport(playSession.UpstreamSessionID); err == nil {
upstreamSessionID := playSession.UpstreamSessionID
defer func() {
_ = h.sessionMgr.EndTransport(upstreamSessionID)
}()
}
}
switch method {
case "remux":
audioTrackIndex := -1
@@ -790,7 +807,35 @@ func (h *PlaybackHandler) handlePlaybackReport(w http.ResponseWriter, r *http.Re
}
playSession, ok := h.playbackStore.Get(req.PlaySessionID)
if !ok || playSession.CompatToken != session.Token || playSession.UpstreamSessionID == "" {
if ok && playSession.CompatToken != session.Token {
playSession, ok = nil, false
}
matchedByRouteOnly := false
if !ok {
// Static=true direct play (Infuse, SenPlayer) skips PlaybackInfo, so the
// client reports progress under its own generated PlaySessionId. The
// stream path recorded that id as an alias on the play session it
// bound; resolve by the alias first, then fall back to the same
// route-scoped lookup the stream path uses (see resolvePlaybackRoute).
// Without either, these reports silently no-op, the admin activity view
// position freezes, and stale cleanup drops the still-active session.
playSession, ok = h.playbackStore.FindByClientPlaySessionID(session.Token, req.PlaySessionID)
if ok && !reportMatchesPlaySession(playSession, req) {
playSession, ok = nil, false
}
}
if !ok {
for _, routeID := range []string{req.ItemID, req.MediaSourceID} {
if routeID == "" {
continue
}
if playSession, _, ok = h.playbackStore.FindByRoute(session.Token, routeID); ok {
matchedByRouteOnly = true
break
}
}
}
if !ok || playSession.UpstreamSessionID == "" {
w.WriteHeader(http.StatusNoContent)
return
}
@@ -804,7 +849,10 @@ func (h *PlaybackHandler) handlePlaybackReport(w http.ResponseWriter, r *http.Re
// causes an hls.js retry loop. Only act when the index actually changes.
if req.AudioStreamIndex != nil && audioSelectionChanged(playSession, req.MediaSourceID, int(*req.AudioStreamIndex)) {
selectedAudioStreamIndex := int(*req.AudioStreamIndex)
updatedPlaySession, updatedSource, updateErr := h.setSelectedAudioStream(req.PlaySessionID, req.MediaSourceID, selectedAudioStreamIndex)
// Key store mutations by the resolved session id: after an alias or
// route fallback, req.PlaySessionID is the client's own id and is not
// a store key.
updatedPlaySession, updatedSource, updateErr := h.setSelectedAudioStream(playSession.ID, req.MediaSourceID, selectedAudioStreamIndex)
if updateErr == nil {
playSession = updatedPlaySession
if resolvedAudioTrackIndex, ok := compatAudioTrackIndex(*updatedSource); ok {
@@ -812,7 +860,7 @@ func (h *PlaybackHandler) handlePlaybackReport(w http.ResponseWriter, r *http.Re
}
if syncErr := h.syncUpstreamAudioSelection(playSession, *updatedSource); syncErr != nil {
slog.Warn("jellycompat audio selection sync failed",
"play_session_id", req.PlaySessionID,
"play_session_id", playSession.ID,
"upstream_session_id", playSession.UpstreamSessionID,
"error", syncErr,
)
@@ -820,14 +868,14 @@ func (h *PlaybackHandler) handlePlaybackReport(w http.ResponseWriter, r *http.Re
restarted, restartErr := h.restartCompatTranscodeForAudioSelection(r.Context(), playSession, *updatedSource, positionSeconds)
if restartErr != nil {
slog.Warn("jellycompat audio selection restart failed",
"play_session_id", req.PlaySessionID,
"play_session_id", playSession.ID,
"upstream_session_id", playSession.UpstreamSessionID,
"error", restartErr,
)
}
audioRestarted = restarted
slog.Info("jellycompat audio selection updated",
"play_session_id", req.PlaySessionID,
"play_session_id", playSession.ID,
"media_source_id", updatedSource.ID,
"audio_stream_index", selectedAudioStreamIndex,
"audio_track_index", audioTrackIndex,
@@ -836,7 +884,17 @@ func (h *PlaybackHandler) handlePlaybackReport(w http.ResponseWriter, r *http.Re
}
}
if positionSeconds > 0 && h.sessionMgr != nil {
_ = h.sessionMgr.UpdateProgress(playSession.UpstreamSessionID, positionSeconds, req.IsPaused)
err := h.sessionMgr.UpdateProgress(playSession.UpstreamSessionID, positionSeconds, req.IsPaused)
if errors.Is(err, playback.ErrSessionNotFound) && !stop {
// The upstream session was reaped as stale (e.g. the client buffered
// far ahead and went quiet between range requests). The report proves
// the client is still playing, so recreate the session instead of
// dropping it from session tracking for the rest of playback.
if revived := h.reviveUpstreamForReport(r.Context(), session, playSession, req.MediaSourceID); revived != nil {
playSession = revived
_ = h.sessionMgr.UpdateProgress(playSession.UpstreamSessionID, positionSeconds, req.IsPaused)
}
}
}
// Persist progress to user store
if positionSeconds > 0 && h.storeProvider != nil && playSession.ItemID != "" {
@@ -854,21 +912,77 @@ func (h *PlaybackHandler) handlePlaybackReport(w http.ResponseWriter, r *http.Re
}
}
}
if stop {
if stop && !matchedByRouteOnly {
// A bare item/source route match is ambiguous when the same item plays
// twice under one token, so never tear down a session the report may
// not own. A session that really stopped emits no further reports or
// transport, so stale cleanup reaps it shortly anyway.
h.teardownPlaySession(r.Context(), playSession)
}
w.WriteHeader(http.StatusNoContent)
}
// reportMatchesPlaySession rejects an alias-resolved session whose item or
// media source contradicts the report, so a stale or reused client id cannot
// route a report (or its teardown) to the wrong play.
func reportMatchesPlaySession(playSession *PlaybackSession, req sessionReportRequest) bool {
if req.ItemID != "" && !mediaSourceIDsEqual(playSession.RouteItemID, req.ItemID) {
return false
}
if req.MediaSourceID != "" && findMediaSource(playSession, req.MediaSourceID) == nil {
return false
}
return true
}
// reviveUpstreamForReport recreates the upstream playback session backing a
// progress report after stale cleanup reaped it. Returns nil when the play
// session has no usable media source or the recreation fails.
func (h *PlaybackHandler) reviveUpstreamForReport(ctx context.Context, session *Session, playSession *PlaybackSession, mediaSourceID string) *PlaybackSession {
if playSession.UpstreamPlayMethod == "" {
return nil
}
source := findMediaSource(playSession, mediaSourceID)
if source == nil {
source = firstMediaSource(playSession)
}
if source == nil {
return nil
}
revived, err := h.ensureUpstreamPlayback(ctx, session, playSession.ID, *source, playSession.UpstreamPlayMethod)
if err != nil {
slog.Warn("jellycompat upstream session revive failed",
"play_session_id", playSession.ID,
"upstream_session_id", playSession.UpstreamSessionID,
"error", err,
)
return nil
}
return revived
}
func (h *PlaybackHandler) ensureUpstreamPlayback(ctx context.Context, compatSession *Session, playSessionID string, source PlaybackMediaSource, method string) (*PlaybackSession, error) {
playSession, ok := h.playbackStore.Get(playSessionID)
if !ok {
return nil, ErrSessionNotFound
}
observedUpstreamID := playSession.UpstreamSessionID
if playSession.UpstreamSessionID != "" && playSession.UpstreamPlayMethod == method {
_ = h.syncUpstreamAudioSelection(playSession, source)
return playSession, nil
upstreamLive := true
if h.sessionMgr != nil {
_, err := h.sessionMgr.GetSession(playSession.UpstreamSessionID)
upstreamLive = err == nil
}
if upstreamLive {
_ = h.syncUpstreamAudioSelection(playSession, source)
return playSession, nil
}
// The upstream session was reaped as stale while the client kept
// playing; recreate it under the same play session. Any transcode
// still keyed to the stale id must go first, or a second ffmpeg
// would start alongside it.
h.closeTranscodeSession(playSession.UpstreamSessionID, "")
}
if h.sessionMgr == nil {
@@ -914,13 +1028,31 @@ func (h *PlaybackHandler) ensureUpstreamPlayback(ctx context.Context, compatSess
UpstreamSessionID: session.ID,
UpstreamPlayMethod: method,
}, source)
if err := h.playbackStore.Update(playSessionID, func(current *PlaybackSession) error {
// Attach the new upstream session only if no concurrent request replaced
// the one we observed (range requests race with progress-report revives).
// The loser stops its session instead of leaving an orphan that counts
// toward the user's stream limits until stale cleanup.
if updateErr := h.playbackStore.Update(playSessionID, func(current *PlaybackSession) error {
if current.UpstreamSessionID != observedUpstreamID {
return errUpstreamReplaced
}
current.UpstreamSessionID = session.ID
current.UpstreamPlayMethod = method
current.TranscodeStarted = false
return nil
}); err != nil {
return nil, err
}); updateErr != nil {
_ = h.sessionMgr.StopSession(session.ID)
if errors.Is(updateErr, errUpstreamReplaced) {
// Adopt the winner only when it serves the same play method;
// otherwise a concurrent method switch made this caller's
// negotiated stream obsolete — surface the conflict rather than
// continuing on a session with mismatched transcode bookkeeping.
if winner, ok := h.playbackStore.Get(playSessionID); ok && winner.UpstreamPlayMethod == method {
return winner, nil
}
return nil, errUpstreamReplaced
}
return nil, updateErr
}
updated, ok := h.playbackStore.Get(playSessionID)
if !ok {
@@ -1174,8 +1306,10 @@ func (h *PlaybackHandler) compatSegmentDuration() int {
}
// createStaticPlaySession builds an on-the-fly play session for Infuse-style
// Static=true direct play requests that skip PlaybackInfo.
func (h *PlaybackHandler) createStaticPlaySession(ctx context.Context, session *Session, routeID, mediaSourceID string) (*PlaybackSession, *PlaybackMediaSource, error) {
// Static=true direct play requests that skip PlaybackInfo. clientPlaySessionID
// is the client's own PlaySessionId (if it sent one) so later playback reports
// carrying it can resolve this session directly.
func (h *PlaybackHandler) createStaticPlaySession(ctx context.Context, session *Session, routeID, mediaSourceID, clientPlaySessionID string) (*PlaybackSession, *PlaybackMediaSource, error) {
contentID, err := decodeContentID(h.codec, routeID)
if err != nil {
return nil, nil, ErrSessionNotFound
@@ -1194,12 +1328,13 @@ func (h *PlaybackHandler) createStaticPlaySession(ctx context.Context, session *
}
ps := &PlaybackSession{
ID: playSessionID,
CompatToken: session.Token,
ItemID: detail.ContentID,
RouteItemID: routeID,
UserID: session.PseudoUserID.String(),
MediaSources: sources,
ID: playSessionID,
CompatToken: session.Token,
ItemID: detail.ContentID,
RouteItemID: routeID,
ClientPlaySessionID: clientPlaySessionID,
UserID: session.PseudoUserID.String(),
MediaSources: sources,
}
h.playbackStore.Put(*ps)
@@ -1214,8 +1349,9 @@ func (h *PlaybackHandler) createStaticPlaySession(ctx context.Context, session *
}
func (h *PlaybackHandler) resolvePlaybackRoute(r *http.Request, compatSession *Session, routeID, mediaSourceID string) (*PlaybackSession, *PlaybackMediaSource, error) {
if playSessionID := newCaseInsensitiveQuery(r.URL.Query()).Get("PlaySessionId"); playSessionID != "" {
if playSession, ok := h.playbackStore.Get(playSessionID); ok && playSession.CompatToken == compatSession.Token {
clientPlaySessionID := newCaseInsensitiveQuery(r.URL.Query()).Get("PlaySessionId")
if clientPlaySessionID != "" {
if playSession, ok := h.playbackStore.Get(clientPlaySessionID); ok && playSession.CompatToken == compatSession.Token {
if mediaSourceID != "" {
source := findMediaSource(playSession, mediaSourceID)
return playSession, source, nil
@@ -1236,6 +1372,16 @@ func (h *PlaybackHandler) resolvePlaybackRoute(r *http.Request, compatSession *S
if !ok {
return nil, nil, ErrSessionNotFound
}
if clientPlaySessionID != "" && playSession.ClientPlaySessionID != clientPlaySessionID {
// Remember the client's own PlaySessionId so playback reports carrying
// it resolve to this session directly instead of by ambiguous route.
if h.playbackStore.Update(playSession.ID, func(current *PlaybackSession) error {
current.ClientPlaySessionID = clientPlaySessionID
return nil
}) == nil {
playSession.ClientPlaySessionID = clientPlaySessionID
}
}
if source == nil && mediaSourceID != "" {
source = findMediaSource(playSession, mediaSourceID)
}