fix(userdata): read user data through to the database
ci/woodpecker/pr/ci Pipeline was successful
ci/woodpecker/push/ci Pipeline was successful

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
This commit is contained in:
2026-09-20 23:51:29 +10:00
parent ec581b5e5f
commit a7919b9bac
14 changed files with 470 additions and 136 deletions
@@ -2332,7 +2332,10 @@ namespace Emby.Server.Implementations.Library
{
IOrderedEnumerable<BaseItem>? 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<BaseItem>? 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<BaseItem> PrefetchUserData(IEnumerable<BaseItem> items, User? user, IReadOnlyList<IBaseItemComparer?> comparers)
{
if (user is null)
{
return items;
}
var userComparers = comparers.OfType<IUserBaseItemComparer>().ToList();
if (userComparers.Count == 0)
{
return items;
}
var itemList = items as IReadOnlyList<BaseItem> ?? items.ToList();
var userData = _userDataManager.GetUserDataBatch(itemList, user);
foreach (var comparer in userComparers)
{
comparer.PrefetchedUserData = userData;
}
return itemList;
}
/// <summary>
/// Gets the comparer.
/// </summary>
@@ -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<JellyfinDbContext> _repository;
private readonly FastConcurrentLru<string, UserItemData> _cache;
/// <summary>
/// Initializes a new instance of the <see cref="UserDataManager"/> class.
@@ -40,7 +37,6 @@ namespace Emby.Server.Implementations.Library
{
_config = config;
_repository = repository;
_cache = new FastConcurrentLru<string, UserItemData>(Environment.ProcessorCount, _config.Configuration.CacheSize, StringComparer.OrdinalIgnoreCase);
}
/// <inheritdoc />
@@ -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
/// <inheritdoc />
public Dictionary<Guid, UserItemData> GetUserDataBatch(IReadOnlyList<BaseItem> items, User user)
{
ArgumentNullException.ThrowIfNull(items);
ArgumentNullException.ThrowIfNull(user);
var result = new Dictionary<Guid, UserItemData>(items.Count);
var itemsNeedingQuery = new List<(BaseItem Item, List<string> 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;
}
/// <summary>
/// Gets the internal key.
/// </summary>
/// <returns>System.String.</returns>
private static string GetCacheKey(long internalUserId, Guid itemId)
{
return internalUserId.ToString(CultureInfo.InvariantCulture) + "-" + itemId.ToString("N", CultureInfo.InvariantCulture);
}
/// <inheritdoc />
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();
}
}
}
@@ -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
/// <value>The user manager.</value>
public IUserManager UserManager { get; set; }
/// <summary>
/// Gets or sets the prefetched user data.
/// </summary>
/// <value>The prefetched user data.</value>
public IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
/// <summary>
/// Gets or sets the user data manager.
/// </summary>
@@ -57,7 +64,7 @@ namespace Emby.Server.Implementations.Sorting
/// <returns>DateTime.</returns>
private DateTime GetDate(BaseItem x)
{
var userdata = UserDataManager.GetUserData(User, x);
var userdata = this.GetUserData(x);
if (userdata is not null && userdata.LastPlayedDate.HasValue)
{
@@ -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
/// <value>The user manager.</value>
public IUserManager UserManager { get; set; }
/// <summary>
/// Gets or sets the prefetched user data.
/// </summary>
/// <value>The prefetched user data.</value>
public IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
/// <summary>
/// Compares the specified x.
/// </summary>
@@ -53,7 +61,7 @@ namespace Emby.Server.Implementations.Sorting
/// <returns>DateTime.</returns>
private int GetValue(BaseItem x)
{
return x.IsFavoriteOrLiked(User, userItemData: null) ? 0 : 1;
return x.IsFavoriteOrLiked(User, this.GetUserData(x)) ? 0 : 1;
}
}
}
@@ -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
/// <value>The user manager.</value>
public IUserManager UserManager { get; set; }
/// <summary>
/// Gets or sets the prefetched user data.
/// </summary>
/// <value>The prefetched user data.</value>
public IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
/// <summary>
/// Compares the specified x.
/// </summary>
@@ -54,7 +62,7 @@ namespace Emby.Server.Implementations.Sorting
/// <returns>DateTime.</returns>
private int GetValue(BaseItem x)
{
return x.IsPlayed(User, userItemData: null) ? 0 : 1;
return x.IsPlayed(User, this.GetUserData(x)) ? 0 : 1;
}
}
}
@@ -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
/// <value>The user manager.</value>
public IUserManager UserManager { get; set; }
/// <summary>
/// Gets or sets the prefetched user data.
/// </summary>
/// <value>The prefetched user data.</value>
public IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
/// <summary>
/// Compares the specified x.
/// </summary>
@@ -54,7 +62,7 @@ namespace Emby.Server.Implementations.Sorting
/// <returns>DateTime.</returns>
private int GetValue(BaseItem x)
{
return x.IsUnplayed(User, userItemData: null) ? 0 : 1;
return x.IsUnplayed(User, this.GetUserData(x)) ? 0 : 1;
}
}
}
@@ -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
/// <value>The user manager.</value>
public IUserManager UserManager { get; set; }
/// <summary>
/// Gets or sets the prefetched user data.
/// </summary>
/// <value>The prefetched user data.</value>
public IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
/// <summary>
/// Compares the specified x.
/// </summary>
@@ -56,7 +64,7 @@ namespace Emby.Server.Implementations.Sorting
/// <returns>DateTime.</returns>
private int GetValue(BaseItem x)
{
var userdata = UserDataManager.GetUserData(User, x);
var userdata = this.GetUserData(x);
return userdata is null ? 0 : userdata.PlayCount;
}
@@ -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<BaseItem> ?? 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<Folder>().Select(f => f.Id).ToList();
var filteredList = filtered.ToList();
var folderIds = filteredList.OfType<Folder>().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<Guid, UserItemData> 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<Guid, UserItemData> 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;
@@ -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
/// </summary>
/// <value>The user data repository.</value>
IUserDataManager UserDataManager { get; set; }
/// <summary>
/// Gets or sets user data for the items being sorted, keyed by item id, read once up front.
/// </summary>
/// <value>The prefetched user data, or <c>null</c> when none was prefetched.</value>
IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
}
}
@@ -0,0 +1,28 @@
#nullable disable
using MediaBrowser.Controller.Entities;
namespace MediaBrowser.Controller.Sorting
{
/// <summary>
/// Helpers shared by the comparers that sort on user data.
/// </summary>
public static class UserBaseItemComparerExtensions
{
/// <summary>
/// Gets the user data for an item, preferring the batch the sort prefetched.
/// </summary>
/// <param name="comparer">The comparer.</param>
/// <param name="item">The item.</param>
/// <returns>The item's user data.</returns>
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);
}
}
}
@@ -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();
}
+18 -3
View File
@@ -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<LiveTvProgram>()
.Select(i => _libraryManager.GetItemById(i.ChannelId))
.OfType<BaseItem>()
.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<BaseItem> 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<Guid, UserItemData> 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)
{
@@ -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<UserData>
{
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<UserData>
{
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<UserData>
{
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<UserData>();
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<UserData>
{
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);
@@ -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;
/// <summary>
/// Two independently constructed <see cref="UserDataManager"/> 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.
/// </summary>
[Trait("Category", "RequiresDocker")]
public sealed class UserDataManagerReplicaTests : IAsyncLifetime
{
private static readonly long _quarterIn = TimeSpan.FromMinutes(20).Ticks;
private PostgreSqlTestServer _server = null!;
/// <inheritdoc/>
public async ValueTask InitializeAsync()
{
_server = await PostgreSqlTestServer.StartAsync().ConfigureAwait(false);
}
/// <inheritdoc/>
public async ValueTask DisposeAsync()
{
await _server.DisposeAsync().ConfigureAwait(false);
}
/// <summary>
/// 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.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[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);
}
/// <summary>
/// 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.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[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);
}
/// <summary>
/// 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.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[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<IServerConfigurationManager>();
config.SetupGet(c => c.Configuration).Returns(new ServerConfiguration());
return new UserDataManager(config.Object, new DataSourceContextFactory(dataSource));
}
private static async Task<ICollection<UserData>> 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<User> 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<JellyfinDbContext>();
var provider = new PostgreSqlDatabaseProvider(dataSource);
provider.Initialise(optionsBuilder, new DatabaseConfigurationOptions { DatabaseType = "PostgreSQL" });
return new JellyfinDbContext(
optionsBuilder.Options,
NullLogger<JellyfinDbContext>.Instance,
provider,
new NoLockBehavior(NullLogger<NoLockBehavior>.Instance));
}
/// <summary>
/// Hands every <see cref="UserDataManager"/> its own context over the one shared database, the way the
/// pooled factory does in the server.
/// </summary>
private sealed class DataSourceContextFactory : IDbContextFactory<JellyfinDbContext>
{
private readonly NpgsqlDataSource _dataSource;
public DataSourceContextFactory(NpgsqlDataSource dataSource)
{
_dataSource = dataSource;
}
public JellyfinDbContext CreateDbContext() => CreateContext(_dataSource);
}
}