From 487dc8482901dcef3572ffb56ee767c30952205c Mon Sep 17 00:00:00 2001 From: Quick104 <31828688+Quick104@users.noreply.github.com> Date: Wed, 22 Jul 2026 22:01:02 -0400 Subject: [PATCH] fix(jellycompat): deduplicate playback negotiations - Replace unstarted negotiations for the same device and item - Apply deduplication atomically across durable store instances --- internal/jellycompat/handlers_playback.go | 7 +- .../playback_negotiation_dedup_test.go | 129 ++++++++++++++++++ internal/jellycompat/playback_sessions.go | 59 +++++++- .../jellycompat/playback_sessions_postgres.go | 113 +++++++++++++++ .../playback_sessions_postgres_test.go | 35 +++++ 5 files changed, 336 insertions(+), 7 deletions(-) create mode 100644 internal/jellycompat/playback_negotiation_dedup_test.go diff --git a/internal/jellycompat/handlers_playback.go b/internal/jellycompat/handlers_playback.go index 6c0c7263..21884293 100644 --- a/internal/jellycompat/handlers_playback.go +++ b/internal/jellycompat/handlers_playback.go @@ -710,9 +710,14 @@ func (h *PlaybackHandler) HandlePlaybackInfo(w http.ResponseWriter, r *http.Requ return } - h.playbackStore.Put(PlaybackSession{ + clientDeviceID := firstNonEmpty( + firstMediaBrowserAuthorizationValue(r, "DeviceId"), + newCaseInsensitiveQuery(r.URL.Query()).Get("DeviceId"), + ) + h.playbackStore.PutNegotiated(PlaybackSession{ ID: playSessionID, CompatToken: session.Token, + ClientDeviceID: clientDeviceID, ItemID: detail.ContentID, RouteItemID: routeItemID, UserID: session.PseudoUserID.String(), diff --git a/internal/jellycompat/playback_negotiation_dedup_test.go b/internal/jellycompat/playback_negotiation_dedup_test.go new file mode 100644 index 00000000..85561930 --- /dev/null +++ b/internal/jellycompat/playback_negotiation_dedup_test.go @@ -0,0 +1,129 @@ +package jellycompat + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/go-chi/chi/v5" +) + +func TestPlaybackSessionStorePutNegotiatedReplacesUnstartedSameDevice(t *testing.T) { + store := NewPlaybackSessionStore(0, nil) + store.PutNegotiated(PlaybackSession{ + ID: "first", + CompatToken: "token", + ClientDeviceID: "web-device", + RouteItemID: "route", + }) + store.PutNegotiated(PlaybackSession{ + ID: "second", + CompatToken: "token", + ClientDeviceID: "web-device", + RouteItemID: "route", + }) + + if _, ok := store.Get("first"); ok { + t.Fatal("superseded unstarted negotiation was retained") + } + if _, ok := store.Get("second"); !ok { + t.Fatal("new negotiation was not stored") + } +} + +func TestPlaybackSessionStorePutNegotiatedPreservesDistinctOrStartedPlays(t *testing.T) { + tests := []struct { + name string + first PlaybackSession + }{ + { + name: "different device", + first: PlaybackSession{ + ID: "first", + CompatToken: "token", + ClientDeviceID: "other-device", + RouteItemID: "route", + }, + }, + { + name: "already started", + first: PlaybackSession{ + ID: "first", + CompatToken: "token", + ClientDeviceID: "web-device", + RouteItemID: "route", + UpstreamSessionID: "upstream-first", + }, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + store := NewPlaybackSessionStore(0, nil) + store.PutNegotiated(tc.first) + store.PutNegotiated(PlaybackSession{ + ID: "second", + CompatToken: "token", + ClientDeviceID: "web-device", + RouteItemID: "route", + }) + + if _, ok := store.Get("first"); !ok { + t.Fatal("distinct or already-started play was replaced") + } + if _, ok := store.Get("second"); !ok { + t.Fatal("new negotiation was not stored") + } + }) + } +} + +func TestHandlePlaybackInfoReplacesDuplicateJellyfinWebNegotiation(t *testing.T) { + handler, routeID := newSubtitleSelectionHandler(t) + first := postPlaybackInfoForDevice(t, handler, routeID, "web-device") + second := postPlaybackInfoForDevice(t, handler, routeID, "web-device") + + store := handler.playbackStore.(*PlaybackSessionStore) + if _, ok := store.Get(first.PlaySessionID); ok { + t.Fatal("first Jellyfin Web negotiation remained routable") + } + stored, ok := store.Get(second.PlaySessionID) + if !ok { + t.Fatal("second Jellyfin Web negotiation was not routable") + } + if stored.ClientDeviceID != "web-device" { + t.Fatalf("ClientDeviceID = %q, want web-device", stored.ClientDeviceID) + } +} + +func postPlaybackInfoForDevice( + t *testing.T, + handler *PlaybackHandler, + routeID string, + deviceID string, +) playbackInfoResponseDTO { + t.Helper() + req := httptest.NewRequest(http.MethodPost, "/Items/"+routeID+"/PlaybackInfo", strings.NewReader(`{}`)) + req.Header.Set( + "X-Emby-Authorization", + `MediaBrowser Client="Jellyfin Web", Device="Chrome", DeviceId="`+deviceID+`", Version="10.11.6"`, + ) + routeCtx := chi.NewRouteContext() + routeCtx.URLParams.Add("id", routeID) + req = req.WithContext(context.WithValue(req.Context(), chi.RouteCtxKey, routeCtx)) + req = req.WithContext(context.WithValue(req.Context(), compatSessionKey, &Session{Token: "token-1"})) + + recorder := httptest.NewRecorder() + handler.HandlePlaybackInfo(recorder, req) + if recorder.Code != http.StatusOK { + t.Fatalf("status = %d, body = %s", recorder.Code, recorder.Body.String()) + } + var response playbackInfoResponseDTO + if err := json.Unmarshal(recorder.Body.Bytes(), &response); err != nil { + t.Fatalf("unmarshal response: %v", err) + } + return response +} diff --git a/internal/jellycompat/playback_sessions.go b/internal/jellycompat/playback_sessions.go index 53468004..fbac880b 100644 --- a/internal/jellycompat/playback_sessions.go +++ b/internal/jellycompat/playback_sessions.go @@ -21,8 +21,13 @@ var ErrTerminalClaimUnavailable = errors.New("compat terminal event claim unavai type PlaybackSession struct { ID string CompatToken string - ItemID string - RouteItemID string + // ClientDeviceID identifies the Jellyfin client installation that created + // this negotiation. Stock Jellyfin Web can issue a second PlaybackInfo + // request for the same play before it starts either response; the newer + // negotiation replaces an older, still-unstarted one from the same device. + ClientDeviceID 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 @@ -81,6 +86,9 @@ type PlaybackMediaSource struct { type CompatPlaybackStore interface { // Put stores or replaces a compat playback session. Put(session PlaybackSession) + // PutNegotiated stores a PlaybackInfo negotiation and atomically replaces + // older, still-unstarted negotiations for the same client device and item. + PutNegotiated(session PlaybackSession) // Get returns a session when it exists and is not expired. Get(id string) (*PlaybackSession, bool) // Delete removes a session. @@ -164,6 +172,16 @@ func (s *PlaybackSessionStore) Put(session PlaybackSession) { s.putNormalized(session) } +// PutNegotiated stores a freshly-created PlaybackInfo session. Jellyfin Web +// may negotiate the same play twice and then request both manifests; retaining +// both creates two native sessions, while only the newer one receives progress +// and Stopped reports. Replacing only unstarted sessions keeps real concurrent +// playback intact while preventing the abandoned negotiation from later +// publishing a stale pause. +func (s *PlaybackSessionStore) PutNegotiated(session PlaybackSession) { + s.putNegotiatedNormalized(session) +} + // putNormalized stores or replaces a compat playback session and returns the // stored copy with normalized timestamps (CreatedAt/UpdatedAt/ExpiresAt). The // durable wrapper uses the return value to persist the same timestamps the cache @@ -173,14 +191,43 @@ func (s *PlaybackSessionStore) putNormalized(session PlaybackSession) PlaybackSe s.mu.Lock() defer s.mu.Unlock() - if session.CreatedAt.IsZero() { - session.CreatedAt = s.now() + session = s.normalizeSession(session) + s.sessions[session.ID] = session + return session +} + +func (s *PlaybackSessionStore) putNegotiatedNormalized(session PlaybackSession) (PlaybackSession, []string) { + s.mu.Lock() + defer s.mu.Unlock() + + session = s.normalizeSession(session) + removed := make([]string, 0, 1) + if session.CompatToken != "" && session.ClientDeviceID != "" && session.RouteItemID != "" { + for id, existing := range s.sessions { + if id == session.ID || existing.Terminal || existing.UpstreamSessionID != "" { + continue + } + if existing.CompatToken == session.CompatToken && + existing.ClientDeviceID == session.ClientDeviceID && + mediaSourceIDsEqual(existing.RouteItemID, session.RouteItemID) { + delete(s.sessions, id) + removed = append(removed, id) + } + } } - session.UpdatedAt = s.now() + s.sessions[session.ID] = session + return session, removed +} + +func (s *PlaybackSessionStore) normalizeSession(session PlaybackSession) PlaybackSession { + now := s.now() + if session.CreatedAt.IsZero() { + session.CreatedAt = now + } + session.UpdatedAt = now if session.ExpiresAt.IsZero() { session.ExpiresAt = session.CreatedAt.Add(s.ttl) } - s.sessions[session.ID] = session return session } diff --git a/internal/jellycompat/playback_sessions_postgres.go b/internal/jellycompat/playback_sessions_postgres.go index 980ada21..61d7707e 100644 --- a/internal/jellycompat/playback_sessions_postgres.go +++ b/internal/jellycompat/playback_sessions_postgres.go @@ -127,6 +127,119 @@ func (d *DurableCompatPlaybackStore) Put(session PlaybackSession) { } } +// PutNegotiated persists a new PlaybackInfo session while removing older, +// unstarted negotiations for the same compat token, client device, and item. +// The advisory transaction lock makes the replacement atomic across Silo +// processes; the in-memory mutation is likewise atomic for the local process. +func (d *DurableCompatPlaybackStore) PutNegotiated(session PlaybackSession) { + if d.pool == nil { + d.mem.PutNegotiated(session) + return + } + + scope := "negotiated\x00" + session.CompatToken + "\x00" + session.ClientDeviceID + "\x00" + session.RouteItemID + unlockSession := d.lockSessionMutation(scope) + defer unlockSession() + d.cacheMutationMu.RLock() + stored, locallyRemoved := d.mem.putNegotiatedNormalized(session) + defer d.finishCacheMutation(stored.ID, stored.CompatToken) + d.markIDValidated(stored.ID) + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + durablyRemoved, err := d.replaceUnstartedNegotiation(ctx, stored) + if err != nil { + d.markUnpersisted(stored.ID) + slog.WarnContext(ctx, "persist negotiated compat playback session failed", + "component", "jellycompat", + "error", err, + "play_session_id", stored.ID, + ) + } else { + d.clearUnpersisted(stored.ID) + } + + removed := make(map[string]struct{}, len(locallyRemoved)+len(durablyRemoved)) + for _, id := range locallyRemoved { + removed[id] = struct{}{} + } + for _, id := range durablyRemoved { + removed[id] = struct{}{} + } + for id := range removed { + d.invalidateValidation(id, "") + d.clearUnpersisted(id) + d.clearPendingUpdates(id) + d.bumpCacheGenerations(id, "") + } + d.invalidateValidation("", stored.CompatToken) +} + +func (d *DurableCompatPlaybackStore) replaceUnstartedNegotiation( + ctx context.Context, + session PlaybackSession, +) ([]string, error) { + tx, err := d.pool.Begin(ctx) + if err != nil { + return nil, err + } + defer func() { _ = tx.Rollback(ctx) }() + + var removed []string + if session.CompatToken != "" && session.ClientDeviceID != "" && session.RouteItemID != "" { + scope := session.CompatToken + "\x00" + session.ClientDeviceID + "\x00" + session.RouteItemID + if _, err := tx.Exec(ctx, `SELECT pg_advisory_xact_lock(hashtextextended($1, 0))`, scope); err != nil { + return nil, err + } + rows, err := tx.Query(ctx, ` + DELETE FROM jellycompat_playback_sessions + WHERE id <> $1 + AND compat_token = $2 + AND data->>'ClientDeviceID' = $3 + AND data->>'RouteItemID' = $4 + AND COALESCE(data->>'UpstreamSessionID', '') = '' + AND COALESCE((data->>'Terminal')::boolean, false) = false + AND expires_at > $5 + RETURNING id + `, session.ID, session.CompatToken, session.ClientDeviceID, session.RouteItemID, d.now()) + if err != nil { + return nil, err + } + for rows.Next() { + var id string + if err := rows.Scan(&id); err != nil { + rows.Close() + return nil, err + } + removed = append(removed, id) + } + if err := rows.Err(); err != nil { + rows.Close() + return nil, err + } + rows.Close() + } + + data, err := json.Marshal(session) + if err != nil { + return nil, err + } + expiresAt := session.ExpiresAt + if expiresAt.IsZero() { + expiresAt = d.now().Add(d.ttl) + } + if _, err := tx.Exec( + ctx, upsertSessionQuery, + session.ID, session.CompatToken, session.UserID, data, expiresAt, + ); err != nil { + return nil, err + } + if err := tx.Commit(ctx); err != nil { + return nil, err + } + return removed, nil +} + // Get periodically revalidates the durable row before returning an active // session. Query failures preserve a still-valid cache entry: a temporary DB // outage must not interrupt an already-playing stream. diff --git a/internal/jellycompat/playback_sessions_postgres_test.go b/internal/jellycompat/playback_sessions_postgres_test.go index f755af21..81d98955 100644 --- a/internal/jellycompat/playback_sessions_postgres_test.go +++ b/internal/jellycompat/playback_sessions_postgres_test.go @@ -86,6 +86,41 @@ func TestDurableCompatPlaybackStore_SurvivesRestart(t *testing.T) { } } +func TestDurableCompatPlaybackStorePutNegotiatedReplacesAcrossInstances(t *testing.T) { + pool := newCompatTestPool(t) + ctx := context.Background() + suffix := fmt.Sprintf("%d", time.Now().UnixNano()) + firstID := "compat-negotiated-first-" + suffix + secondID := "compat-negotiated-second-" + suffix + t.Cleanup(func() { + _, _ = pool.Exec(ctx, `DELETE FROM jellycompat_playback_sessions WHERE id = ANY($1)`, []string{firstID, secondID}) + }) + + first := NewDurableCompatPlaybackStore(pool, time.Hour, nil) + first.PutNegotiated(PlaybackSession{ + ID: firstID, + CompatToken: "negotiated-token-" + suffix, + ClientDeviceID: "web-device", + RouteItemID: "route-1", + }) + + second := NewDurableCompatPlaybackStore(pool, time.Hour, nil) + second.PutNegotiated(PlaybackSession{ + ID: secondID, + CompatToken: "negotiated-token-" + suffix, + ClientDeviceID: "web-device", + RouteItemID: "route-1", + }) + + fresh := NewDurableCompatPlaybackStore(pool, time.Hour, nil) + if _, ok := fresh.Get(firstID); ok { + t.Fatal("superseded negotiation remained durable") + } + if _, ok := fresh.Get(secondID); !ok { + t.Fatal("replacement negotiation was not durable") + } +} + // M5: an empty compat token must never trigger a DB scan — FindByRoute returns // the in-memory result only. With a non-nil pool but no live DB, a scan attempt // would block/error on the pool; instead the empty-token path returns cleanly