1c59e6afcb
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
244 lines
11 KiB
C#
244 lines
11 KiB
C#
using System;
|
|
using System.Collections.Generic;
|
|
using System.Globalization;
|
|
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 : IClassFixture<UserDataManagerReplicaTests.DatabaseFixture>
|
|
{
|
|
private static readonly long _quarterIn = TimeSpan.FromMinutes(20).Ticks;
|
|
|
|
private readonly NpgsqlDataSource _dataSource;
|
|
|
|
public UserDataManagerReplicaTests(DatabaseFixture fixture)
|
|
{
|
|
_dataSource = fixture.DataSource;
|
|
}
|
|
|
|
/// <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 itemId = Guid.NewGuid();
|
|
var user = await CreateUserAndItemAsync(_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 itemId = Guid.NewGuid();
|
|
var user = await CreateUserAndItemAsync(_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 itemId = Guid.NewGuid();
|
|
var user = await CreateUserAndItemAsync(_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> CreateUserAndItemAsync(NpgsqlDataSource dataSource, Guid itemId, CancellationToken cancellationToken)
|
|
{
|
|
var context = CreateContext(dataSource);
|
|
await using (context.ConfigureAwait(false))
|
|
{
|
|
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);
|
|
|
|
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);
|
|
}
|
|
|
|
/// <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);
|
|
}
|
|
}
|
|
}
|