perf(nextup): batch the user data reads the next episode selection makes
Next Up walked every series in an unpaginated loop and read user data one episode at a time, so the Home screen row cost two queries per series once the cache was gone. - read the played state of every candidate episode in one query - read the versions the resume check and the last played date need in one query each - default PrefetchedUserData on IUserBaseItemComparer so plugin comparers still compile - cover the bounded query count and the plugin comparer with tests - share one database across the replica tests
This commit is contained in:
@@ -13,6 +13,7 @@ using MediaBrowser.Controller.Configuration;
|
||||
using MediaBrowser.Controller.Dto;
|
||||
using MediaBrowser.Controller.Entities;
|
||||
using MediaBrowser.Controller.Library;
|
||||
using MediaBrowser.Controller.Persistence;
|
||||
using MediaBrowser.Controller.TV;
|
||||
using MediaBrowser.Model.Querying;
|
||||
using Episode = MediaBrowser.Controller.Entities.TV.Episode;
|
||||
@@ -124,53 +125,100 @@ namespace Emby.Server.Implementations.TV
|
||||
|
||||
var batchResult = _libraryManager.GetNextUpEpisodesBatch(query, seriesKeys, includeSpecials, includeRewatching);
|
||||
|
||||
var nextUpList = new List<(DateTime LastWatchedDate, Episode Episode)>();
|
||||
|
||||
var results = new List<NextUpEpisodeBatchResult>(seriesKeys.Count);
|
||||
foreach (var seriesKey in seriesKeys)
|
||||
{
|
||||
if (!batchResult.TryGetValue(seriesKey, out var result))
|
||||
if (batchResult.TryGetValue(seriesKey, out var result))
|
||||
{
|
||||
continue;
|
||||
results.Add(result);
|
||||
}
|
||||
}
|
||||
|
||||
var nextEpisode = DetermineNextEpisode(result, user, includeSpecials, request.EnableResumable, false);
|
||||
// The selection below tests the played state of every episode it considers, so read the whole
|
||||
// series batch in one query rather than one query per series.
|
||||
var selectionCandidates = new List<BaseItem>();
|
||||
foreach (var result in results)
|
||||
{
|
||||
AddCandidate(selectionCandidates, result.NextUp);
|
||||
AddCandidate(selectionCandidates, result.LastWatched);
|
||||
AddCandidate(selectionCandidates, result.NextPlayedForRewatching);
|
||||
AddCandidate(selectionCandidates, result.LastWatchedForRewatching);
|
||||
|
||||
if (result.Specials is not null)
|
||||
{
|
||||
selectionCandidates.AddRange(result.Specials);
|
||||
}
|
||||
}
|
||||
|
||||
var selectionUserData = _userDataManager.GetUserDataBatch(selectionCandidates, user);
|
||||
|
||||
var candidates = new List<NextUpCandidate>();
|
||||
foreach (var result in results)
|
||||
{
|
||||
var nextEpisode = SelectNextEpisode(result, user, includeSpecials, includePlayed: false, selectionUserData);
|
||||
if (nextEpisode is not null)
|
||||
{
|
||||
// The last played date and the version that was actually played live on the version item's user data
|
||||
// The played state propagated to the sibling versions carries no date
|
||||
var (playedVersion, lastPlayedDate) = GetMostRecentlyPlayedVersion(result.LastWatched, user);
|
||||
nextEpisode = GetPreferredVersion(nextEpisode, result.LastWatched, playedVersion);
|
||||
|
||||
DateTime lastWatchedDate = DateTime.MinValue;
|
||||
if (result.LastWatched is not null)
|
||||
{
|
||||
lastWatchedDate = lastPlayedDate ?? DateTime.MinValue.AddDays(1);
|
||||
}
|
||||
|
||||
nextUpList.Add((lastWatchedDate, nextEpisode));
|
||||
candidates.Add(new NextUpCandidate(nextEpisode, result.LastWatched, !request.EnableResumable));
|
||||
}
|
||||
|
||||
if (includeRewatching)
|
||||
{
|
||||
var nextPlayedEpisode = DetermineNextEpisodeForRewatching(result, user, includeSpecials);
|
||||
|
||||
var nextPlayedEpisode = SelectNextEpisode(result, user, includeSpecials, includePlayed: true, selectionUserData);
|
||||
if (nextPlayedEpisode is not null)
|
||||
{
|
||||
var (playedVersion, lastPlayedDate) = GetMostRecentlyPlayedVersion(result.LastWatchedForRewatching, user);
|
||||
nextPlayedEpisode = GetPreferredVersion(nextPlayedEpisode, result.LastWatchedForRewatching, playedVersion);
|
||||
|
||||
DateTime rewatchLastWatchedDate = DateTime.MinValue;
|
||||
if (result.LastWatchedForRewatching is not null)
|
||||
{
|
||||
rewatchLastWatchedDate = lastPlayedDate ?? DateTime.MinValue.AddDays(1);
|
||||
}
|
||||
|
||||
nextUpList.Add((rewatchLastWatchedDate, nextPlayedEpisode));
|
||||
// A rewatch suggestion is dropped once it has been resumed, whatever the request asked for.
|
||||
candidates.Add(new NextUpCandidate(nextPlayedEpisode, result.LastWatchedForRewatching, true));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The resume progress may live on an alternate version, so read every version in one query.
|
||||
var episodeVersions = new List<BaseItem>();
|
||||
foreach (var candidate in candidates)
|
||||
{
|
||||
if (candidate.DropWhenResumed)
|
||||
{
|
||||
candidate.EpisodeVersions = candidate.Episode.GetAllVersions();
|
||||
episodeVersions.AddRange(candidate.EpisodeVersions);
|
||||
}
|
||||
}
|
||||
|
||||
if (episodeVersions.Count > 0)
|
||||
{
|
||||
var resumeUserData = _userDataManager.GetUserDataBatch(episodeVersions, user);
|
||||
candidates.RemoveAll(candidate => candidate.EpisodeVersions
|
||||
.Any(version => GetUserData(user, version, resumeUserData)?.PlaybackPositionTicks > 0));
|
||||
}
|
||||
|
||||
// The last played date and the version that was actually played live on the version item's user data
|
||||
// The played state propagated to the sibling versions carries no date
|
||||
var lastWatchedVersions = new List<BaseItem>();
|
||||
foreach (var candidate in candidates)
|
||||
{
|
||||
if (candidate.LastWatched is Video lastWatchedVideo)
|
||||
{
|
||||
candidate.LastWatchedVersions = lastWatchedVideo.GetAllVersions();
|
||||
lastWatchedVersions.AddRange(candidate.LastWatchedVersions);
|
||||
}
|
||||
}
|
||||
|
||||
var lastWatchedUserData = _userDataManager.GetUserDataBatch(lastWatchedVersions, user);
|
||||
|
||||
var nextUpList = new List<(DateTime LastWatchedDate, Episode Episode)>(candidates.Count);
|
||||
foreach (var candidate in candidates)
|
||||
{
|
||||
var (playedVersion, lastPlayedDate) = GetMostRecentlyPlayedVersion(candidate.LastWatchedVersions, user, lastWatchedUserData);
|
||||
var nextEpisode = GetPreferredVersion(candidate.Episode, candidate.LastWatched, playedVersion);
|
||||
|
||||
DateTime lastWatchedDate = DateTime.MinValue;
|
||||
if (candidate.LastWatched is not null)
|
||||
{
|
||||
lastWatchedDate = lastPlayedDate ?? DateTime.MinValue.AddDays(1);
|
||||
}
|
||||
|
||||
nextUpList.Add((lastWatchedDate, nextEpisode));
|
||||
}
|
||||
|
||||
var sortedEpisodes = nextUpList
|
||||
.OrderByDescending(x => x.LastWatchedDate)
|
||||
.Select(x => (BaseItem)x.Episode);
|
||||
@@ -178,12 +226,25 @@ namespace Emby.Server.Implementations.TV
|
||||
return GetResult(sortedEpisodes, request);
|
||||
}
|
||||
|
||||
private Episode? DetermineNextEpisode(
|
||||
MediaBrowser.Controller.Persistence.NextUpEpisodeBatchResult result,
|
||||
private static void AddCandidate(List<BaseItem> candidates, BaseItem? item)
|
||||
{
|
||||
if (item is not null)
|
||||
{
|
||||
candidates.Add(item);
|
||||
}
|
||||
}
|
||||
|
||||
private UserItemData? GetUserData(User user, BaseItem item, IReadOnlyDictionary<Guid, UserItemData> prefetchedUserData)
|
||||
=> prefetchedUserData.TryGetValue(item.Id, out var userData)
|
||||
? userData
|
||||
: _userDataManager.GetUserData(user, item);
|
||||
|
||||
private Episode? SelectNextEpisode(
|
||||
NextUpEpisodeBatchResult result,
|
||||
User user,
|
||||
bool includeSpecials,
|
||||
bool includeResumable,
|
||||
bool includePlayed)
|
||||
bool includePlayed,
|
||||
IReadOnlyDictionary<Guid, UserItemData> prefetchedUserData)
|
||||
{
|
||||
var nextEpisode = (includePlayed ? result.NextPlayedForRewatching : result.NextUp) as Episode;
|
||||
var lastWatchedEpisode = (includePlayed ? result.LastWatchedForRewatching : result.LastWatched) as Episode;
|
||||
@@ -217,60 +278,41 @@ namespace Emby.Server.Implementations.TV
|
||||
|
||||
if (!includePlayed)
|
||||
{
|
||||
sortedEpisodes = sortedEpisodes.Where(episode => _userDataManager.GetUserData(user, episode) is not { Played: true });
|
||||
sortedEpisodes = sortedEpisodes.Where(episode => GetUserData(user, episode, prefetchedUserData) is not { Played: true });
|
||||
}
|
||||
|
||||
nextEpisode = sortedEpisodes.FirstOrDefault();
|
||||
}
|
||||
}
|
||||
|
||||
if (nextEpisode is not null && !includeResumable)
|
||||
{
|
||||
// The resume progress may live on an alternate version
|
||||
foreach (var version in nextEpisode.GetAllVersions())
|
||||
{
|
||||
if (_userDataManager.GetUserData(user, version)?.PlaybackPositionTicks > 0)
|
||||
{
|
||||
return null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return nextEpisode;
|
||||
}
|
||||
|
||||
private Episode? DetermineNextEpisodeForRewatching(
|
||||
MediaBrowser.Controller.Persistence.NextUpEpisodeBatchResult result,
|
||||
User user,
|
||||
bool includeSpecials)
|
||||
{
|
||||
return DetermineNextEpisode(result, user, includeSpecials, includeResumable: false, includePlayed: true);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Gets the version of the last watched episode that was actually played, together with its last played date.
|
||||
/// The version that was played carries the most recent LastPlayedDate.
|
||||
/// dates.
|
||||
/// </summary>
|
||||
/// <param name="lastWatched">The last watched episode (any version).</param>
|
||||
/// <param name="versions">The versions of the last watched episode.</param>
|
||||
/// <param name="user">The user.</param>
|
||||
/// <param name="prefetchedUserData">User data read for every version up front.</param>
|
||||
/// <returns>The played version and its last played date.</returns>
|
||||
private (Video? PlayedVersion, DateTime? LastPlayedDate) GetMostRecentlyPlayedVersion(BaseItem? lastWatched, User user)
|
||||
private (Video? PlayedVersion, DateTime? LastPlayedDate) GetMostRecentlyPlayedVersion(
|
||||
IReadOnlyList<Video> versions,
|
||||
User user,
|
||||
IReadOnlyDictionary<Guid, UserItemData> prefetchedUserData)
|
||||
{
|
||||
if (lastWatched is not Video lastWatchedVideo)
|
||||
if (versions.Count == 0)
|
||||
{
|
||||
return (null, null);
|
||||
}
|
||||
|
||||
var versions = lastWatchedVideo.GetAllVersions();
|
||||
var userDataByVersion = _userDataManager.GetUserDataBatch(versions, user);
|
||||
|
||||
var playedVersion = VersionPlaybackSelector.SelectMostRecentlyPlayed(
|
||||
versions,
|
||||
version => userDataByVersion.GetValueOrDefault(version.Id),
|
||||
version => GetUserData(user, version, prefetchedUserData),
|
||||
data => data.LastPlayedDate.HasValue);
|
||||
|
||||
return (playedVersion, playedVersion is null ? null : userDataByVersion[playedVersion.Id].LastPlayedDate);
|
||||
return (playedVersion, playedVersion is null ? null : GetUserData(user, playedVersion, prefetchedUserData)?.LastPlayedDate);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
@@ -346,5 +388,28 @@ namespace Emby.Server.Implementations.TV
|
||||
totalCount,
|
||||
items.ToArray());
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// An episode picked for Next Up, together with the versions its user data is read from.
|
||||
/// </summary>
|
||||
private sealed class NextUpCandidate
|
||||
{
|
||||
public NextUpCandidate(Episode episode, BaseItem? lastWatched, bool dropWhenResumed)
|
||||
{
|
||||
Episode = episode;
|
||||
LastWatched = lastWatched;
|
||||
DropWhenResumed = dropWhenResumed;
|
||||
}
|
||||
|
||||
public Episode Episode { get; }
|
||||
|
||||
public BaseItem? LastWatched { get; }
|
||||
|
||||
public bool DropWhenResumed { get; }
|
||||
|
||||
public IReadOnlyList<Video> EpisodeVersions { get; set; } = [];
|
||||
|
||||
public IReadOnlyList<Video> LastWatchedVersions { get; set; } = [];
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -33,8 +33,13 @@ namespace MediaBrowser.Controller.Sorting
|
||||
|
||||
/// <summary>
|
||||
/// Gets or sets user data for the items being sorted, keyed by item id, read once up front.
|
||||
/// A comparer that does not store it reads its user data one item at a time instead.
|
||||
/// </summary>
|
||||
/// <value>The prefetched user data, or <c>null</c> when none was prefetched.</value>
|
||||
IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
|
||||
IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData
|
||||
{
|
||||
get => null;
|
||||
set { }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@ using Emby.Naming.Common;
|
||||
using Emby.Server.Implementations.Library;
|
||||
using Emby.Server.Implementations.Sorting;
|
||||
using Jellyfin.Data.Enums;
|
||||
using Jellyfin.Database.Implementations.Entities;
|
||||
using Jellyfin.Database.Implementations.Enums;
|
||||
using MediaBrowser.Controller.Configuration;
|
||||
using MediaBrowser.Controller.Entities;
|
||||
@@ -63,13 +64,49 @@ public class LibraryManagerSortTests
|
||||
Assert.Equal(new[] { "Alpha", "Mike", "Zulu" }, sorted.Select(i => i.Name));
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public void Sort_ComparerThatIgnoresPrefetchedUserData_StillSortsFromLiveReads()
|
||||
{
|
||||
var alpha = new Audio { Name = "Alpha", SortName = "Alpha", Id = Guid.NewGuid() };
|
||||
var zulu = new Audio { Name = "Zulu", SortName = "Zulu", Id = Guid.NewGuid() };
|
||||
var playCounts = new Dictionary<Guid, int> { [alpha.Id] = 1, [zulu.Id] = 9 };
|
||||
|
||||
var userDataManager = new Mock<IUserDataManager>();
|
||||
userDataManager
|
||||
.Setup(u => u.GetUserData(It.IsAny<User>(), It.IsAny<BaseItem>()))
|
||||
.Returns<User, BaseItem>((_, item) => new UserItemData { Key = item.Id.ToString("N"), PlayCount = playCounts[item.Id] });
|
||||
userDataManager
|
||||
.Setup(u => u.GetUserDataBatch(It.IsAny<IReadOnlyList<BaseItem>>(), It.IsAny<User>()))
|
||||
.Returns(new Dictionary<Guid, UserItemData>());
|
||||
|
||||
var libraryManager = CreateLibraryManager(
|
||||
new IBaseItemComparer[] { new PluginPlayCountComparer() },
|
||||
userDataManager);
|
||||
|
||||
var sorted = libraryManager.Sort(
|
||||
new BaseItem[] { alpha, zulu },
|
||||
new User("sorter", "provider", "provider"),
|
||||
new[] { (ItemSortBy.PlayCount, SortOrder.Descending) }).ToArray();
|
||||
|
||||
Assert.Equal(new[] { "Zulu", "Alpha" }, sorted.Select(i => i.Name));
|
||||
userDataManager.Verify(u => u.GetUserData(It.IsAny<User>(), It.IsAny<BaseItem>()), Times.AtLeastOnce);
|
||||
}
|
||||
|
||||
private static Folder MakeFolder(string name, DateTime dateLastMediaAdded)
|
||||
=> new() { Name = name, Id = Guid.NewGuid(), DateLastMediaAdded = dateLastMediaAdded };
|
||||
|
||||
private static Emby.Server.Implementations.Library.LibraryManager CreateLibraryManager(IReadOnlyCollection<IBaseItemComparer> comparers)
|
||||
private static Emby.Server.Implementations.Library.LibraryManager CreateLibraryManager(
|
||||
IReadOnlyCollection<IBaseItemComparer> comparers,
|
||||
Mock<IUserDataManager>? userDataManager = null)
|
||||
{
|
||||
var fixture = new Fixture().Customize(new AutoMoqCustomization());
|
||||
fixture.Register(() => new NamingOptions());
|
||||
|
||||
if (userDataManager is not null)
|
||||
{
|
||||
fixture.Inject(userDataManager.Object);
|
||||
}
|
||||
|
||||
var configMock = fixture.Freeze<Mock<IServerConfigurationManager>>();
|
||||
configMock.Setup(c => c.ApplicationPaths.ProgramDataPath).Returns("/data");
|
||||
BaseItem.ConfigurationManager ??= configMock.Object;
|
||||
@@ -86,4 +123,22 @@ public class LibraryManagerSortTests
|
||||
fixture.Create<IEnumerable<ILibraryPostScanTask>>()))
|
||||
.Create();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// A comparer of the shape a third-party plugin ships: it implements
|
||||
/// <see cref="IUserBaseItemComparer"/> without ever mentioning PrefetchedUserData.
|
||||
/// </summary>
|
||||
public sealed class PluginPlayCountComparer : IUserBaseItemComparer
|
||||
{
|
||||
public User User { get; set; } = null!;
|
||||
|
||||
public IUserManager UserManager { get; set; } = null!;
|
||||
|
||||
public IUserDataManager UserDataManager { get; set; } = null!;
|
||||
|
||||
public ItemSortBy Type => ItemSortBy.PlayCount;
|
||||
|
||||
public int Compare(BaseItem? x, BaseItem? y)
|
||||
=> UserDataManager.GetUserData(User, x!)!.PlayCount.CompareTo(UserDataManager.GetUserData(User, y!)!.PlayCount);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,94 @@
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Globalization;
|
||||
using System.Linq;
|
||||
using Emby.Server.Implementations.TV;
|
||||
using Jellyfin.Database.Implementations.Entities;
|
||||
using MediaBrowser.Controller.Configuration;
|
||||
using MediaBrowser.Controller.Dto;
|
||||
using MediaBrowser.Controller.Entities;
|
||||
using MediaBrowser.Controller.Entities.TV;
|
||||
using MediaBrowser.Controller.Library;
|
||||
using MediaBrowser.Controller.Persistence;
|
||||
using MediaBrowser.Model.Configuration;
|
||||
using MediaBrowser.Model.Querying;
|
||||
using Moq;
|
||||
using Xunit;
|
||||
|
||||
namespace Jellyfin.Server.Implementations.Tests.TV;
|
||||
|
||||
public class TVSeriesManagerNextUpTests
|
||||
{
|
||||
[Theory]
|
||||
[InlineData(1)]
|
||||
[InlineData(25)]
|
||||
[InlineData(200)]
|
||||
public void GetNextUp_ReadsUserDataInABoundedNumberOfQueries(int seriesCount)
|
||||
{
|
||||
var user = new User("next-up", "provider", "provider");
|
||||
var libraryManager = new Mock<ILibraryManager>();
|
||||
var userDataManager = new Mock<IUserDataManager>();
|
||||
|
||||
var seriesKeys = Enumerable.Range(0, seriesCount)
|
||||
.Select(i => i.ToString(CultureInfo.InvariantCulture))
|
||||
.ToList();
|
||||
|
||||
var batch = seriesKeys.ToDictionary(
|
||||
key => key,
|
||||
key => new NextUpEpisodeBatchResult
|
||||
{
|
||||
NextUp = new Episode { Id = Guid.NewGuid(), Name = "Next " + key },
|
||||
LastWatched = new Episode { Id = Guid.NewGuid(), Name = "Watched " + key }
|
||||
});
|
||||
|
||||
libraryManager
|
||||
.Setup(l => l.GetNextUpSeriesKeys(It.IsAny<InternalItemsQuery>(), It.IsAny<IReadOnlyCollection<BaseItem>>(), It.IsAny<DateTime>()))
|
||||
.Returns(seriesKeys);
|
||||
libraryManager
|
||||
.Setup(l => l.GetNextUpEpisodesBatch(It.IsAny<InternalItemsQuery>(), It.IsAny<IReadOnlyList<string>>(), It.IsAny<bool>(), It.IsAny<bool>()))
|
||||
.Returns(batch);
|
||||
libraryManager.Setup(l => l.GetLinkedAlternateVersions(It.IsAny<Video>())).Returns([]);
|
||||
libraryManager.Setup(l => l.GetLocalAlternateVersionIds(It.IsAny<Video>())).Returns([]);
|
||||
|
||||
var batchReads = 0;
|
||||
userDataManager
|
||||
.Setup(u => u.GetUserDataBatch(It.IsAny<IReadOnlyList<BaseItem>>(), It.IsAny<User>()))
|
||||
.Returns<IReadOnlyList<BaseItem>, User>((items, _) =>
|
||||
{
|
||||
batchReads++;
|
||||
return items.DistinctBy(i => i.Id).ToDictionary(
|
||||
i => i.Id,
|
||||
i => new UserItemData { Key = i.Id.ToString("N", CultureInfo.InvariantCulture) });
|
||||
});
|
||||
|
||||
var previousLibraryManager = BaseItem.LibraryManager;
|
||||
BaseItem.LibraryManager = libraryManager.Object;
|
||||
|
||||
try
|
||||
{
|
||||
var manager = new TVSeriesManager(userDataManager.Object, libraryManager.Object, CreateConfigurationManager());
|
||||
|
||||
var result = manager.GetNextUp(
|
||||
new NextUpQuery { User = user, EnableTotalRecordCount = true },
|
||||
[],
|
||||
new DtoOptions(false));
|
||||
|
||||
Assert.Equal(seriesCount, result.TotalRecordCount);
|
||||
|
||||
// Selection, the resume check and the last played date: three reads whatever the library holds.
|
||||
Assert.Equal(3, batchReads);
|
||||
userDataManager.Verify(u => u.GetUserData(It.IsAny<User>(), It.IsAny<BaseItem>()), Times.Never);
|
||||
}
|
||||
finally
|
||||
{
|
||||
BaseItem.LibraryManager = previousLibraryManager;
|
||||
}
|
||||
}
|
||||
|
||||
private static IServerConfigurationManager CreateConfigurationManager()
|
||||
{
|
||||
var configurationManager = new Mock<IServerConfigurationManager>();
|
||||
configurationManager.SetupGet(c => c.Configuration).Returns(new ServerConfiguration());
|
||||
return configurationManager.Object;
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,6 @@
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Globalization;
|
||||
using System.Linq;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
@@ -28,22 +29,15 @@ namespace Jellyfin.Server.Tests.Library;
|
||||
/// 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
|
||||
public sealed class UserDataManagerReplicaTests : IClassFixture<UserDataManagerReplicaTests.DatabaseFixture>
|
||||
{
|
||||
private static readonly long _quarterIn = TimeSpan.FromMinutes(20).Ticks;
|
||||
|
||||
private PostgreSqlTestServer _server = null!;
|
||||
private readonly NpgsqlDataSource _dataSource;
|
||||
|
||||
/// <inheritdoc/>
|
||||
public async ValueTask InitializeAsync()
|
||||
public UserDataManagerReplicaTests(DatabaseFixture fixture)
|
||||
{
|
||||
_server = await PostgreSqlTestServer.StartAsync().ConfigureAwait(false);
|
||||
}
|
||||
|
||||
/// <inheritdoc/>
|
||||
public async ValueTask DisposeAsync()
|
||||
{
|
||||
await _server.DisposeAsync().ConfigureAwait(false);
|
||||
_dataSource = fixture.DataSource;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
@@ -56,14 +50,11 @@ public sealed class UserDataManagerReplicaTests : IAsyncLifetime
|
||||
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 user = await CreateUserAndItemAsync(_dataSource, itemId, cancellationToken);
|
||||
|
||||
var replicaA = CreateManager(dataSource);
|
||||
var replicaB = CreateManager(dataSource);
|
||||
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" };
|
||||
|
||||
@@ -72,7 +63,7 @@ public sealed class UserDataManagerReplicaTests : IAsyncLifetime
|
||||
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);
|
||||
itemOnB.UserData = await LoadUserDataAsync(_dataSource, itemId, cancellationToken);
|
||||
|
||||
var later = replicaA.GetUserData(user, itemOnA)!;
|
||||
later.PlaybackPositionTicks = _quarterIn;
|
||||
@@ -91,14 +82,11 @@ public sealed class UserDataManagerReplicaTests : IAsyncLifetime
|
||||
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 user = await CreateUserAndItemAsync(_dataSource, itemId, cancellationToken);
|
||||
|
||||
var replicaA = CreateManager(dataSource);
|
||||
var replicaB = CreateManager(dataSource);
|
||||
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" };
|
||||
|
||||
@@ -107,7 +95,7 @@ public sealed class UserDataManagerReplicaTests : IAsyncLifetime
|
||||
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);
|
||||
itemOnB.UserData = await LoadUserDataAsync(_dataSource, itemId, cancellationToken);
|
||||
|
||||
var favourited = replicaA.GetUserData(user, itemOnA)!;
|
||||
favourited.IsFavorite = true;
|
||||
@@ -131,14 +119,11 @@ public sealed class UserDataManagerReplicaTests : IAsyncLifetime
|
||||
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 user = await CreateUserAndItemAsync(_dataSource, itemId, cancellationToken);
|
||||
|
||||
var replicaA = CreateManager(dataSource);
|
||||
var replicaB = CreateManager(dataSource);
|
||||
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" };
|
||||
|
||||
@@ -146,7 +131,7 @@ public sealed class UserDataManagerReplicaTests : IAsyncLifetime
|
||||
seed.PlaybackPositionTicks = TimeSpan.FromMinutes(5).Ticks;
|
||||
replicaA.SaveUserData(user, itemOnA, seed, UserDataSaveReason.PlaybackProgress, cancellationToken);
|
||||
|
||||
itemOnB.UserData = await LoadUserDataAsync(dataSource, itemId, 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)!;
|
||||
@@ -182,14 +167,12 @@ public sealed class UserDataManagerReplicaTests : IAsyncLifetime
|
||||
}
|
||||
}
|
||||
|
||||
private static async Task<User> CreateSchemaWithUserAndItemAsync(NpgsqlDataSource dataSource, Guid itemId, CancellationToken cancellationToken)
|
||||
private static async Task<User> CreateUserAndItemAsync(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");
|
||||
var user = new User("replica-user-" + itemId.ToString("N", CultureInfo.InvariantCulture), "provider", "provider");
|
||||
context.Users.Add(user);
|
||||
context.BaseItems.Add(new BaseItemEntity { Id = itemId, Type = typeof(AudioBook).FullName! });
|
||||
await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false);
|
||||
@@ -225,4 +208,36 @@ public sealed class UserDataManagerReplicaTests : IAsyncLifetime
|
||||
|
||||
public JellyfinDbContext CreateDbContext() => CreateContext(_dataSource);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Builds the schema once for the whole class. Every test keeps to its own user and item, so one
|
||||
/// database serves all of them and the shared server is spared three schema builds.
|
||||
/// </summary>
|
||||
public sealed class DatabaseFixture : IAsyncLifetime
|
||||
{
|
||||
private PostgreSqlTestServer _server = null!;
|
||||
|
||||
public NpgsqlDataSource DataSource { get; private set; } = null!;
|
||||
|
||||
/// <inheritdoc/>
|
||||
public async ValueTask InitializeAsync()
|
||||
{
|
||||
_server = await PostgreSqlTestServer.StartAsync().ConfigureAwait(false);
|
||||
var connectionString = await _server.CreateDatabaseAsync("userdata_replica", CancellationToken.None).ConfigureAwait(false);
|
||||
DataSource = new NpgsqlDataSourceBuilder(connectionString).Build();
|
||||
|
||||
var context = CreateContext(DataSource);
|
||||
await using (context.ConfigureAwait(false))
|
||||
{
|
||||
await context.Database.EnsureCreatedAsync(CancellationToken.None).ConfigureAwait(false);
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc/>
|
||||
public async ValueTask DisposeAsync()
|
||||
{
|
||||
await DataSource.DisposeAsync().ConfigureAwait(false);
|
||||
await _server.DisposeAsync().ConfigureAwait(false);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user