fix(jellycompat): deduplicate playback negotiations

- Replace unstarted negotiations for the same device and item
- Apply deduplication atomically across durable store instances
This commit is contained in:
Quick104
2026-07-22 22:01:02 -04:00
parent 166c5ef32f
commit 487dc84829
5 changed files with 336 additions and 7 deletions
+6 -1
View File
@@ -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(),
@@ -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
}
+53 -6
View File
@@ -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
}
@@ -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.
@@ -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