From a7919b9baca5cb5e5ab4519bf3635abdb735a2a5 Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Sun, 20 Sep 2026 23:51:29 +1000 Subject: [PATCH] fix(userdata): read user data through to the database The user data cache and the item's in-memory rows were both filled once per replica and never invalidated, so a pod serving a playback tick read its own stale resume position and wrote it back over the position another pod had just saved, silently losing resume points, played state, favourites and ratings. - drop the user item data cache - read single and batched user data from the database on every query - prefetch user data for in-memory sorts and filters so each stays one query - cover the lost update against real PostgreSQL with two manager instances --- .../Library/LibraryManager.cs | 44 +++- .../Library/UserDataManager.cs | 105 +++----- .../Sorting/DatePlayedComparer.cs | 9 +- .../Sorting/IsFavoriteOrLikeComparer.cs | 10 +- .../Sorting/IsPlayedComparer.cs | 10 +- .../Sorting/IsUnplayedComparer.cs | 10 +- .../Sorting/PlayCountComparer.cs | 10 +- .../Entities/UserViewBuilder.cs | 46 +++- .../Sorting/IUserBaseItemComparer.cs | 9 + .../Sorting/UserBaseItemComparerExtensions.cs | 28 +++ .../Channels/ChannelManager.cs | 3 +- src/Jellyfin.LiveTv/LiveTvManager.cs | 21 +- .../Library/UserDataManagerTests.cs | 73 +++--- .../Library/UserDataManagerReplicaTests.cs | 228 ++++++++++++++++++ 14 files changed, 470 insertions(+), 136 deletions(-) create mode 100644 MediaBrowser.Controller/Sorting/UserBaseItemComparerExtensions.cs create mode 100644 tests/Jellyfin.Server.Tests/Library/UserDataManagerReplicaTests.cs diff --git a/Emby.Server.Implementations/Library/LibraryManager.cs b/Emby.Server.Implementations/Library/LibraryManager.cs index 0044fcd4dc..97e03b45c9 100644 --- a/Emby.Server.Implementations/Library/LibraryManager.cs +++ b/Emby.Server.Implementations/Library/LibraryManager.cs @@ -2332,7 +2332,10 @@ namespace Emby.Server.Implementations.Library { IOrderedEnumerable? orderedItems = null; - foreach (var orderBy in sortBy.Select(o => GetComparer(o, user)).Where(c => c is not null)) + var comparers = sortBy.Select(o => GetComparer(o, user)).Where(c => c is not null).ToList(); + items = PrefetchUserData(items, user, comparers); + + foreach (var orderBy in comparers) { if (orderBy is RandomComparer) { @@ -2364,14 +2367,14 @@ namespace Emby.Server.Implementations.Library { IOrderedEnumerable? orderedItems = null; - foreach (var (name, sortOrder) in orderBy) - { - var comparer = GetComparer(name, user); - if (comparer is null) - { - continue; - } + var comparers = orderBy + .Select(o => (Comparer: GetComparer(o.OrderBy, user), o.SortOrder)) + .Where(c => c.Comparer is not null) + .ToList(); + items = PrefetchUserData(items, user, comparers.Select(c => c.Comparer).ToList()); + foreach (var (comparer, sortOrder) in comparers) + { if (comparer is RandomComparer) { var randomItems = items.ToArray(); @@ -2397,6 +2400,31 @@ namespace Emby.Server.Implementations.Library return orderedItems ?? items; } + // The user comparers read user data per item, so without one batched read up front an + // in-memory sort would issue a database round trip per comparison. + private IEnumerable PrefetchUserData(IEnumerable items, User? user, IReadOnlyList comparers) + { + if (user is null) + { + return items; + } + + var userComparers = comparers.OfType().ToList(); + if (userComparers.Count == 0) + { + return items; + } + + var itemList = items as IReadOnlyList ?? items.ToList(); + var userData = _userDataManager.GetUserDataBatch(itemList, user); + foreach (var comparer in userComparers) + { + comparer.PrefetchedUserData = userData; + } + + return itemList; + } + /// /// Gets the comparer. /// diff --git a/Emby.Server.Implementations/Library/UserDataManager.cs b/Emby.Server.Implementations/Library/UserDataManager.cs index 0680046c11..4ddbd60671 100644 --- a/Emby.Server.Implementations/Library/UserDataManager.cs +++ b/Emby.Server.Implementations/Library/UserDataManager.cs @@ -2,10 +2,8 @@ using System; using System.Collections.Generic; -using System.Globalization; using System.Linq; using System.Threading; -using BitFaster.Caching.Lru; using Jellyfin.Database.Implementations; using Jellyfin.Database.Implementations.Entities; using MediaBrowser.Controller.Configuration; @@ -27,7 +25,6 @@ namespace Emby.Server.Implementations.Library { private readonly IServerConfigurationManager _config; private readonly IDbContextFactory _repository; - private readonly FastConcurrentLru _cache; /// /// Initializes a new instance of the class. @@ -40,7 +37,6 @@ namespace Emby.Server.Implementations.Library { _config = config; _repository = repository; - _cache = new FastConcurrentLru(Environment.ProcessorCount, _config.Configuration.CacheSize, StringComparer.OrdinalIgnoreCase); } /// @@ -77,11 +73,6 @@ namespace Emby.Server.Implementations.Library dbContext.SaveChanges(); transaction.Commit(); - var userId = user.InternalId; - var cacheKey = GetCacheKey(userId, item.Id); - _cache.AddOrUpdate(cacheKey, userData); - item.UserData = dbContext.UserData.Where(e => e.ItemId == item.Id).AsNoTracking().ToArray(); // rehydrate the cached userdata - UserDataSaved?.Invoke(this, new UserDataSaveEventArgs { Keys = keys, @@ -180,64 +171,41 @@ namespace Emby.Server.Implementations.Library /// public Dictionary GetUserDataBatch(IReadOnlyList items, User user) { + ArgumentNullException.ThrowIfNull(items); + ArgumentNullException.ThrowIfNull(user); + var result = new Dictionary(items.Count); - var itemsNeedingQuery = new List<(BaseItem Item, List Keys)>(); - - foreach (var item in items) - { - var cacheKey = GetCacheKey(user.InternalId, item.Id); - if (_cache.TryGet(cacheKey, out var cachedData)) - { - result[item.Id] = cachedData; - } - else - { - var userDataRow = ResolveUserDataRow(item, item.UserData?.Where(e => e.UserId.Equals(user.Id))); - var userData = userDataRow is not null ? Map(userDataRow) : null; - if (userData is not null) - { - result[item.Id] = userData; - _cache.AddOrUpdate(cacheKey, userData); - } - else - { - var keys = item.GetUserDataKeys(); - itemsNeedingQuery.Add((item, keys)); - } - } - } - - if (itemsNeedingQuery.Count == 0) + if (items.Count == 0) { return result; } - // Build a single query for all missing items. Fetch rows by item alone so rows kept - // under keys from older metadata resolve the same way as the in-memory path. - var allItemIds = itemsNeedingQuery.Select(x => x.Item.Id).ToList(); + // Fetch rows by item alone so rows kept under keys from older metadata resolve the same + // way as the single item path. + var itemIds = items.Select(e => e.Id).Distinct().ToList(); using var context = _repository.CreateDbContext(); - var userDataArray = context.UserData + var userDataByItem = context.UserData .AsNoTracking() .Where(e => e.UserId.Equals(user.Id)) - .WhereOneOrMany(allItemIds, e => e.ItemId) - .ToArray(); + .WhereOneOrMany(itemIds, e => e.ItemId) + .ToArray() + .GroupBy(e => e.ItemId) + .ToDictionary(g => g.Key, g => g.ToArray()); - var userDataByItem = userDataArray.GroupBy(e => e.ItemId).ToDictionary(g => g.Key, g => g.ToArray()); - foreach (var (item, keys) in itemsNeedingQuery) + foreach (var item in items) { - UserItemData userData; - if (userDataByItem.TryGetValue(item.Id, out var itemUserData) && itemUserData.Length > 0) + if (result.ContainsKey(item.Id)) { - userData = Map(ResolveUserDataRow(item, itemUserData)!); - } - else - { - userData = new UserItemData { Key = keys.Count > 0 ? keys[0] : string.Empty }; + continue; } - result[item.Id] = userData; - var cacheKey = GetCacheKey(user.InternalId, item.Id); - _cache.AddOrUpdate(cacheKey, userData); + var row = userDataByItem.TryGetValue(item.Id, out var itemUserData) + ? ResolveUserDataRow(item, itemUserData) + : null; + + result[item.Id] = row is not null + ? Map(row) + : new UserItemData { Key = item.GetUserDataKeys().FirstOrDefault() ?? string.Empty }; } return result; @@ -340,20 +308,19 @@ namespace Emby.Server.Implementations.Library return result; } - /// - /// Gets the internal key. - /// - /// System.String. - private static string GetCacheKey(long internalUserId, Guid itemId) - { - return internalUserId.ToString(CultureInfo.InvariantCulture) + "-" + itemId.ToString("N", CultureInfo.InvariantCulture); - } - /// public UserItemData? GetUserData(User user, BaseItem item) { ArgumentNullException.ThrowIfNull(user); - var row = ResolveUserDataRow(item, item.UserData?.Where(e => e.UserId.Equals(user.Id))); + ArgumentNullException.ThrowIfNull(item); + + using var dbContext = _repository.CreateDbContext(); + var rows = dbContext.UserData + .AsNoTracking() + .Where(e => e.ItemId == item.Id && e.UserId == user.Id) + .ToArray(); + + var row = ResolveUserDataRow(item, rows); return row is not null ? Map(row) : new UserItemData() { Key = item.GetUserDataKeys()[0], @@ -536,16 +503,6 @@ namespace Emby.Server.Implementations.Library } dbContext.SaveChanges(); - - var cacheKey = GetCacheKey(user.InternalId, item.Id); - if (_cache.TryGet(cacheKey, out var cached)) - { - cached.AudioStreamIndex = null; - cached.SubtitleStreamIndex = null; - _cache.AddOrUpdate(cacheKey, cached); - } - - item.UserData = dbContext.UserData.Where(e => e.ItemId == item.Id).AsNoTracking().ToArray(); } } } diff --git a/Emby.Server.Implementations/Sorting/DatePlayedComparer.cs b/Emby.Server.Implementations/Sorting/DatePlayedComparer.cs index 2c8e2b37d0..917485b457 100644 --- a/Emby.Server.Implementations/Sorting/DatePlayedComparer.cs +++ b/Emby.Server.Implementations/Sorting/DatePlayedComparer.cs @@ -1,6 +1,7 @@ #nullable disable using System; +using System.Collections.Generic; using Jellyfin.Data.Enums; using Jellyfin.Database.Implementations.Entities; using MediaBrowser.Controller.Entities; @@ -27,6 +28,12 @@ namespace Emby.Server.Implementations.Sorting /// The user manager. public IUserManager UserManager { get; set; } + /// + /// Gets or sets the prefetched user data. + /// + /// The prefetched user data. + public IReadOnlyDictionary PrefetchedUserData { get; set; } + /// /// Gets or sets the user data manager. /// @@ -57,7 +64,7 @@ namespace Emby.Server.Implementations.Sorting /// DateTime. private DateTime GetDate(BaseItem x) { - var userdata = UserDataManager.GetUserData(User, x); + var userdata = this.GetUserData(x); if (userdata is not null && userdata.LastPlayedDate.HasValue) { diff --git a/Emby.Server.Implementations/Sorting/IsFavoriteOrLikeComparer.cs b/Emby.Server.Implementations/Sorting/IsFavoriteOrLikeComparer.cs index 86d08ed27b..9eafca1391 100644 --- a/Emby.Server.Implementations/Sorting/IsFavoriteOrLikeComparer.cs +++ b/Emby.Server.Implementations/Sorting/IsFavoriteOrLikeComparer.cs @@ -1,6 +1,8 @@ #nullable disable #pragma warning disable CS1591 +using System; +using System.Collections.Generic; using Jellyfin.Data.Enums; using Jellyfin.Database.Implementations.Entities; using MediaBrowser.Controller.Entities; @@ -35,6 +37,12 @@ namespace Emby.Server.Implementations.Sorting /// The user manager. public IUserManager UserManager { get; set; } + /// + /// Gets or sets the prefetched user data. + /// + /// The prefetched user data. + public IReadOnlyDictionary PrefetchedUserData { get; set; } + /// /// Compares the specified x. /// @@ -53,7 +61,7 @@ namespace Emby.Server.Implementations.Sorting /// DateTime. private int GetValue(BaseItem x) { - return x.IsFavoriteOrLiked(User, userItemData: null) ? 0 : 1; + return x.IsFavoriteOrLiked(User, this.GetUserData(x)) ? 0 : 1; } } } diff --git a/Emby.Server.Implementations/Sorting/IsPlayedComparer.cs b/Emby.Server.Implementations/Sorting/IsPlayedComparer.cs index 9faa02f1fd..b4e3787ffc 100644 --- a/Emby.Server.Implementations/Sorting/IsPlayedComparer.cs +++ b/Emby.Server.Implementations/Sorting/IsPlayedComparer.cs @@ -2,6 +2,8 @@ #pragma warning disable CS1591 +using System; +using System.Collections.Generic; using Jellyfin.Data.Enums; using Jellyfin.Database.Implementations.Entities; using MediaBrowser.Controller.Entities; @@ -36,6 +38,12 @@ namespace Emby.Server.Implementations.Sorting /// The user manager. public IUserManager UserManager { get; set; } + /// + /// Gets or sets the prefetched user data. + /// + /// The prefetched user data. + public IReadOnlyDictionary PrefetchedUserData { get; set; } + /// /// Compares the specified x. /// @@ -54,7 +62,7 @@ namespace Emby.Server.Implementations.Sorting /// DateTime. private int GetValue(BaseItem x) { - return x.IsPlayed(User, userItemData: null) ? 0 : 1; + return x.IsPlayed(User, this.GetUserData(x)) ? 0 : 1; } } } diff --git a/Emby.Server.Implementations/Sorting/IsUnplayedComparer.cs b/Emby.Server.Implementations/Sorting/IsUnplayedComparer.cs index 6f177c4637..3b27b8092a 100644 --- a/Emby.Server.Implementations/Sorting/IsUnplayedComparer.cs +++ b/Emby.Server.Implementations/Sorting/IsUnplayedComparer.cs @@ -2,6 +2,8 @@ #pragma warning disable CS1591 +using System; +using System.Collections.Generic; using Jellyfin.Data.Enums; using Jellyfin.Database.Implementations.Entities; using MediaBrowser.Controller.Entities; @@ -36,6 +38,12 @@ namespace Emby.Server.Implementations.Sorting /// The user manager. public IUserManager UserManager { get; set; } + /// + /// Gets or sets the prefetched user data. + /// + /// The prefetched user data. + public IReadOnlyDictionary PrefetchedUserData { get; set; } + /// /// Compares the specified x. /// @@ -54,7 +62,7 @@ namespace Emby.Server.Implementations.Sorting /// DateTime. private int GetValue(BaseItem x) { - return x.IsUnplayed(User, userItemData: null) ? 0 : 1; + return x.IsUnplayed(User, this.GetUserData(x)) ? 0 : 1; } } } diff --git a/Emby.Server.Implementations/Sorting/PlayCountComparer.cs b/Emby.Server.Implementations/Sorting/PlayCountComparer.cs index 26e28b03bc..568e9e69f1 100644 --- a/Emby.Server.Implementations/Sorting/PlayCountComparer.cs +++ b/Emby.Server.Implementations/Sorting/PlayCountComparer.cs @@ -1,5 +1,7 @@ #nullable disable +using System; +using System.Collections.Generic; using Jellyfin.Data.Enums; using Jellyfin.Database.Implementations.Entities; using MediaBrowser.Controller.Entities; @@ -38,6 +40,12 @@ namespace Emby.Server.Implementations.Sorting /// The user manager. public IUserManager UserManager { get; set; } + /// + /// Gets or sets the prefetched user data. + /// + /// The prefetched user data. + public IReadOnlyDictionary PrefetchedUserData { get; set; } + /// /// Compares the specified x. /// @@ -56,7 +64,7 @@ namespace Emby.Server.Implementations.Sorting /// DateTime. private int GetValue(BaseItem x) { - var userdata = UserDataManager.GetUserData(User, x); + var userdata = this.GetUserData(x); return userdata is null ? 0 : userdata.PlayCount; } diff --git a/MediaBrowser.Controller/Entities/UserViewBuilder.cs b/MediaBrowser.Controller/Entities/UserViewBuilder.cs index f9ad2d86e6..69ca1a351f 100644 --- a/MediaBrowser.Controller/Entities/UserViewBuilder.cs +++ b/MediaBrowser.Controller/Entities/UserViewBuilder.cs @@ -449,19 +449,26 @@ namespace MediaBrowser.Controller.Entities IUserDataManager userDataManager, ILibraryManager libraryManager) { - var filtered = items.Where(i => Filter(i, user, query, userDataManager, libraryManager)); + var itemList = items as IReadOnlyList ?? items.ToList(); + + // The user data checks below run per item, so read them all in one query up front. + var userDataBatch = user is not null && RequiresUserData(query) + ? userDataManager.GetUserDataBatch(itemList, user) + : null; + + var filtered = itemList.Where(i => Filter(i, user, query, userDataManager, libraryManager, userDataBatch)); if (query.IsPlayed.HasValue && user is not null) { - var itemList = filtered.ToList(); - var folderIds = itemList.OfType().Select(f => f.Id).ToList(); + var filteredList = filtered.ToList(); + var folderIds = filteredList.OfType().Select(f => f.Id).ToList(); if (folderIds.Count > 0) { var counts = libraryManager.GetPlayedAndTotalCountBatch(folderIds, user); var isPlayedValue = query.IsPlayed.Value; - return itemList.Where(item => + return filteredList.Where(item => { if (item is Folder) { @@ -473,7 +480,7 @@ namespace MediaBrowser.Controller.Entities }); } - return itemList; + return filteredList; } return filtered; @@ -515,12 +522,29 @@ namespace MediaBrowser.Controller.Entities itemsArray); } + private static bool RequiresUserData(InternalItemsQuery query) + => query.IsLiked.HasValue + || query.IsFavoriteOrLiked.HasValue + || query.IsFavorite.HasValue + || query.IsResumable.HasValue + || query.IsPlayed.HasValue; + + private static UserItemData GetUserData( + IUserDataManager userDataManager, + User user, + BaseItem item, + Dictionary userDataBatch) + => userDataBatch is not null && userDataBatch.TryGetValue(item.Id, out var userData) + ? userData + : userDataManager.GetUserData(user, item); + private static bool Filter( BaseItem item, User user, InternalItemsQuery query, IUserDataManager userDataManager, - ILibraryManager libraryManager) + ILibraryManager libraryManager, + Dictionary userDataBatch) { if (!string.IsNullOrEmpty(query.NameStartsWith) && !item.SortName.StartsWith(query.NameStartsWith, StringComparison.InvariantCultureIgnoreCase)) { @@ -568,7 +592,7 @@ namespace MediaBrowser.Controller.Entities if (query.IsLiked.HasValue) { - userData = userDataManager.GetUserData(user, item); + userData = GetUserData(userDataManager, user, item, userDataBatch); if (!userData.Likes.HasValue || userData.Likes != query.IsLiked.Value) { return false; @@ -577,7 +601,7 @@ namespace MediaBrowser.Controller.Entities if (query.IsFavoriteOrLiked.HasValue) { - userData ??= userDataManager.GetUserData(user, item); + userData ??= GetUserData(userDataManager, user, item, userDataBatch); var isFavoriteOrLiked = userData.IsFavorite || (userData.Likes ?? false); if (isFavoriteOrLiked != query.IsFavoriteOrLiked.Value) @@ -588,7 +612,7 @@ namespace MediaBrowser.Controller.Entities if (query.IsFavorite.HasValue) { - userData ??= userDataManager.GetUserData(user, item); + userData ??= GetUserData(userDataManager, user, item, userDataBatch); if (userData.IsFavorite != query.IsFavorite.Value) { return false; @@ -597,7 +621,7 @@ namespace MediaBrowser.Controller.Entities if (query.IsResumable.HasValue) { - userData ??= userDataManager.GetUserData(user, item); + userData ??= GetUserData(userDataManager, user, item, userDataBatch); var isResumable = userData.PlaybackPositionTicks > 0; if (isResumable != query.IsResumable.Value) @@ -612,7 +636,7 @@ namespace MediaBrowser.Controller.Entities // Folders are batch-filtered by the collection Filter() overload. if (!item.IsFolder) { - userData ??= userDataManager.GetUserData(user, item); + userData ??= GetUserData(userDataManager, user, item, userDataBatch); if (item.IsPlayed(user, userData) != query.IsPlayed.Value) { return false; diff --git a/MediaBrowser.Controller/Sorting/IUserBaseItemComparer.cs b/MediaBrowser.Controller/Sorting/IUserBaseItemComparer.cs index 2206a021a8..1b0a414916 100644 --- a/MediaBrowser.Controller/Sorting/IUserBaseItemComparer.cs +++ b/MediaBrowser.Controller/Sorting/IUserBaseItemComparer.cs @@ -1,6 +1,9 @@ #nullable disable +using System; +using System.Collections.Generic; using Jellyfin.Database.Implementations.Entities; +using MediaBrowser.Controller.Entities; using MediaBrowser.Controller.Library; namespace MediaBrowser.Controller.Sorting @@ -27,5 +30,11 @@ namespace MediaBrowser.Controller.Sorting /// /// The user data repository. IUserDataManager UserDataManager { get; set; } + + /// + /// Gets or sets user data for the items being sorted, keyed by item id, read once up front. + /// + /// The prefetched user data, or null when none was prefetched. + IReadOnlyDictionary PrefetchedUserData { get; set; } } } diff --git a/MediaBrowser.Controller/Sorting/UserBaseItemComparerExtensions.cs b/MediaBrowser.Controller/Sorting/UserBaseItemComparerExtensions.cs new file mode 100644 index 0000000000..c4fb395e4f --- /dev/null +++ b/MediaBrowser.Controller/Sorting/UserBaseItemComparerExtensions.cs @@ -0,0 +1,28 @@ +#nullable disable + +using MediaBrowser.Controller.Entities; + +namespace MediaBrowser.Controller.Sorting +{ + /// + /// Helpers shared by the comparers that sort on user data. + /// + public static class UserBaseItemComparerExtensions + { + /// + /// Gets the user data for an item, preferring the batch the sort prefetched. + /// + /// The comparer. + /// The item. + /// The item's user data. + public static UserItemData GetUserData(this IUserBaseItemComparer comparer, BaseItem item) + { + if (comparer.PrefetchedUserData is not null && comparer.PrefetchedUserData.TryGetValue(item.Id, out var userData)) + { + return userData; + } + + return comparer.UserDataManager.GetUserData(comparer.User, item); + } + } +} diff --git a/src/Jellyfin.LiveTv/Channels/ChannelManager.cs b/src/Jellyfin.LiveTv/Channels/ChannelManager.cs index ed02fe6a1d..c0fef94f0b 100644 --- a/src/Jellyfin.LiveTv/Channels/ChannelManager.cs +++ b/src/Jellyfin.LiveTv/Channels/ChannelManager.cs @@ -212,7 +212,8 @@ namespace Jellyfin.LiveTv.Channels if (query.IsFavorite.HasValue) { var val = query.IsFavorite.Value; - channels = channels.Where(i => _userDataManager.GetUserData(user, i).IsFavorite == val) + var userData = _userDataManager.GetUserDataBatch(channels, user); + channels = channels.Where(i => userData.TryGetValue(i.Id, out var data) && data.IsFavorite == val) .ToList(); } diff --git a/src/Jellyfin.LiveTv/LiveTvManager.cs b/src/Jellyfin.LiveTv/LiveTvManager.cs index 2edf7681db..ed389d9da0 100644 --- a/src/Jellyfin.LiveTv/LiveTvManager.cs +++ b/src/Jellyfin.LiveTv/LiveTvManager.cs @@ -304,8 +304,17 @@ namespace Jellyfin.LiveTv if (query.IsAiring ?? false) { + // Scoring reads the channel's user data per program, so read every channel's in one query. + var channels = programList + .Cast() + .Select(i => _libraryManager.GetItemById(i.ChannelId)) + .OfType() + .DistinctBy(i => i.Id) + .ToList(); + var channelUserData = _userDataManager.GetUserDataBatch(channels, user); + orderedPrograms = orderedPrograms - .ThenByDescending(i => GetRecommendationScore(i, user, true)); + .ThenByDescending(i => GetRecommendationScore(i, user, true, channelUserData)); } IEnumerable programs = orderedPrograms; @@ -338,7 +347,11 @@ namespace Jellyfin.LiveTv _dtoService.GetBaseItemDtos(internalResult.Items, options, query.User))); } - private int GetRecommendationScore(LiveTvProgram program, User user, bool factorChannelWatchCount) + private int GetRecommendationScore( + LiveTvProgram program, + User user, + bool factorChannelWatchCount, + IReadOnlyDictionary channelUserData) { var score = 0; @@ -359,7 +372,9 @@ namespace Jellyfin.LiveTv return score; } - var channelUserdata = _userDataManager.GetUserData(user, channel); + var channelUserdata = channelUserData.TryGetValue(channel.Id, out var cached) + ? cached + : _userDataManager.GetUserData(user, channel); if (channelUserdata.Likes.HasValue) { diff --git a/tests/Jellyfin.Server.Implementations.Tests/Library/UserDataManagerTests.cs b/tests/Jellyfin.Server.Implementations.Tests/Library/UserDataManagerTests.cs index ba3127bc08..6810ba64de 100644 --- a/tests/Jellyfin.Server.Implementations.Tests/Library/UserDataManagerTests.cs +++ b/tests/Jellyfin.Server.Implementations.Tests/Library/UserDataManagerTests.cs @@ -1,5 +1,4 @@ using System; -using System.Collections.Generic; using Emby.Server.Implementations.Library; using Jellyfin.Database.Implementations; using Jellyfin.Database.Implementations.Entities; @@ -49,6 +48,12 @@ public sealed class UserDataManagerTests : IDisposable { Id = Guid.NewGuid() }; + + using (var ctx = CreateDbContext()) + { + ctx.Users.Add(_user); + ctx.SaveChanges(); + } } public void Dispose() @@ -78,6 +83,23 @@ public sealed class UserDataManagerTests : IDisposable }; } + private void Seed(AudioBook item, params UserData[] rows) + { + using var ctx = CreateDbContext(); + ctx.BaseItems.Add(new BaseItemEntity { Id = item.Id, Type = typeof(AudioBook).FullName! }); + ctx.UserData.AddRange(rows); + ctx.SaveChanges(); + } + + private User CreateOtherUser() + { + var user = new User("other", "auth-provider", "reset-provider") { Id = Guid.NewGuid() }; + using var ctx = CreateDbContext(); + ctx.Users.Add(user); + ctx.SaveChanges(); + return user; + } + private UserData CreateUserDataRow(AudioBook item, string key, long positionTicks) { return new UserData @@ -98,11 +120,10 @@ public sealed class UserDataManagerTests : IDisposable var currentKey = item.GetUserDataKeys()[0]; // the retired-key row comes first to ensure selection is by key, not row order - item.UserData = new List - { + Seed( + item, CreateUserDataRow(item, "Author-Old Album-0001Old File Name", 111), - CreateUserDataRow(item, currentKey, 222) - }; + CreateUserDataRow(item, currentKey, 222)); var userData = _userDataManager.GetUserData(_user, item); @@ -117,11 +138,10 @@ public sealed class UserDataManagerTests : IDisposable var item = CreateAudioBook(); var idKey = item.GetUserDataKeys()[1]; - item.UserData = new List - { + Seed( + item, CreateUserDataRow(item, "Author-Old Album-0001Old File Name", 111), - CreateUserDataRow(item, idKey, 333) - }; + CreateUserDataRow(item, idKey, 333)); var userData = _userDataManager.GetUserData(_user, item); @@ -135,10 +155,7 @@ public sealed class UserDataManagerTests : IDisposable { var item = CreateAudioBook(); - item.UserData = new List - { - CreateUserDataRow(item, "Author-Old Album-0001Old File Name", 111) - }; + Seed(item, CreateUserDataRow(item, "Author-Old Album-0001Old File Name", 111)); var userData = _userDataManager.GetUserData(_user, item); @@ -150,7 +167,7 @@ public sealed class UserDataManagerTests : IDisposable public void GetUserData_NoRows_ReturnsDefaultWithPrimaryKey() { var item = CreateAudioBook(); - item.UserData = new List(); + Seed(item); var userData = _userDataManager.GetUserData(_user, item); @@ -166,13 +183,9 @@ public sealed class UserDataManagerTests : IDisposable var currentKey = item.GetUserDataKeys()[0]; var otherUserRow = CreateUserDataRow(item, currentKey, 999); - otherUserRow.UserId = Guid.NewGuid(); + otherUserRow.UserId = CreateOtherUser().Id; - item.UserData = new List - { - otherUserRow, - CreateUserDataRow(item, currentKey, 222) - }; + Seed(item, otherUserRow, CreateUserDataRow(item, currentKey, 222)); var userData = _userDataManager.GetUserData(_user, item); @@ -183,23 +196,15 @@ public sealed class UserDataManagerTests : IDisposable [Fact] public void GetUserDataBatch_DatabaseFallback_ResolvesRowsByKeyOrder() { - // no preloaded navigation data, so the batch takes the database fallback var fossilItem = CreateAudioBook(); var retiredItem = CreateAudioBook(); - using (var ctx = CreateDbContext()) - { - ctx.Users.Add(_user); - ctx.BaseItems.Add(new BaseItemEntity { Id = fossilItem.Id, Type = typeof(AudioBook).FullName! }); - ctx.BaseItems.Add(new BaseItemEntity { Id = retiredItem.Id, Type = typeof(AudioBook).FullName! }); - - // the stale id-key row is inserted first so selection by row order would return it - ctx.UserData.AddRange( - CreateUserDataRow(fossilItem, fossilItem.GetUserDataKeys()[1], 111), - CreateUserDataRow(fossilItem, fossilItem.GetUserDataKeys()[0], 222), - CreateUserDataRow(retiredItem, "Author-Old Album-0001Old File Name", 333)); - ctx.SaveChanges(); - } + // the stale id-key row is inserted first so selection by row order would return it + Seed( + fossilItem, + CreateUserDataRow(fossilItem, fossilItem.GetUserDataKeys()[1], 111), + CreateUserDataRow(fossilItem, fossilItem.GetUserDataKeys()[0], 222)); + Seed(retiredItem, CreateUserDataRow(retiredItem, "Author-Old Album-0001Old File Name", 333)); var result = _userDataManager.GetUserDataBatch([fossilItem, retiredItem], _user); diff --git a/tests/Jellyfin.Server.Tests/Library/UserDataManagerReplicaTests.cs b/tests/Jellyfin.Server.Tests/Library/UserDataManagerReplicaTests.cs new file mode 100644 index 0000000000..020f86f6b0 --- /dev/null +++ b/tests/Jellyfin.Server.Tests/Library/UserDataManagerReplicaTests.cs @@ -0,0 +1,228 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Emby.Server.Implementations.Library; +using Jellyfin.Database.Implementations; +using Jellyfin.Database.Implementations.DbConfiguration; +using Jellyfin.Database.Implementations.Entities; +using Jellyfin.Database.Implementations.Locking; +using Jellyfin.Database.Providers.PostgreSQL; +using Jellyfin.Server.Tests.Migrations; +using MediaBrowser.Controller.Configuration; +using MediaBrowser.Model.Configuration; +using MediaBrowser.Model.Entities; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Logging.Abstractions; +using Moq; +using Npgsql; +using Xunit; +using AudioBook = MediaBrowser.Controller.Entities.AudioBook; + +namespace Jellyfin.Server.Tests.Library; + +/// +/// Two independently constructed instances over one PostgreSQL database are the +/// in-process stand-in for two replicas sharing one database: what either of them writes, the other has to +/// see on its very next read, and a read-modify-write on one must not roll back the other's. +/// +[Trait("Category", "RequiresDocker")] +public sealed class UserDataManagerReplicaTests : IAsyncLifetime +{ + private static readonly long _quarterIn = TimeSpan.FromMinutes(20).Ticks; + + private PostgreSqlTestServer _server = null!; + + /// + public async ValueTask InitializeAsync() + { + _server = await PostgreSqlTestServer.StartAsync().ConfigureAwait(false); + } + + /// + public async ValueTask DisposeAsync() + { + await _server.DisposeAsync().ConfigureAwait(false); + } + + /// + /// A resume position written by the replica serving the playback tick has to be the position the next + /// request reads, whichever replica it lands on - both through the single item read the write path uses + /// and through the batch read the library pages render from. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task ResumePositionWrittenOnOneReplica_IsReadOnAnother() + { + var cancellationToken = TestContext.Current.CancellationToken; + var connectionString = await _server.CreateDatabaseAsync("userdata_replica_resume", cancellationToken); + + await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build(); + var itemId = Guid.NewGuid(); + var user = await CreateSchemaWithUserAndItemAsync(dataSource, itemId, cancellationToken); + + var replicaA = CreateManager(dataSource); + var replicaB = CreateManager(dataSource); + var itemOnA = new AudioBook { Id = itemId, Name = "Replica Book" }; + var itemOnB = new AudioBook { Id = itemId, Name = "Replica Book" }; + + var early = replicaA.GetUserData(user, itemOnA)!; + early.PlaybackPositionTicks = TimeSpan.FromMinutes(5).Ticks; + replicaA.SaveUserData(user, itemOnA, early, UserDataSaveReason.PlaybackProgress, cancellationToken); + + // Replica B materialised the item before the later tick, so it holds the earlier row in memory. + itemOnB.UserData = await LoadUserDataAsync(dataSource, itemId, cancellationToken); + + var later = replicaA.GetUserData(user, itemOnA)!; + later.PlaybackPositionTicks = _quarterIn; + replicaA.SaveUserData(user, itemOnA, later, UserDataSaveReason.PlaybackProgress, cancellationToken); + + Assert.Equal(_quarterIn, replicaB.GetUserData(user, itemOnB)!.PlaybackPositionTicks); + Assert.Equal(_quarterIn, replicaB.GetUserDataBatch([itemOnB], user)[itemId].PlaybackPositionTicks); + } + + /// + /// The playback tick is a read-modify-write of the whole row, so a tick served by one replica must build + /// on the favourite another replica just recorded instead of writing it back out. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task PlaybackTickOnOneReplica_KeepsFavouriteSetOnAnother() + { + var cancellationToken = TestContext.Current.CancellationToken; + var connectionString = await _server.CreateDatabaseAsync("userdata_replica_favourite", cancellationToken); + + await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build(); + var itemId = Guid.NewGuid(); + var user = await CreateSchemaWithUserAndItemAsync(dataSource, itemId, cancellationToken); + + var replicaA = CreateManager(dataSource); + var replicaB = CreateManager(dataSource); + var itemOnA = new AudioBook { Id = itemId, Name = "Replica Book" }; + var itemOnB = new AudioBook { Id = itemId, Name = "Replica Book" }; + + var seed = replicaA.GetUserData(user, itemOnA)!; + seed.PlaybackPositionTicks = TimeSpan.FromMinutes(5).Ticks; + replicaA.SaveUserData(user, itemOnA, seed, UserDataSaveReason.PlaybackProgress, cancellationToken); + + // Replica B is serving the playback session and read the item before the favourite was recorded. + itemOnB.UserData = await LoadUserDataAsync(dataSource, itemId, cancellationToken); + + var favourited = replicaA.GetUserData(user, itemOnA)!; + favourited.IsFavorite = true; + replicaA.SaveUserData(user, itemOnA, favourited, UserDataSaveReason.UpdateUserRating, cancellationToken); + + var tick = replicaB.GetUserData(user, itemOnB)!; + tick.PlaybackPositionTicks = _quarterIn; + replicaB.SaveUserData(user, itemOnB, tick, UserDataSaveReason.PlaybackProgress, cancellationToken); + + var stored = replicaA.GetUserData(user, itemOnA)!; + Assert.True(stored.IsFavorite); + Assert.Equal(_quarterIn, stored.PlaybackPositionTicks); + } + + /// + /// A tick that lands on the other replica has to carry the position forward from where the session + /// actually is, not from the position that replica happened to have in memory. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task PlaybackTickOnOneReplica_ResumesFromThePositionAnotherWrote() + { + var cancellationToken = TestContext.Current.CancellationToken; + var connectionString = await _server.CreateDatabaseAsync("userdata_replica_lost_update", cancellationToken); + + await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build(); + var itemId = Guid.NewGuid(); + var user = await CreateSchemaWithUserAndItemAsync(dataSource, itemId, cancellationToken); + + var replicaA = CreateManager(dataSource); + var replicaB = CreateManager(dataSource); + var itemOnA = new AudioBook { Id = itemId, Name = "Replica Book" }; + var itemOnB = new AudioBook { Id = itemId, Name = "Replica Book" }; + + var seed = replicaA.GetUserData(user, itemOnA)!; + seed.PlaybackPositionTicks = TimeSpan.FromMinutes(5).Ticks; + replicaA.SaveUserData(user, itemOnA, seed, UserDataSaveReason.PlaybackProgress, cancellationToken); + + itemOnB.UserData = await LoadUserDataAsync(dataSource, itemId, cancellationToken); + + // The viewer seeks forward and the tick reporting it lands on replica A. + var seeked = replicaA.GetUserData(user, itemOnA)!; + seeked.PlaybackPositionTicks = _quarterIn; + replicaA.SaveUserData(user, itemOnA, seeked, UserDataSaveReason.PlaybackProgress, cancellationToken); + + // The next tick lands on replica B, which adds ten seconds to whatever it reads. + var tick = replicaB.GetUserData(user, itemOnB)!; + tick.PlaybackPositionTicks += TimeSpan.FromSeconds(10).Ticks; + replicaB.SaveUserData(user, itemOnB, tick, UserDataSaveReason.PlaybackProgress, cancellationToken); + + var stored = replicaA.GetUserData(user, itemOnA)!; + Assert.Equal(_quarterIn + TimeSpan.FromSeconds(10).Ticks, stored.PlaybackPositionTicks); + } + + private static UserDataManager CreateManager(NpgsqlDataSource dataSource) + { + var config = new Mock(); + config.SetupGet(c => c.Configuration).Returns(new ServerConfiguration()); + return new UserDataManager(config.Object, new DataSourceContextFactory(dataSource)); + } + + private static async Task> LoadUserDataAsync(NpgsqlDataSource dataSource, Guid itemId, CancellationToken cancellationToken) + { + var context = CreateContext(dataSource); + await using (context.ConfigureAwait(false)) + { + return await context.UserData + .AsNoTracking() + .Where(e => e.ItemId.Equals(itemId)) + .ToArrayAsync(cancellationToken) + .ConfigureAwait(false); + } + } + + private static async Task CreateSchemaWithUserAndItemAsync(NpgsqlDataSource dataSource, Guid itemId, CancellationToken cancellationToken) + { + var context = CreateContext(dataSource); + await using (context.ConfigureAwait(false)) + { + await context.Database.EnsureCreatedAsync(cancellationToken).ConfigureAwait(false); + + var user = new User("replica-user", "provider", "provider"); + context.Users.Add(user); + context.BaseItems.Add(new BaseItemEntity { Id = itemId, Type = typeof(AudioBook).FullName! }); + await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); + + return user; + } + } + + private static JellyfinDbContext CreateContext(NpgsqlDataSource dataSource) + { + var optionsBuilder = new DbContextOptionsBuilder(); + var provider = new PostgreSqlDatabaseProvider(dataSource); + provider.Initialise(optionsBuilder, new DatabaseConfigurationOptions { DatabaseType = "PostgreSQL" }); + return new JellyfinDbContext( + optionsBuilder.Options, + NullLogger.Instance, + provider, + new NoLockBehavior(NullLogger.Instance)); + } + + /// + /// Hands every its own context over the one shared database, the way the + /// pooled factory does in the server. + /// + private sealed class DataSourceContextFactory : IDbContextFactory + { + private readonly NpgsqlDataSource _dataSource; + + public DataSourceContextFactory(NpgsqlDataSource dataSource) + { + _dataSource = dataSource; + } + + public JellyfinDbContext CreateDbContext() => CreateContext(_dataSource); + } +}