Files
unkin-agent 1c59e6afcb
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
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
2026-09-21 00:36:39 +10:00

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);
}
}
}