From 56919b956555513fa25d6142f9d940c4acfd28ba Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Sun, 13 Sep 2026 13:08:47 +1000 Subject: [PATCH 1/3] fix(devices): read devices and device options through to the database The device cache was filled once at construction, so a token minted by one replica was unknown to every other replica already running and a token revoked on one replica stayed valid on the others until they restarted. - drop the eager device and device options dictionaries - read devices and device options from the database on every query - push the device query filters and ordering into SQL - open the request's database context only for the api key fallback - cover both directions against real PostgreSQL with two manager instances --- .../Session/SessionManager.cs | 26 +-- Jellyfin.Api/Controllers/DevicesController.cs | 21 +- .../Devices/DeviceManager.cs | 175 +++++++++------- .../Security/AuthorizationContext.cs | 144 ++++++------- .../Users/DeviceAccessHost.cs | 4 +- .../Devices/IDeviceManager.cs | 16 +- .../Devices/DeviceManagerReplicaTests.cs | 197 ++++++++++++++++++ 7 files changed, 406 insertions(+), 177 deletions(-) create mode 100644 tests/Jellyfin.Server.Tests/Devices/DeviceManagerReplicaTests.cs diff --git a/Emby.Server.Implementations/Session/SessionManager.cs b/Emby.Server.Implementations/Session/SessionManager.cs index 94215bed79..13bf42f437 100644 --- a/Emby.Server.Implementations/Session/SessionManager.cs +++ b/Emby.Server.Implementations/Session/SessionManager.cs @@ -256,7 +256,7 @@ namespace Emby.Server.Implementations.Session ArgumentException.ThrowIfNullOrEmpty(deviceId); var activityDate = DateTime.UtcNow; - var session = GetSessionInfo(appName, appVersion, deviceId, deviceName, remoteEndPoint, user); + var session = await GetSessionInfo(appName, appVersion, deviceId, deviceName, remoteEndPoint, user).ConfigureAwait(false); var lastActivityDate = session.LastActivityDate; session.LastActivityDate = activityDate; @@ -491,7 +491,7 @@ namespace Emby.Server.Implementations.Session /// The remote end point. /// The user. /// SessionInfo. - private SessionInfo GetSessionInfo( + private async Task GetSessionInfo( string appName, string appVersion, string deviceId, @@ -504,7 +504,7 @@ namespace Emby.Server.Implementations.Session ArgumentException.ThrowIfNullOrEmpty(deviceId); var key = GetSessionKey(appName, deviceId, user?.Id ?? Guid.Empty); - SessionInfo newSession = CreateSessionInfo(key, appName, appVersion, deviceId, deviceName, remoteEndPoint, user); + SessionInfo newSession = await CreateSessionInfo(key, appName, appVersion, deviceId, deviceName, remoteEndPoint, user).ConfigureAwait(false); SessionInfo sessionInfo = _activeConnections.GetOrAdd(key, newSession); if (ReferenceEquals(newSession, sessionInfo)) { @@ -532,7 +532,7 @@ namespace Emby.Server.Implementations.Session return sessionInfo; } - private SessionInfo CreateSessionInfo( + private async Task CreateSessionInfo( string key, string appName, string appVersion, @@ -562,7 +562,7 @@ namespace Emby.Server.Implementations.Session deviceName = "Network Device"; } - var deviceOptions = _deviceManager.GetDeviceOptions(deviceId) ?? new() + var deviceOptions = await _deviceManager.GetDeviceOptions(deviceId).ConfigureAwait(false) ?? new() { DeviceId = deviceId }; @@ -1768,12 +1768,12 @@ namespace Emby.Server.Implementations.Session // This should be validated above, but if it isn't don't delete all tokens. ArgumentException.ThrowIfNullOrEmpty(deviceId); - var existing = _deviceManager.GetDevices( + var existing = (await _deviceManager.GetDevices( new DeviceQuery { DeviceId = deviceId, UserId = user.Id - }).Items; + }).ConfigureAwait(false)).Items; foreach (var auth in existing) { @@ -1801,12 +1801,12 @@ namespace Emby.Server.Implementations.Session ArgumentException.ThrowIfNullOrEmpty(accessToken); - var existing = _deviceManager.GetDevices( + var existing = (await _deviceManager.GetDevices( new DeviceQuery { Limit = 1, AccessToken = accessToken - }).Items; + }).ConfigureAwait(false)).Items; if (existing.Count > 0) { @@ -1845,10 +1845,10 @@ namespace Emby.Server.Implementations.Session { CheckDisposed(); - var existing = _deviceManager.GetDevices(new DeviceQuery + var existing = await _deviceManager.GetDevices(new DeviceQuery { UserId = userId - }); + }).ConfigureAwait(false); foreach (var info in existing.Items) { @@ -2045,11 +2045,11 @@ namespace Emby.Server.Implementations.Session /// public async Task GetSessionByAuthenticationToken(string token, string deviceId, string remoteEndpoint) { - var items = _deviceManager.GetDevices(new DeviceQuery + var items = (await _deviceManager.GetDevices(new DeviceQuery { AccessToken = token, Limit = 1 - }).Items; + }).ConfigureAwait(false)).Items; if (items.Count == 0) { diff --git a/Jellyfin.Api/Controllers/DevicesController.cs b/Jellyfin.Api/Controllers/DevicesController.cs index 2bbfeb40b8..59244588bb 100644 --- a/Jellyfin.Api/Controllers/DevicesController.cs +++ b/Jellyfin.Api/Controllers/DevicesController.cs @@ -50,10 +50,10 @@ public class DevicesController : BaseJellyfinApiController /// An containing the list of devices. [HttpGet] [ProducesResponseType(StatusCodes.Status200OK)] - public ActionResult> GetDevices([FromQuery] Guid? userId) + public async Task>> GetDevices([FromQuery] Guid? userId) { userId = RequestHelpers.GetUserId(User, userId); - return _deviceManager.GetDevicesForUser(userId); + return await _deviceManager.GetDevicesForUser(userId).ConfigureAwait(false); } /// @@ -66,9 +66,9 @@ public class DevicesController : BaseJellyfinApiController [HttpGet("Info")] [ProducesResponseType(StatusCodes.Status200OK)] [ProducesResponseType(StatusCodes.Status404NotFound)] - public ActionResult GetDeviceInfo([FromQuery, Required] string id) + public async Task> GetDeviceInfo([FromQuery, Required] string id) { - var deviceInfo = _deviceManager.GetDevice(id); + var deviceInfo = await _deviceManager.GetDevice(id).ConfigureAwait(false); if (deviceInfo is null) { return NotFound(); @@ -87,9 +87,9 @@ public class DevicesController : BaseJellyfinApiController [HttpGet("Options")] [ProducesResponseType(StatusCodes.Status200OK)] [ProducesResponseType(StatusCodes.Status404NotFound)] - public ActionResult GetDeviceOptions([FromQuery, Required] string id) + public async Task> GetDeviceOptions([FromQuery, Required] string id) { - var deviceInfo = _deviceManager.GetDeviceOptions(id); + var deviceInfo = await _deviceManager.GetDeviceOptions(id).ConfigureAwait(false); if (deviceInfo is null) { return NotFound(); @@ -127,7 +127,12 @@ public class DevicesController : BaseJellyfinApiController [ProducesResponseType(StatusCodes.Status400BadRequest)] public async Task DeleteDevice([FromQuery] string[] id) { - var devices = id.Select(_deviceManager.GetDevice).ToArray(); + var devices = new List(id.Length); + foreach (var deviceId in id) + { + devices.Add(await _deviceManager.GetDevice(deviceId).ConfigureAwait(false)); + } + if (devices.Any(f => f is null)) { return BadRequest(); @@ -135,7 +140,7 @@ public class DevicesController : BaseJellyfinApiController foreach (var device in devices) { - var sessions = _deviceManager.GetDevices(new DeviceQuery { DeviceId = device!.Id }); + var sessions = await _deviceManager.GetDevices(new DeviceQuery { DeviceId = device!.Id }).ConfigureAwait(false); foreach (var session in sessions.Items) { diff --git a/Jellyfin.Server.Implementations/Devices/DeviceManager.cs b/Jellyfin.Server.Implementations/Devices/DeviceManager.cs index d0d52a23fb..672693138f 100644 --- a/Jellyfin.Server.Implementations/Devices/DeviceManager.cs +++ b/Jellyfin.Server.Implementations/Devices/DeviceManager.cs @@ -31,8 +31,6 @@ namespace Jellyfin.Server.Implementations.Devices private readonly IDbContextFactory _dbProvider; private readonly IUserManager _userManager; private readonly ConcurrentDictionary _capabilitiesMap = new(); - private readonly ConcurrentDictionary _devices; - private readonly ConcurrentDictionary _deviceOptions; /// /// Initializes a new instance of the class. @@ -43,23 +41,6 @@ namespace Jellyfin.Server.Implementations.Devices { _dbProvider = dbProvider; _userManager = userManager; - _devices = new ConcurrentDictionary(); - _deviceOptions = new ConcurrentDictionary(); - - using var dbContext = _dbProvider.CreateDbContext(); - foreach (var device in dbContext.Devices - .OrderBy(d => d.Id) - .AsEnumerable()) - { - _devices.TryAdd(device.Id, device); - } - - foreach (var deviceOption in dbContext.DeviceOptions - .OrderBy(d => d.Id) - .AsEnumerable()) - { - _deviceOptions.TryAdd(deviceOption.DeviceId, deviceOption); - } } /// @@ -89,8 +70,6 @@ namespace Jellyfin.Server.Implementations.Devices await dbContext.SaveChangesAsync().ConfigureAwait(false); } - _deviceOptions[deviceId] = deviceOptions; - DeviceOptionsUpdated?.Invoke(this, new GenericEventArgs>(new Tuple(deviceId, deviceOptions))); } @@ -102,21 +81,24 @@ namespace Jellyfin.Server.Implementations.Devices { dbContext.Devices.Add(device); await dbContext.SaveChangesAsync().ConfigureAwait(false); - _devices.TryAdd(device.Id, device); } return device; } /// - public DeviceOptionsDto? GetDeviceOptions(string deviceId) + public async Task GetDeviceOptions(string deviceId) { - if (_deviceOptions.TryGetValue(deviceId, out var deviceOptions)) + var dbContext = await _dbProvider.CreateDbContextAsync().ConfigureAwait(false); + await using (dbContext.ConfigureAwait(false)) { - return ToDeviceOptionsDto(deviceOptions); - } + var deviceOptions = await dbContext.DeviceOptions + .AsNoTracking() + .FirstOrDefaultAsync(dev => dev.DeviceId == deviceId) + .ConfigureAwait(false); - return null; + return deviceOptions is null ? null : ToDeviceOptionsDto(deviceOptions); + } } /// @@ -133,43 +115,79 @@ namespace Jellyfin.Server.Implementations.Devices } /// - public DeviceInfoDto? GetDevice(string id) + public async Task GetDevice(string id) { - var device = _devices.Values.Where(d => d.DeviceId == id).OrderByDescending(d => d.DateLastActivity).FirstOrDefault(); - _deviceOptions.TryGetValue(id, out var deviceOption); + var dbContext = await _dbProvider.CreateDbContextAsync().ConfigureAwait(false); + await using (dbContext.ConfigureAwait(false)) + { + var device = await dbContext.Devices + .AsNoTracking() + .Where(d => d.DeviceId == id) + .OrderByDescending(d => d.DateLastActivity) + .FirstOrDefaultAsync() + .ConfigureAwait(false); - var deviceInfo = device is null ? null : ToDeviceInfo(device, deviceOption); - return deviceInfo is null ? null : ToDeviceInfoDto(deviceInfo); + if (device is null) + { + return null; + } + + var deviceOption = await dbContext.DeviceOptions + .AsNoTracking() + .FirstOrDefaultAsync(dev => dev.DeviceId == id) + .ConfigureAwait(false); + + return ToDeviceInfoDto(ToDeviceInfo(device, deviceOption)); + } } /// - public QueryResult GetDevices(DeviceQuery query) + public async Task> GetDevices(DeviceQuery query) { - IEnumerable devices = _devices.Values - .Where(device => !query.UserId.HasValue || device.UserId.Equals(query.UserId.Value)) - .Where(device => query.DeviceId is null || device.DeviceId == query.DeviceId) - .Where(device => query.AccessToken is null || device.AccessToken == query.AccessToken) - .OrderBy(d => d.Id) - .ToList(); - var count = devices.Count(); - - if (query.Skip.HasValue) + var dbContext = await _dbProvider.CreateDbContextAsync().ConfigureAwait(false); + await using (dbContext.ConfigureAwait(false)) { - devices = devices.Skip(query.Skip.Value); - } + IQueryable filtered = dbContext.Devices.AsNoTracking(); - if (query.Limit.HasValue && query.Limit.Value > 0) - { - devices = devices.Take(query.Limit.Value); - } + if (query.UserId.HasValue) + { + filtered = filtered.Where(device => device.UserId.Equals(query.UserId.Value)); + } - return new QueryResult(query.Skip, count, devices.ToList()); + if (query.DeviceId is not null) + { + filtered = filtered.Where(device => device.DeviceId == query.DeviceId); + } + + if (query.AccessToken is not null) + { + filtered = filtered.Where(device => device.AccessToken == query.AccessToken); + } + + // Every filter is an exact match on an indexed column and the table holds one row per + // user and device, so paging the materialised set costs less than a second round trip. + var matched = await filtered.OrderBy(d => d.Id).ToListAsync().ConfigureAwait(false); + + IEnumerable devices = matched; + + if (query.Skip.HasValue) + { + devices = devices.Skip(query.Skip.Value); + } + + if (query.Limit.HasValue && query.Limit.Value > 0) + { + devices = devices.Take(query.Limit.Value); + } + + return new QueryResult(query.Skip, matched.Count, devices.ToList()); + } } /// - public QueryResult GetDeviceInfos(DeviceQuery query) + public async Task> GetDeviceInfos(DeviceQuery query) { - var devices = GetDevices(query); + var devices = await GetDevices(query).ConfigureAwait(false); return new QueryResult( devices.StartIndex, @@ -178,38 +196,49 @@ namespace Jellyfin.Server.Implementations.Devices } /// - public QueryResult GetDevicesForUser(Guid? userId) + public async Task> GetDevicesForUser(Guid? userId) { - IEnumerable devices = _devices.Values - .OrderByDescending(d => d.DateLastActivity) - .ThenBy(d => d.DeviceId); - - if (!userId.IsNullOrEmpty()) + var dbContext = await _dbProvider.CreateDbContextAsync().ConfigureAwait(false); + await using (dbContext.ConfigureAwait(false)) { - var user = _userManager.GetUserById(userId.Value); - if (user is null) + IEnumerable devices = await dbContext.Devices + .AsNoTracking() + .OrderByDescending(d => d.DateLastActivity) + .ThenBy(d => d.DeviceId) + .ToListAsync() + .ConfigureAwait(false); + + if (!userId.IsNullOrEmpty()) { - throw new ResourceNotFoundException(); + var user = _userManager.GetUserById(userId.Value); + if (user is null) + { + throw new ResourceNotFoundException(); + } + + devices = devices.Where(i => CanAccessDevice(user, i.DeviceId)); } - devices = devices.Where(i => CanAccessDevice(user, i.DeviceId)); + var options = await dbContext.DeviceOptions + .AsNoTracking() + .ToDictionaryAsync(dev => dev.DeviceId) + .ConfigureAwait(false); + + var array = devices.Select(device => + { + options.TryGetValue(device.DeviceId, out var option); + return ToDeviceInfo(device, option); + }) + .Select(ToDeviceInfoDto) + .ToArray(); + + return new QueryResult(array); } - - var array = devices.Select(device => - { - _deviceOptions.TryGetValue(device.DeviceId, out var option); - return ToDeviceInfo(device, option); - }) - .Select(ToDeviceInfoDto) - .ToArray(); - - return new QueryResult(array); } /// public async Task DeleteDevice(Device device) { - _devices.TryRemove(device.Id, out _); var dbContext = await _dbProvider.CreateDbContextAsync().ConfigureAwait(false); await using (dbContext.ConfigureAwait(false)) { @@ -229,8 +258,6 @@ namespace Jellyfin.Server.Implementations.Devices dbContext.Devices.Update(device); await dbContext.SaveChangesAsync().ConfigureAwait(false); } - - _devices[device.Id] = device; } /// diff --git a/Jellyfin.Server.Implementations/Security/AuthorizationContext.cs b/Jellyfin.Server.Implementations/Security/AuthorizationContext.cs index 8657cb7dbb..8a50082e07 100644 --- a/Jellyfin.Server.Implementations/Security/AuthorizationContext.cs +++ b/Jellyfin.Server.Implementations/Security/AuthorizationContext.cs @@ -126,95 +126,95 @@ namespace Jellyfin.Server.Implementations.Security return authInfo; } + var device = (await _deviceManager.GetDevices( + new DeviceQuery { AccessToken = token }).ConfigureAwait(false)).Items.FirstOrDefault(); + + if (device is not null) + { + authInfo.IsAuthenticated = true; + var updateToken = false; + + // TODO: Remove these checks for IsNullOrWhiteSpace + if (string.IsNullOrWhiteSpace(authInfo.Client)) + { + authInfo.Client = device.AppName; + } + + if (string.IsNullOrWhiteSpace(authInfo.DeviceId)) + { + authInfo.DeviceId = device.DeviceId; + } + + // Temporary. TODO - allow clients to specify that the token has been shared with a casting device + var allowTokenInfoUpdate = !authInfo.Client.Contains("chromecast", StringComparison.OrdinalIgnoreCase); + + if (string.IsNullOrWhiteSpace(authInfo.Device)) + { + authInfo.Device = device.DeviceName; + } + else if (!string.Equals(authInfo.Device, device.DeviceName, StringComparison.OrdinalIgnoreCase)) + { + if (allowTokenInfoUpdate) + { + updateToken = true; + device.DeviceName = authInfo.Device; + } + } + + if (string.IsNullOrWhiteSpace(authInfo.Version)) + { + authInfo.Version = device.AppVersion; + } + else if (!string.Equals(authInfo.Version, device.AppVersion, StringComparison.OrdinalIgnoreCase)) + { + if (allowTokenInfoUpdate) + { + updateToken = true; + device.AppVersion = authInfo.Version; + } + } + + if ((DateTime.UtcNow - device.DateLastActivity).TotalMinutes > 3) + { + device.DateLastActivity = DateTime.UtcNow; + updateToken = true; + } + + authInfo.User = _userManager.GetUserById(device.UserId); + + if (updateToken) + { + await _deviceManager.UpdateDevice(device).ConfigureAwait(false); + } + + return authInfo; + } + var dbContext = await _jellyfinDbProvider.CreateDbContextAsync().ConfigureAwait(false); await using (dbContext.ConfigureAwait(false)) { - var device = _deviceManager.GetDevices( - new DeviceQuery { AccessToken = token }).Items.FirstOrDefault(); - - if (device is not null) + var key = await dbContext.ApiKeys.FirstOrDefaultAsync(apiKey => apiKey.AccessToken == token).ConfigureAwait(false); + if (key is not null) { authInfo.IsAuthenticated = true; - var updateToken = false; - - // TODO: Remove these checks for IsNullOrWhiteSpace - if (string.IsNullOrWhiteSpace(authInfo.Client)) - { - authInfo.Client = device.AppName; - } - + authInfo.Client = key.Name; + authInfo.Token = key.AccessToken; if (string.IsNullOrWhiteSpace(authInfo.DeviceId)) { - authInfo.DeviceId = device.DeviceId; + authInfo.DeviceId = _serverApplicationHost.SystemId; } - // Temporary. TODO - allow clients to specify that the token has been shared with a casting device - var allowTokenInfoUpdate = !authInfo.Client.Contains("chromecast", StringComparison.OrdinalIgnoreCase); - if (string.IsNullOrWhiteSpace(authInfo.Device)) { - authInfo.Device = device.DeviceName; - } - else if (!string.Equals(authInfo.Device, device.DeviceName, StringComparison.OrdinalIgnoreCase)) - { - if (allowTokenInfoUpdate) - { - updateToken = true; - device.DeviceName = authInfo.Device; - } + authInfo.Device = _serverApplicationHost.Name; } if (string.IsNullOrWhiteSpace(authInfo.Version)) { - authInfo.Version = device.AppVersion; - } - else if (!string.Equals(authInfo.Version, device.AppVersion, StringComparison.OrdinalIgnoreCase)) - { - if (allowTokenInfoUpdate) - { - updateToken = true; - device.AppVersion = authInfo.Version; - } + authInfo.Version = _serverApplicationHost.ApplicationVersionString; } - if ((DateTime.UtcNow - device.DateLastActivity).TotalMinutes > 3) - { - device.DateLastActivity = DateTime.UtcNow; - updateToken = true; - } - - authInfo.User = _userManager.GetUserById(device.UserId); - - if (updateToken) - { - await _deviceManager.UpdateDevice(device).ConfigureAwait(false); - } - } - else - { - var key = await dbContext.ApiKeys.FirstOrDefaultAsync(apiKey => apiKey.AccessToken == token).ConfigureAwait(false); - if (key is not null) - { - authInfo.IsAuthenticated = true; - authInfo.Client = key.Name; - authInfo.Token = key.AccessToken; - if (string.IsNullOrWhiteSpace(authInfo.DeviceId)) - { - authInfo.DeviceId = _serverApplicationHost.SystemId; - } - - if (string.IsNullOrWhiteSpace(authInfo.Device)) - { - authInfo.Device = _serverApplicationHost.Name; - } - - if (string.IsNullOrWhiteSpace(authInfo.Version)) - { - authInfo.Version = _serverApplicationHost.ApplicationVersionString; - } - - authInfo.IsApiKey = true; - } + authInfo.IsApiKey = true; } return authInfo; diff --git a/Jellyfin.Server.Implementations/Users/DeviceAccessHost.cs b/Jellyfin.Server.Implementations/Users/DeviceAccessHost.cs index 92e2bb4fa7..1f53b52f6a 100644 --- a/Jellyfin.Server.Implementations/Users/DeviceAccessHost.cs +++ b/Jellyfin.Server.Implementations/Users/DeviceAccessHost.cs @@ -61,10 +61,10 @@ public sealed class DeviceAccessHost : IHostedService private async Task UpdateDeviceAccess(User user) { - var existing = _deviceManager.GetDevices(new DeviceQuery + var existing = (await _deviceManager.GetDevices(new DeviceQuery { UserId = user.Id - }).Items; + }).ConfigureAwait(false)).Items; foreach (var device in existing) { diff --git a/MediaBrowser.Controller/Devices/IDeviceManager.cs b/MediaBrowser.Controller/Devices/IDeviceManager.cs index ea38950d32..df127e2f7b 100644 --- a/MediaBrowser.Controller/Devices/IDeviceManager.cs +++ b/MediaBrowser.Controller/Devices/IDeviceManager.cs @@ -47,29 +47,29 @@ public interface IDeviceManager /// Gets the device information. /// /// The identifier. - /// DeviceInfoDto. - DeviceInfoDto? GetDevice(string id); + /// A representing the retrieval of the device information. + Task GetDevice(string id); /// /// Gets devices based on the provided query. /// /// The device query. /// A representing the retrieval of the devices. - QueryResult GetDevices(DeviceQuery query); + Task> GetDevices(DeviceQuery query); /// /// Gets device information based on the provided query. /// /// The device query. /// A representing the retrieval of the device information. - QueryResult GetDeviceInfos(DeviceQuery query); + Task> GetDeviceInfos(DeviceQuery query); /// /// Gets the device information. /// /// The user's id, or null. - /// IEnumerable<DeviceInfoDto>. - QueryResult GetDevicesForUser(Guid? userId); + /// A representing the retrieval of the device information. + Task> GetDevicesForUser(Guid? userId); /// /// Deletes a device. @@ -105,8 +105,8 @@ public interface IDeviceManager /// Gets the options of a device. /// /// The device id. - /// of the device. - DeviceOptionsDto? GetDeviceOptions(string deviceId); + /// A representing the retrieval of the of the device. + Task GetDeviceOptions(string deviceId); /// /// Gets the dto for client capabilities. diff --git a/tests/Jellyfin.Server.Tests/Devices/DeviceManagerReplicaTests.cs b/tests/Jellyfin.Server.Tests/Devices/DeviceManagerReplicaTests.cs new file mode 100644 index 0000000000..9647a3fc09 --- /dev/null +++ b/tests/Jellyfin.Server.Tests/Devices/DeviceManagerReplicaTests.cs @@ -0,0 +1,197 @@ +using System; +using System.Threading; +using System.Threading.Tasks; +using Jellyfin.Data.Queries; +using Jellyfin.Database.Implementations; +using Jellyfin.Database.Implementations.DbConfiguration; +using Jellyfin.Database.Implementations.Entities; +using Jellyfin.Database.Implementations.Entities.Security; +using Jellyfin.Database.Implementations.Locking; +using Jellyfin.Database.Providers.PostgreSQL; +using Jellyfin.Server.Implementations.Devices; +using Jellyfin.Server.Tests.Migrations; +using MediaBrowser.Controller.Library; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Logging.Abstractions; +using Moq; +using Npgsql; +using Xunit; + +namespace Jellyfin.Server.Tests.Devices; + +/// +/// 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. +/// +[Trait("Category", "RequiresDocker")] +public sealed class DeviceManagerReplicaTests : IAsyncLifetime +{ + private PostgreSqlTestServer _server = null!; + + /// + public async ValueTask InitializeAsync() + { + _server = await PostgreSqlTestServer.StartAsync().ConfigureAwait(false); + } + + /// + public async ValueTask DisposeAsync() + { + await _server.DisposeAsync().ConfigureAwait(false); + } + + /// + /// A client that logs in against one replica has to be authenticated by every other replica that was + /// already running when the token was minted - a rolling update or a scale-up must not 401 it. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task TokenMintedOnOneReplica_IsValidOnAnother() + { + var cancellationToken = TestContext.Current.CancellationToken; + var connectionString = await _server.CreateDatabaseAsync("device_replica_create", cancellationToken); + + await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build(); + var user = await CreateSchemaWithUserAsync(dataSource, cancellationToken); + + // Both replicas start before the login, so neither has the device in hand when it is created. + var replicaA = CreateManager(dataSource, user); + var replicaB = CreateManager(dataSource, user); + + var device = await replicaA.CreateDevice(new Device(user.Id, "Jellyfin Web", "1.0.0", "Living Room TV", "device-1")); + + var seenByB = await replicaB.GetDevices(new DeviceQuery { AccessToken = device.AccessToken }); + + Assert.Equal(device.Id, Assert.Single(seenByB.Items).Id); + Assert.Equal(user.Id, seenByB.Items[0].UserId); + } + + /// + /// Revocation has to propagate at least as fast as creation: a token logged out on one replica must not + /// still authenticate on another. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task TokenRevokedOnOneReplica_IsInvalidOnAnother() + { + var cancellationToken = TestContext.Current.CancellationToken; + var connectionString = await _server.CreateDatabaseAsync("device_replica_revoke", cancellationToken); + + await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build(); + var user = await CreateSchemaWithUserAsync(dataSource, cancellationToken); + + var replicaA = CreateManager(dataSource, user); + var device = await replicaA.CreateDevice(new Device(user.Id, "Jellyfin Web", "1.0.0", "Living Room TV", "device-1")); + + // Replica B comes up after the login, so it starts out agreeing that the token is valid. + var replicaB = CreateManager(dataSource, user); + Assert.Single((await replicaB.GetDevices(new DeviceQuery { AccessToken = device.AccessToken })).Items); + + await replicaA.DeleteDevice(device); + + Assert.Empty((await replicaB.GetDevices(new DeviceQuery { AccessToken = device.AccessToken })).Items); + } + + /// + /// A device renamed on one replica has to be reported under its new name by the others, and the rename has + /// to survive being read back through a replica that never saw the write. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task DeviceOptionsWrittenOnOneReplica_AreReadOnAnother() + { + var cancellationToken = TestContext.Current.CancellationToken; + var connectionString = await _server.CreateDatabaseAsync("device_replica_options", cancellationToken); + + await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build(); + var user = await CreateSchemaWithUserAsync(dataSource, cancellationToken); + + var replicaA = CreateManager(dataSource, user); + var replicaB = CreateManager(dataSource, user); + + await replicaA.CreateDevice(new Device(user.Id, "Jellyfin Web", "1.0.0", "Living Room TV", "device-1")); + await replicaA.UpdateDeviceOptions("device-1", "Kitchen TV"); + + var options = await replicaB.GetDeviceOptions("device-1"); + Assert.Equal("Kitchen TV", options?.CustomName); + + var info = await replicaB.GetDevice("device-1"); + Assert.Equal("Kitchen TV", info?.CustomName); + } + + /// + /// Activity written by the replica serving the request has to be visible to the others, because the next + /// request from the same client can land anywhere. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task DeviceUpdatedOnOneReplica_IsReadOnAnother() + { + var cancellationToken = TestContext.Current.CancellationToken; + var connectionString = await _server.CreateDatabaseAsync("device_replica_update", cancellationToken); + + await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build(); + var user = await CreateSchemaWithUserAsync(dataSource, cancellationToken); + + var replicaA = CreateManager(dataSource, user); + var replicaB = CreateManager(dataSource, user); + + var device = await replicaA.CreateDevice(new Device(user.Id, "Jellyfin Web", "1.0.0", "Living Room TV", "device-1")); + device.AppVersion = "2.0.0"; + await replicaA.UpdateDevice(device); + + var seenByB = Assert.Single((await replicaB.GetDevices(new DeviceQuery { DeviceId = "device-1" })).Items); + Assert.Equal("2.0.0", seenByB.AppVersion); + } + + private static DeviceManager CreateManager(NpgsqlDataSource dataSource, User user) + { + var userManager = new Mock(); + userManager.Setup(manager => manager.GetUserById(user.Id)).Returns(user); + return new DeviceManager(new DataSourceContextFactory(dataSource), userManager.Object); + } + + private static async Task CreateSchemaWithUserAsync(NpgsqlDataSource dataSource, 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); + 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); + } +} From 95220e2ed64109a9953cecd3a2be11745b71425a Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Sun, 13 Sep 2026 13:10:52 +1000 Subject: [PATCH 2/3] fix(ha): read transcode store config from the variable form deployments set The startup configuration only reads JELLYFIN_ prefixed environment variables, so the bare Jellyfin__TranscodeStore__* form used by the chart, the manifests and the README is dropped and the Redis store is never registered. Nothing logs the selected store, so the fallback is invisible. - Read bare Jellyfin__* variables into the Jellyfin:* configuration root - Keep an explicit JELLYFIN_ variable winning over the bare form - Log the selected transcode session store at startup - Ping Redis once at startup and log an unreachable store at Error - Log endpoints only, never the connection string --- .../TranscodeStoreConnectivityProbe.cs | 56 ++++++++++ Jellyfin.Server/CoreAppHost.cs | 28 +---- .../ConfigurationBuilderExtensions.cs | 51 +++++++++ ...anscodeStoreServiceCollectionExtensions.cs | 87 +++++++++++++++ Jellyfin.Server/Program.cs | 2 + .../MediaEncoding/TranscodeStoreOptions.cs | 10 ++ README.md | 11 +- docs/FORK-DIFF.md | 3 + .../TranscodeStoreWiringTests.cs | 103 +++++++++++++++++ .../JellyfinSectionConfigurationTests.cs | 95 ++++++++++++++++ .../HighAvailability/RecordingLogger.cs | 42 +++++++ .../TranscodeStoreConnectivityProbeTests.cs | 75 +++++++++++++ .../TranscodeStoreRegistrationTests.cs | 104 ++++++++++++++++++ 13 files changed, 640 insertions(+), 27 deletions(-) create mode 100644 Emby.Server.Implementations/MediaEncoding/TranscodeStoreConnectivityProbe.cs create mode 100644 Jellyfin.Server/Extensions/ConfigurationBuilderExtensions.cs create mode 100644 Jellyfin.Server/Extensions/TranscodeStoreServiceCollectionExtensions.cs create mode 100644 tests/Jellyfin.Server.Implementations.Tests/MediaEncoding/TranscodeStoreWiringTests.cs create mode 100644 tests/Jellyfin.Server.Tests/HighAvailability/JellyfinSectionConfigurationTests.cs create mode 100644 tests/Jellyfin.Server.Tests/HighAvailability/RecordingLogger.cs create mode 100644 tests/Jellyfin.Server.Tests/HighAvailability/TranscodeStoreConnectivityProbeTests.cs create mode 100644 tests/Jellyfin.Server.Tests/HighAvailability/TranscodeStoreRegistrationTests.cs diff --git a/Emby.Server.Implementations/MediaEncoding/TranscodeStoreConnectivityProbe.cs b/Emby.Server.Implementations/MediaEncoding/TranscodeStoreConnectivityProbe.cs new file mode 100644 index 0000000000..2c41875ced --- /dev/null +++ b/Emby.Server.Implementations/MediaEncoding/TranscodeStoreConnectivityProbe.cs @@ -0,0 +1,56 @@ +using System; +using System.Threading; +using System.Threading.Tasks; +using MediaBrowser.Controller.MediaEncoding; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; +using StackExchange.Redis; + +namespace Emby.Server.Implementations.MediaEncoding; + +/// +/// Pings the configured Redis transcode session store once at startup so an unreachable store is +/// reported there instead of being discovered as a silent loss of cross-pod takeover. +/// +public sealed class TranscodeStoreConnectivityProbe : IHostedService +{ + private readonly IServiceProvider _serviceProvider; + private readonly ILogger _logger; + + /// + /// Initializes a new instance of the class. + /// + /// The service provider used to resolve the Redis connection. + /// The logger. + public TranscodeStoreConnectivityProbe(IServiceProvider serviceProvider, ILogger logger) + { + _serviceProvider = serviceProvider; + _logger = logger; + } + + /// + public async Task StartAsync(CancellationToken cancellationToken) + { + try + { + // Resolved here rather than injected: connecting must not be able to abort startup. + var redis = _serviceProvider.GetRequiredService(); + var roundTrip = await redis.GetDatabase().PingAsync().ConfigureAwait(false); + + _logger.LogInformation( + "Redis transcode session store is reachable ({RoundTripMs}ms round trip). HA transcode takeover is active.", + (long)roundTrip.TotalMilliseconds); + } + catch (Exception ex) + { + _logger.LogError( + ex, + "Redis transcode session store is configured but UNREACHABLE. HA transcode takeover is not working: sessions stay local to this instance and are lost when it restarts. Check {Key}.", + TranscodeStoreOptions.RedisConnectionStringKey); + } + } + + /// + public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask; +} diff --git a/Jellyfin.Server/CoreAppHost.cs b/Jellyfin.Server/CoreAppHost.cs index 0077b52bf1..1dffd75401 100644 --- a/Jellyfin.Server/CoreAppHost.cs +++ b/Jellyfin.Server/CoreAppHost.cs @@ -2,7 +2,6 @@ using System; using System.Collections.Generic; using System.Reflection; using Emby.Server.Implementations; -using Emby.Server.Implementations.MediaEncoding; using Emby.Server.Implementations.ScheduledTasks; using Emby.Server.Implementations.Session; using Jellyfin.Api.WebSocketListeners; @@ -10,6 +9,7 @@ using Jellyfin.Database.Implementations; using Jellyfin.Drawing; using Jellyfin.Drawing.Skia; using Jellyfin.LiveTv; +using Jellyfin.Server.Extensions; using Jellyfin.Server.Implementations.Activity; using Jellyfin.Server.Implementations.Devices; using Jellyfin.Server.Implementations.Events; @@ -35,7 +35,6 @@ using MediaBrowser.Providers.Lyric; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; -using StackExchange.Redis; namespace Jellyfin.Server { @@ -107,29 +106,8 @@ namespace Jellyfin.Server serviceCollection.AddScoped(); // Transcode session store: Redis-backed when configured, no-op otherwise. - serviceCollection.Configure(_startupConfig.GetSection("Jellyfin:TranscodeStore")); - var redisConnectionString = _startupConfig["Jellyfin:TranscodeStore:RedisConnectionString"]; - if (!string.IsNullOrEmpty(redisConnectionString)) - { - serviceCollection.AddSingleton(sp => - { - try - { - return ConnectionMultiplexer.Connect(redisConnectionString); - } - catch (Exception ex) - { - sp.GetRequiredService>() - .LogError(ex, "Failed to connect to Redis. Check the Jellyfin:TranscodeStore:RedisConnectionString configuration."); - throw; - } - }); - serviceCollection.AddSingleton(); - } - else - { - serviceCollection.AddSingleton(); - } + var redisConnectionString = _startupConfig[TranscodeStoreOptions.RedisConnectionStringKey]; + serviceCollection.AddTranscodeSessionStore(_startupConfig, Logger); // Scan-leader lease: gates periodic library-mutating scheduled tasks to a single leader // instance. Redis-backed when enabled and a Redis connection is configured, no-op otherwise. diff --git a/Jellyfin.Server/Extensions/ConfigurationBuilderExtensions.cs b/Jellyfin.Server/Extensions/ConfigurationBuilderExtensions.cs new file mode 100644 index 0000000000..90a5035aa4 --- /dev/null +++ b/Jellyfin.Server/Extensions/ConfigurationBuilderExtensions.cs @@ -0,0 +1,51 @@ +using System.Collections.Generic; +using System.Linq; +using Microsoft.Extensions.Configuration; + +namespace Jellyfin.Server.Extensions; + +/// +/// Extensions for building the application configuration. +/// +public static class ConfigurationBuilderExtensions +{ + /// + /// The environment variable prefix that maps onto this fork's own Jellyfin:* configuration + /// keys, for example Jellyfin__TranscodeStore__RedisConnectionString. + /// + public const string JellyfinSectionEnvironmentPrefix = "Jellyfin__"; + + /// + /// The configuration section root this fork keeps its own settings under. + /// + public const string JellyfinSectionRoot = "Jellyfin"; + + /// + /// Adds environment variables named Jellyfin__Section__Key as the configuration keys + /// Jellyfin:Section:Key. + /// + /// + /// The base configuration only reads JELLYFIN_ prefixed environment variables, so the + /// unprefixed form every manifest, chart and document uses would otherwise be dropped and the + /// feature it configures would stay off with no error. + /// + /// The configuration builder. + /// The updated configuration builder. + public static IConfigurationBuilder AddJellyfinSectionEnvironmentVariables(this IConfigurationBuilder builder) + { + // Read through the framework provider so "__" to ":" normalisation and case handling match the + // prefixed form exactly; the prefix it strips is then put back as the section root. + var scoped = new ConfigurationBuilder() + .AddEnvironmentVariables(JellyfinSectionEnvironmentPrefix) + .Build(); + + var entries = scoped.AsEnumerable() + .Where(entry => entry.Value is not null) + .Select(entry => new KeyValuePair( + ConfigurationPath.Combine(JellyfinSectionRoot, entry.Key), + entry.Value)) + .ToList(); + + return builder.AddInMemoryCollection(entries); + } +} diff --git a/Jellyfin.Server/Extensions/TranscodeStoreServiceCollectionExtensions.cs b/Jellyfin.Server/Extensions/TranscodeStoreServiceCollectionExtensions.cs new file mode 100644 index 0000000000..60ca77817d --- /dev/null +++ b/Jellyfin.Server/Extensions/TranscodeStoreServiceCollectionExtensions.cs @@ -0,0 +1,87 @@ +using System; +using System.Linq; +using Emby.Server.Implementations.MediaEncoding; +using MediaBrowser.Controller.MediaEncoding; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using StackExchange.Redis; + +namespace Jellyfin.Server.Extensions; + +/// +/// Extensions for registering the transcode session store. +/// +public static class TranscodeStoreServiceCollectionExtensions +{ + /// + /// Registers the transcode session store, Redis-backed when a connection string is configured and + /// no-op otherwise, and reports the selected store at . + /// + /// The service collection. + /// The configuration to read Jellyfin:TranscodeStore from. + /// The logger to report the selected store on. + /// The updated service collection. + public static IServiceCollection AddTranscodeSessionStore( + this IServiceCollection serviceCollection, + IConfiguration configuration, + ILogger logger) + { + ArgumentNullException.ThrowIfNull(configuration); + ArgumentNullException.ThrowIfNull(logger); + + serviceCollection.Configure(configuration.GetSection(TranscodeStoreOptions.ConfigurationSection)); + + var redisConnectionString = configuration[TranscodeStoreOptions.RedisConnectionStringKey]; + if (string.IsNullOrEmpty(redisConnectionString)) + { + logger.LogInformation( + "Transcode session store: {Store}. Cross-pod transcode takeover is off; set {Key} to enable it.", + nameof(NullTranscodeSessionStore), + TranscodeStoreOptions.RedisConnectionStringKey); + + return serviceCollection.AddSingleton(); + } + + logger.LogInformation( + "Transcode session store: {Store} on {Endpoints}.", + nameof(RedisTranscodeSessionStore), + DescribeEndpoints(redisConnectionString)); + + serviceCollection.AddSingleton(sp => + { + try + { + return ConnectionMultiplexer.Connect(redisConnectionString); + } + catch (Exception ex) + { + sp.GetRequiredService>() + .LogError(ex, "Failed to connect to Redis. Check the {Key} configuration.", TranscodeStoreOptions.RedisConnectionStringKey); + throw; + } + }); + serviceCollection.AddSingleton(); + serviceCollection.AddHostedService(); + + return serviceCollection; + } + + /// + /// Renders the endpoints of a connection string for logging. The connection string itself is never + /// logged because it can carry a password. + /// + private static string DescribeEndpoints(string redisConnectionString) + { + try + { + return string.Join( + ',', + ConfigurationOptions.Parse(redisConnectionString).EndPoints.Select(endpoint => endpoint.ToString())); + } + catch (ArgumentException) + { + return "(unparsable connection string)"; + } + } +} diff --git a/Jellyfin.Server/Program.cs b/Jellyfin.Server/Program.cs index 8390c91313..5530cfae93 100644 --- a/Jellyfin.Server/Program.cs +++ b/Jellyfin.Server/Program.cs @@ -391,6 +391,8 @@ namespace Jellyfin.Server .AddInMemoryCollection(inMemoryDefaultConfig) .AddJsonFile(LoggingConfigFileDefault, optional: false, reloadOnChange: true) .AddJsonFile(LoggingConfigFileSystem, optional: true, reloadOnChange: true) + // Added before the prefixed source so an explicit JELLYFIN_ variable still wins. + .AddJellyfinSectionEnvironmentVariables() .AddEnvironmentVariables("JELLYFIN_") .AddInMemoryCollection(commandLineOpts.ConvertToConfig()); } diff --git a/MediaBrowser.Controller/MediaEncoding/TranscodeStoreOptions.cs b/MediaBrowser.Controller/MediaEncoding/TranscodeStoreOptions.cs index 65a7b31348..ea527b80cf 100644 --- a/MediaBrowser.Controller/MediaEncoding/TranscodeStoreOptions.cs +++ b/MediaBrowser.Controller/MediaEncoding/TranscodeStoreOptions.cs @@ -5,6 +5,16 @@ namespace MediaBrowser.Controller.MediaEncoding; /// public sealed class TranscodeStoreOptions { + /// + /// The configuration section these options bind from. + /// + public const string ConfigurationSection = "Jellyfin:TranscodeStore"; + + /// + /// The configuration key holding the Redis connection string. + /// + public const string RedisConnectionStringKey = ConfigurationSection + ":RedisConnectionString"; + /// /// Gets or sets the Redis connection string. /// A null or empty value indicates single-instance mode, where diff --git a/README.md b/README.md index 816aba126b..7ac2845e7d 100644 --- a/README.md +++ b/README.md @@ -72,7 +72,7 @@ dotnet run --project Jellyfin.Server/Jellyfin.Server.csproj -- \ ### HA mode with Redis -Set the `Jellyfin:TranscodeStore:RedisConnectionString` configuration key. You can pass it as an environment variable, a `DOTNET_` prefixed env var, or in a JSON config file. +Set the `Jellyfin:TranscodeStore:RedisConnectionString` configuration key. You can pass it as a `Jellyfin__TranscodeStore__RedisConnectionString` environment variable, as the equivalent `JELLYFIN_` prefixed variable, or in a JSON config file. **Environment variable:** @@ -98,7 +98,14 @@ dotnet run --project Jellyfin.Server/Jellyfin.Server.csproj -- \ } ``` -When `RedisConnectionString` is set, `RedisTranscodeSessionStore` is registered in DI. If the Redis connection fails at startup, the server throws and refuses to start — this is intentional so you don't silently fall back to broken HA behavior. +The selected store is logged at startup, so HA transcoding is never on or off without a signal: + +``` +Transcode session store: RedisTranscodeSessionStore on valkey:6379. +Redis transcode session store is reachable (2ms round trip). HA transcode takeover is active. +``` + +Without a connection string the line reads `Transcode session store: NullTranscodeSessionStore`. A configured but unreachable store is logged at `Error`; the server keeps serving with per-instance sessions rather than refusing to start. --- diff --git a/docs/FORK-DIFF.md b/docs/FORK-DIFF.md index 0483ed27ff..b7656f4093 100644 --- a/docs/FORK-DIFF.md +++ b/docs/FORK-DIFF.md @@ -37,6 +37,9 @@ auth or plugin logic is rewritten. | `MediaBrowser.Controller/MediaEncoding/TranscodeStoreOptions.cs` | `RedisConnectionString`, `LeaseDurationSeconds` and `SessionRetentionSeconds` | | `MediaBrowser.Controller/MediaEncoding/NullTranscodeSessionStore.cs` | No-op store used when no Redis connection is configured | | `Emby.Server.Implementations/MediaEncoding/RedisTranscodeSessionStore.cs` | Redis store; sessions under `jellyfin:transcode:{playSessionId}`, key TTL is the retention window so an orphaned session outlives its lease | +| `Emby.Server.Implementations/MediaEncoding/TranscodeStoreConnectivityProbe.cs` | Startup ping; an unreachable configured store is logged at `Error` instead of failing open silently | +| `Jellyfin.Server/Extensions/TranscodeStoreServiceCollectionExtensions.cs` | Store selection, logged at `Information` so the active store is visible at startup | +| `Jellyfin.Server/Extensions/ConfigurationBuilderExtensions.cs` | Reads bare `Jellyfin__*` environment variables into the `Jellyfin:*` configuration root | Lease takeover and renewal each run as a single Lua script, so concurrent pods cannot both claim an expired lease and a renewal cannot revert a takeover. The expiry is stored as unix diff --git a/tests/Jellyfin.Server.Implementations.Tests/MediaEncoding/TranscodeStoreWiringTests.cs b/tests/Jellyfin.Server.Implementations.Tests/MediaEncoding/TranscodeStoreWiringTests.cs new file mode 100644 index 0000000000..61df768846 --- /dev/null +++ b/tests/Jellyfin.Server.Implementations.Tests/MediaEncoding/TranscodeStoreWiringTests.cs @@ -0,0 +1,103 @@ +using System; +using System.IO; +using System.Threading.Tasks; +using Emby.Server.Implementations.MediaEncoding; +using Jellyfin.Server; +using Jellyfin.Server.Extensions; +using MediaBrowser.Common.Configuration; +using MediaBrowser.Controller.MediaEncoding; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging.Abstractions; +using Moq; +using StackExchange.Redis; +using Testcontainers.Redis; +using Xunit; + +namespace Jellyfin.Server.Implementations.Tests.MediaEncoding; + +/// +/// Drives the whole configuration path a deployment uses: a bare Jellyfin__TranscodeStore__* +/// environment variable, the server's own configuration builder, the store registration, and a +/// session round-trip against a real Valkey server. +/// +[Trait("Category", "RequiresDocker")] +public sealed class TranscodeStoreWiringTests : IAsyncLifetime +{ + private const string RedisConnectionStringVariable = "Jellyfin__TranscodeStore__RedisConnectionString"; + + private readonly RedisContainer _container; + private string _configDirectory = string.Empty; + + /// + /// Initializes a new instance of the class. + /// + public TranscodeStoreWiringTests() + { + _container = new RedisBuilder("valkey/valkey:8-alpine").Build(); + } + + /// + /// Starts Valkey and lays out the configuration directory the server reads at startup. + /// + /// A representing the asynchronous operation. + public async ValueTask InitializeAsync() + { + await _container.StartAsync().ConfigureAwait(false); + _configDirectory = Directory.CreateTempSubdirectory("jellyfin-wiring-test").FullName; + await File.WriteAllTextAsync(Path.Combine(_configDirectory, "logging.default.json"), "{}").ConfigureAwait(false); + } + + /// + /// Removes the environment variable, configuration directory and container. + /// + /// A representing the asynchronous operation. + public async ValueTask DisposeAsync() + { + Environment.SetEnvironmentVariable(RedisConnectionStringVariable, null); + + if (_configDirectory.Length > 0) + { + Directory.Delete(_configDirectory, true); + } + + await _container.DisposeAsync().ConfigureAwait(false); + } + + /// + /// The variable form deployments set selects the Redis store and that store really talks to Valkey. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task ManifestStyleEnvironmentVariable_Should_Reach_Valkey() + { + Environment.SetEnvironmentVariable( + RedisConnectionStringVariable, + _container.GetConnectionString() + ",abortConnect=false"); + + var appPaths = new Mock(); + appPaths.Setup(paths => paths.ConfigurationDirectoryPath).Returns(_configDirectory); + var configuration = Program.CreateAppConfiguration(new StartupOptions(), appPaths.Object); + + var services = new ServiceCollection(); + services.AddLogging(); + services.AddTranscodeSessionStore(configuration, NullLogger.Instance); + + await using var provider = services.BuildServiceProvider(); + + var store = provider.GetRequiredService(); + Assert.IsType(store); + + var playSessionId = Guid.NewGuid().ToString("N"); + await store.SetAsync( + TranscodeSession.CreateForPlaylist(playSessionId, "media-1", "pod-a", "/transcodes/abc.m3u8", TimeSpan.FromSeconds(30)), + TestContext.Current.CancellationToken); + + var stored = await store.TryGetAsync(playSessionId, TestContext.Current.CancellationToken); + + Assert.NotNull(stored); + Assert.Equal("pod-a", stored.OwnerPod); + + var redis = provider.GetRequiredService(); + Assert.True(await redis.GetDatabase().KeyExistsAsync("jellyfin:transcode:" + playSessionId)); + } +} diff --git a/tests/Jellyfin.Server.Tests/HighAvailability/JellyfinSectionConfigurationTests.cs b/tests/Jellyfin.Server.Tests/HighAvailability/JellyfinSectionConfigurationTests.cs new file mode 100644 index 0000000000..4cfc9518ed --- /dev/null +++ b/tests/Jellyfin.Server.Tests/HighAvailability/JellyfinSectionConfigurationTests.cs @@ -0,0 +1,95 @@ +using System; +using System.IO; +using MediaBrowser.Common.Configuration; +using Microsoft.Extensions.Configuration; +using Moq; +using Xunit; + +namespace Jellyfin.Server.Tests.HighAvailability; + +/// +/// The fork's own settings live under the Jellyfin:* configuration root and every manifest, +/// chart and document sets them as bare Jellyfin__Section__Key environment variables. The +/// startup configuration the host reads them from must therefore accept that form; when it does not, +/// a correctly set variable is dropped and the feature it configures stays off without any error. +/// +public sealed class JellyfinSectionConfigurationTests : IDisposable +{ + private const string RedisKey = "Jellyfin:TranscodeStore:RedisConnectionString"; + private const string UnprefixedVariable = "Jellyfin__TranscodeStore__RedisConnectionString"; + private const string PrefixedVariable = "JELLYFIN_Jellyfin__TranscodeStore__RedisConnectionString"; + private const string LeaseUnprefixedVariable = "Jellyfin__TranscodeStore__LeaseDurationSeconds"; + + private readonly string _configDirectory; + + public JellyfinSectionConfigurationTests() + { + _configDirectory = Directory.CreateTempSubdirectory("jellyfin-config-test").FullName; + File.WriteAllText(Path.Combine(_configDirectory, "logging.default.json"), "{}"); + } + + public void Dispose() + { + Environment.SetEnvironmentVariable(UnprefixedVariable, null); + Environment.SetEnvironmentVariable(PrefixedVariable, null); + Environment.SetEnvironmentVariable(LeaseUnprefixedVariable, null); + Directory.Delete(_configDirectory, true); + } + + [Fact] + public void CreateAppConfiguration_Should_Read_UnprefixedVariable() + { + Environment.SetEnvironmentVariable(UnprefixedVariable, "valkey-cheeztv-valkey:6379,abortConnect=false"); + + var config = CreateConfiguration(); + + Assert.Equal("valkey-cheeztv-valkey:6379,abortConnect=false", config[RedisKey]); + } + + [Fact] + public void CreateAppConfiguration_Should_Read_UnprefixedNonStringValue() + { + Environment.SetEnvironmentVariable(LeaseUnprefixedVariable, "45"); + + var config = CreateConfiguration(); + + Assert.Equal("45", config["Jellyfin:TranscodeStore:LeaseDurationSeconds"]); + } + + [Fact] + public void CreateAppConfiguration_Should_Read_PrefixedVariable() + { + Environment.SetEnvironmentVariable(PrefixedVariable, "redis:6379"); + + var config = CreateConfiguration(); + + Assert.Equal("redis:6379", config[RedisKey]); + } + + [Fact] + public void CreateAppConfiguration_Should_Prefer_PrefixedVariable() + { + Environment.SetEnvironmentVariable(UnprefixedVariable, "unprefixed:6379"); + Environment.SetEnvironmentVariable(PrefixedVariable, "prefixed:6379"); + + var config = CreateConfiguration(); + + Assert.Equal("prefixed:6379", config[RedisKey]); + } + + [Fact] + public void CreateAppConfiguration_Should_Leave_Key_Unset_Without_Variables() + { + var config = CreateConfiguration(); + + Assert.Null(config[RedisKey]); + } + + private IConfiguration CreateConfiguration() + { + var appPaths = new Mock(); + appPaths.Setup(paths => paths.ConfigurationDirectoryPath).Returns(_configDirectory); + + return Program.CreateAppConfiguration(new StartupOptions(), appPaths.Object); + } +} diff --git a/tests/Jellyfin.Server.Tests/HighAvailability/RecordingLogger.cs b/tests/Jellyfin.Server.Tests/HighAvailability/RecordingLogger.cs new file mode 100644 index 0000000000..3ebd00b256 --- /dev/null +++ b/tests/Jellyfin.Server.Tests/HighAvailability/RecordingLogger.cs @@ -0,0 +1,42 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using Microsoft.Extensions.Logging; + +namespace Jellyfin.Server.Tests.HighAvailability; + +/// +/// An that keeps every formatted entry so tests can assert on the startup +/// signals operators rely on. +/// +/// The category type. +internal sealed class RecordingLogger : ILogger +{ + private readonly List<(LogLevel Level, string Message, Exception? Exception)> _entries = new(); + + public IReadOnlyList<(LogLevel Level, string Message, Exception? Exception)> Entries => _entries; + + public IDisposable BeginScope(TState state) + where TState : notnull + => NoopScope.Instance; + + public bool IsEnabled(LogLevel logLevel) => true; + + public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func formatter) + { + ArgumentNullException.ThrowIfNull(formatter); + _entries.Add((logLevel, formatter(state, exception), exception)); + } + + public bool HasEntry(LogLevel level, string substring) + => _entries.Any(entry => entry.Level == level && entry.Message.Contains(substring, StringComparison.Ordinal)); + + private sealed class NoopScope : IDisposable + { + public static readonly NoopScope Instance = new(); + + public void Dispose() + { + } + } +} diff --git a/tests/Jellyfin.Server.Tests/HighAvailability/TranscodeStoreConnectivityProbeTests.cs b/tests/Jellyfin.Server.Tests/HighAvailability/TranscodeStoreConnectivityProbeTests.cs new file mode 100644 index 0000000000..3f952a6001 --- /dev/null +++ b/tests/Jellyfin.Server.Tests/HighAvailability/TranscodeStoreConnectivityProbeTests.cs @@ -0,0 +1,75 @@ +using System; +using System.Threading; +using System.Threading.Tasks; +using Emby.Server.Implementations.MediaEncoding; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using Moq; +using StackExchange.Redis; +using Xunit; + +namespace Jellyfin.Server.Tests.HighAvailability; + +/// +/// A Redis store that cannot be reached degrades silently: the client is configured not to abort the +/// connection and every call site swallows failures. The probe is the only startup signal, so both of +/// its outcomes are pinned here. +/// +public sealed class TranscodeStoreConnectivityProbeTests +{ + [Fact] + public async Task StartAsync_Should_Log_Information_When_Reachable() + { + var database = new Mock(); + database.Setup(db => db.PingAsync(It.IsAny())).ReturnsAsync(TimeSpan.FromMilliseconds(3)); + + var (logger, probe) = CreateProbe(database.Object); + + await probe.StartAsync(CancellationToken.None); + + Assert.True(logger.HasEntry(LogLevel.Information, "reachable")); + Assert.DoesNotContain(logger.Entries, entry => entry.Level >= LogLevel.Warning); + } + + [Fact] + public async Task StartAsync_Should_Log_Error_When_Unreachable() + { + var database = new Mock(); + database.Setup(db => db.PingAsync(It.IsAny())) + .ThrowsAsync(new RedisConnectionException(ConnectionFailureType.UnableToConnect, "no route to host")); + + var (logger, probe) = CreateProbe(database.Object); + + await probe.StartAsync(CancellationToken.None); + + Assert.True(logger.HasEntry(LogLevel.Error, "UNREACHABLE")); + } + + [Fact] + public async Task StartAsync_Should_Log_Error_Instead_Of_Aborting_Startup() + { + var services = new ServiceCollection(); + services.AddSingleton(_ => throw new RedisConnectionException(ConnectionFailureType.UnableToConnect, "no route to host")); + using var provider = services.BuildServiceProvider(); + + var logger = new RecordingLogger(); + var probe = new TranscodeStoreConnectivityProbe(provider, logger); + + await probe.StartAsync(CancellationToken.None); + + Assert.True(logger.HasEntry(LogLevel.Error, "UNREACHABLE")); + } + + private static (RecordingLogger Logger, TranscodeStoreConnectivityProbe Probe) CreateProbe(IDatabase database) + { + var multiplexer = new Mock(); + multiplexer.Setup(redis => redis.GetDatabase(It.IsAny(), It.IsAny())).Returns(database); + + var services = new ServiceCollection(); + services.AddSingleton(multiplexer.Object); + var provider = services.BuildServiceProvider(); + + var logger = new RecordingLogger(); + return (logger, new TranscodeStoreConnectivityProbe(provider, logger)); + } +} diff --git a/tests/Jellyfin.Server.Tests/HighAvailability/TranscodeStoreRegistrationTests.cs b/tests/Jellyfin.Server.Tests/HighAvailability/TranscodeStoreRegistrationTests.cs new file mode 100644 index 0000000000..42f1563a4c --- /dev/null +++ b/tests/Jellyfin.Server.Tests/HighAvailability/TranscodeStoreRegistrationTests.cs @@ -0,0 +1,104 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using Emby.Server.Implementations.MediaEncoding; +using Jellyfin.Server.Extensions; +using MediaBrowser.Controller.MediaEncoding; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; +using Xunit; + +namespace Jellyfin.Server.Tests.HighAvailability; + +/// +/// Which transcode session store was selected is invisible at runtime — the Redis client does not +/// abort on an unreachable server and every call site swallows failures — so the selection is +/// asserted here together with the startup log line that reports it. +/// +public sealed class TranscodeStoreRegistrationTests +{ + [Fact] + public void AddTranscodeSessionStore_Should_Select_RedisStore_From_UnprefixedEnvironmentKeyShape() + { + var services = new ServiceCollection(); + var logger = new RecordingLogger(); + + services.AddTranscodeSessionStore(BuildConfiguration("valkey-cheeztv-valkey:6379,abortConnect=false"), logger); + + Assert.Equal(typeof(RedisTranscodeSessionStore), StoreImplementation(services)); + Assert.True(logger.HasEntry(LogLevel.Information, nameof(RedisTranscodeSessionStore))); + } + + [Fact] + public void AddTranscodeSessionStore_Should_Select_NullStore_Without_ConnectionString() + { + var services = new ServiceCollection(); + var logger = new RecordingLogger(); + + services.AddTranscodeSessionStore(BuildConfiguration(null), logger); + + Assert.Equal(typeof(NullTranscodeSessionStore), StoreImplementation(services)); + Assert.True(logger.HasEntry(LogLevel.Information, nameof(NullTranscodeSessionStore))); + } + + [Fact] + public void AddTranscodeSessionStore_Should_Log_Endpoints_Without_Password() + { + var services = new ServiceCollection(); + var logger = new RecordingLogger(); + + services.AddTranscodeSessionStore(BuildConfiguration("valkey:6379,password=hunter2"), logger); + + var message = Assert.Single(logger.Entries, entry => entry.Level == LogLevel.Information).Message; + Assert.Contains("valkey:6379", message, StringComparison.Ordinal); + Assert.DoesNotContain("hunter2", message, StringComparison.Ordinal); + } + + [Fact] + public void AddTranscodeSessionStore_Should_Register_ConnectivityProbe_Only_With_Redis() + { + var withRedis = new ServiceCollection(); + withRedis.AddTranscodeSessionStore(BuildConfiguration("valkey:6379"), new RecordingLogger()); + + var withoutRedis = new ServiceCollection(); + withoutRedis.AddTranscodeSessionStore(BuildConfiguration(null), new RecordingLogger()); + + Assert.Contains(withRedis, descriptor => descriptor.ImplementationType == typeof(TranscodeStoreConnectivityProbe)); + Assert.DoesNotContain(withoutRedis, descriptor => descriptor.ServiceType == typeof(IHostedService)); + } + + [Fact] + public void AddTranscodeSessionStore_Should_Bind_Options() + { + var services = new ServiceCollection(); + var configuration = new ConfigurationBuilder() + .AddInMemoryCollection(new Dictionary + { + [TranscodeStoreOptions.RedisConnectionStringKey] = "valkey:6379", + ["Jellyfin:TranscodeStore:LeaseDurationSeconds"] = "45" + }) + .Build(); + + services.AddTranscodeSessionStore(configuration, new RecordingLogger()); + + using var provider = services.BuildServiceProvider(); + var options = provider.GetRequiredService>().Value; + + Assert.Equal(45, options.LeaseDurationSeconds); + Assert.Equal("valkey:6379", options.RedisConnectionString); + } + + private static IConfiguration BuildConfiguration(string? redisConnectionString) + => new ConfigurationBuilder() + .AddInMemoryCollection(new Dictionary + { + [TranscodeStoreOptions.RedisConnectionStringKey] = redisConnectionString + }) + .Build(); + + private static Type? StoreImplementation(IServiceCollection services) + => services.Single(descriptor => descriptor.ServiceType == typeof(ITranscodeSessionStore)).ImplementationType; +} From 6e405af1dc41e26301c465be4e4c43ebcab8d39a Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Sun, 13 Sep 2026 13:18:11 +1000 Subject: [PATCH 3/3] fix(db): collapse presentation-key groups without min(uuid) PostgreSQL has no min(uuid) aggregate, so every query that picked a group representative with MIN over the item id failed with 42883: library browse, search, Recently Added, the by-name endpoints and Upcoming. - Pick the representative with an anti-join on (primary version, id) - Cover the collapse with a repository test against a real PostgreSQL --- .../Item/BaseItemRepository.ByName.cs | 6 +- .../Item/BaseItemRepository.QueryBuilding.cs | 24 +- .../PostgreSqlPresentationKeyGroupingTests.cs | 213 ++++++++++++++++++ 3 files changed, 236 insertions(+), 7 deletions(-) create mode 100644 tests/Jellyfin.Server.Tests/Item/PostgreSqlPresentationKeyGroupingTests.cs diff --git a/Jellyfin.Server.Implementations/Item/BaseItemRepository.ByName.cs b/Jellyfin.Server.Implementations/Item/BaseItemRepository.ByName.cs index cdc8744642..2a4b9c516b 100644 --- a/Jellyfin.Server.Implementations/Item/BaseItemRepository.ByName.cs +++ b/Jellyfin.Server.Implementations/Item/BaseItemRepository.ByName.cs @@ -219,9 +219,11 @@ public sealed partial class BaseItemRepository } else { + // The representative is the row no other row in its group sorts before, not MIN(Id): + // PostgreSQL has no min(uuid) aggregate, while comparing two uuids is supported everywhere. representativeIds = masterQuery - .GroupBy(e => e.PresentationUniqueKey) - .Select(g => g.Min(e => e.Id)) + .Where(e => !masterQuery.Any(o => o.PresentationUniqueKey == e.PresentationUniqueKey && o.Id.CompareTo(e.Id) < 0)) + .Select(e => e.Id) .ToList(); } diff --git a/Jellyfin.Server.Implementations/Item/BaseItemRepository.QueryBuilding.cs b/Jellyfin.Server.Implementations/Item/BaseItemRepository.QueryBuilding.cs index c0067d8392..99f3577725 100644 --- a/Jellyfin.Server.Implementations/Item/BaseItemRepository.QueryBuilding.cs +++ b/Jellyfin.Server.Implementations/Item/BaseItemRepository.QueryBuilding.cs @@ -96,22 +96,36 @@ public sealed partial class BaseItemRepository // primary version (PrimaryVersionId is null) so detail pages and actions target it instead // of an arbitrary alternate. Keep the grouped ids as an IQueryable sub-select; materializing // to a List would inline one bound parameter per id and hit SQLite's variable cap. + // The representative is the row no other row in its group sorts before, not MIN(Id): + // PostgreSQL has no min(uuid) aggregate, while comparing two uuids is supported everywhere. + // The anti-join reads the filtered set twice, so it has to close over a local that the + // reassignment below cannot reach - capturing dbQuery itself makes the tree self-referential. + var candidates = dbQuery; var enableGroupByPresentationUniqueKey = EnableGroupByPresentationUniqueKey(filter); if (enableGroupByPresentationUniqueKey && filter.GroupBySeriesPresentationUniqueKey) { - var groupedIds = dbQuery.GroupBy(e => new { e.PresentationUniqueKey, e.SeriesPresentationUniqueKey }) - .Select(g => g.Where(e => e.PrimaryVersionId == null).Min(e => (Guid?)e.Id) ?? g.Min(e => (Guid?)e.Id)); + var groupedIds = candidates + .Where(e => !candidates.Any(o => o.PresentationUniqueKey == e.PresentationUniqueKey + && o.SeriesPresentationUniqueKey == e.SeriesPresentationUniqueKey + && ((o.PrimaryVersionId == null && e.PrimaryVersionId != null) + || ((o.PrimaryVersionId == null) == (e.PrimaryVersionId == null) && o.Id.CompareTo(e.Id) < 0)))) + .Select(e => e.Id); dbQuery = context.BaseItems.AsNoTracking().Where(e => groupedIds.Contains(e.Id)); } else if (enableGroupByPresentationUniqueKey) { - var groupedIds = dbQuery.GroupBy(e => e.PresentationUniqueKey) - .Select(g => g.Where(e => e.PrimaryVersionId == null).Min(e => (Guid?)e.Id) ?? g.Min(e => (Guid?)e.Id)); + var groupedIds = candidates + .Where(e => !candidates.Any(o => o.PresentationUniqueKey == e.PresentationUniqueKey + && ((o.PrimaryVersionId == null && e.PrimaryVersionId != null) + || ((o.PrimaryVersionId == null) == (e.PrimaryVersionId == null) && o.Id.CompareTo(e.Id) < 0)))) + .Select(e => e.Id); dbQuery = context.BaseItems.AsNoTracking().Where(e => groupedIds.Contains(e.Id)); } else if (filter.GroupBySeriesPresentationUniqueKey) { - var groupedIds = dbQuery.GroupBy(e => e.SeriesPresentationUniqueKey).Select(e => e.Min(x => x.Id)); + var groupedIds = candidates + .Where(e => !candidates.Any(o => o.SeriesPresentationUniqueKey == e.SeriesPresentationUniqueKey && o.Id.CompareTo(e.Id) < 0)) + .Select(e => e.Id); dbQuery = context.BaseItems.AsNoTracking().Where(e => groupedIds.Contains(e.Id)); } else diff --git a/tests/Jellyfin.Server.Tests/Item/PostgreSqlPresentationKeyGroupingTests.cs b/tests/Jellyfin.Server.Tests/Item/PostgreSqlPresentationKeyGroupingTests.cs new file mode 100644 index 0000000000..343542266c --- /dev/null +++ b/tests/Jellyfin.Server.Tests/Item/PostgreSqlPresentationKeyGroupingTests.cs @@ -0,0 +1,213 @@ +using System; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Emby.Server.Implementations.Data; +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.Implementations.Item; +using Jellyfin.Server.Tests.Migrations; +using MediaBrowser.Controller; +using MediaBrowser.Controller.Configuration; +using MediaBrowser.Controller.Entities; +using MediaBrowser.Model.Configuration; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Logging.Abstractions; +using Moq; +using Npgsql; +using Xunit; +using BaseItemKind = Jellyfin.Data.Enums.BaseItemKind; +using User = Jellyfin.Database.Implementations.Entities.User; + +namespace Jellyfin.Server.Tests.Item; + +/// +/// Runs the presentation-key collapse that library browse, search and the by-name endpoints all go +/// through against a real PostgreSQL. SQLite accepts MIN over any column type, PostgreSQL has no +/// min(uuid) aggregate, so a representative picked with an aggregate over the id only ever fails +/// here - with 42883 function min(uuid) does not exist. +/// +[Trait("Category", "RequiresDocker")] +public sealed class PostgreSqlPresentationKeyGroupingTests : IAsyncLifetime +{ + // The alternate sorts before the primary, and the second duplicate genre before the first, so a + // representative that ignores the primary-version preference or the id order picks the wrong row. + private static readonly Guid _primaryMovieId = Guid.Parse("eeeeeeee-0000-0000-0000-000000000001"); + private static readonly Guid _alternateMovieId = Guid.Parse("11111111-0000-0000-0000-000000000002"); + private static readonly Guid _secondAlternateMovieId = Guid.Parse("22222222-0000-0000-0000-000000000003"); + private static readonly Guid _firstOrphanId = Guid.Parse("33333333-0000-0000-0000-000000000004"); + private static readonly Guid _secondOrphanId = Guid.Parse("dddddddd-0000-0000-0000-000000000005"); + private static readonly Guid _orphanPrimaryId = Guid.Parse("cccccccc-0000-0000-0000-000000000006"); + private static readonly Guid _standaloneMovieId = Guid.Parse("44444444-0000-0000-0000-000000000007"); + + private static readonly Guid _firstGenreId = Guid.Parse("bbbbbbbb-0000-0000-0000-000000000001"); + private static readonly Guid _secondGenreId = Guid.Parse("aaaaaaaa-0000-0000-0000-000000000002"); + private static readonly Guid _genreMovieId = Guid.Parse("66666666-0000-0000-0000-000000000003"); + private static readonly Guid _genreValueId = Guid.Parse("77777777-0000-0000-0000-000000000004"); + + private readonly ItemTypeLookup _itemTypeLookup = new(); + + private PostgreSqlTestServer _server = null!; + private NpgsqlDataSource _dataSource = null!; + private BaseItemRepository _repository = null!; + + public async ValueTask InitializeAsync() + { + _server = await PostgreSqlTestServer.StartAsync().ConfigureAwait(false); + var connectionString = await _server.CreateDatabaseAsync("presentation_key_grouping", TestContext.Current.CancellationToken).ConfigureAwait(false); + _dataSource = new NpgsqlDataSourceBuilder(connectionString).Build(); + + var context = CreateDbContext(); + await using (context.ConfigureAwait(false)) + { + await context.Database.EnsureCreatedAsync(TestContext.Current.CancellationToken).ConfigureAwait(false); + } + + var serverConfigurationManager = new Mock(); + serverConfigurationManager.Setup(c => c.Configuration).Returns(new ServerConfiguration()); + + var factory = new Mock>(); + factory.Setup(f => f.CreateDbContext()).Returns(CreateDbContext); + factory.Setup(f => f.CreateDbContextAsync(It.IsAny())).ReturnsAsync(CreateDbContext); + + _repository = new BaseItemRepository( + factory.Object, + new Mock().Object, + _itemTypeLookup, + serverConfigurationManager.Object, + NullLogger.Instance); + + await SeedAsync().ConfigureAwait(false); + } + + public async ValueTask DisposeAsync() + { + await _dataSource.DisposeAsync().ConfigureAwait(false); + await _server.DisposeAsync().ConfigureAwait(false); + } + + /// + /// The collapse behind library browse and search: one row per presentation key, the primary version + /// when the group has one, the lowest id otherwise. + /// + [Fact] + public void GetItemList_CollapsesPresentationKeyGroups() + { + var items = _repository.GetItemList(new InternalItemsQuery(new User("grouping", "auth", "reset")) + { + IncludeItemTypes = [BaseItemKind.Movie], + IncludeOwnedItems = true + }); + + Assert.Equal( + new[] { _primaryMovieId, _firstOrphanId, _standaloneMovieId, _genreMovieId }.OrderBy(id => id).ToArray(), + items.Select(i => i.Id).OrderBy(id => id).ToArray()); + } + + /// + /// The same collapse on the by-name path, which picks the lowest id per group without a + /// primary-version preference. + /// + [Fact] + public void GetGenres_CollapsesPresentationKeyGroups() + { + var result = _repository.GetGenres(new InternalItemsQuery(new User("genres", "auth", "reset"))); + + var item = Assert.Single(result.Items); + Assert.Equal(_secondGenreId, item.Item.Id); + Assert.Equal(1, result.TotalRecordCount); + } + + private JellyfinDbContext CreateDbContext() + { + 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)); + } + + private async Task SeedAsync() + { + var context = CreateDbContext(); + await using (context.ConfigureAwait(false)) + { + // One version group: a primary plus two alternates, both sorting before it by id. + var versionKey = _primaryMovieId.ToString("N"); + context.BaseItems.Add(CreateMovie(_primaryMovieId, "Movie", versionKey, null)); + context.BaseItems.Add(CreateMovie(_alternateMovieId, "Movie - 1080p", versionKey, _primaryMovieId)); + context.BaseItems.Add(CreateMovie(_secondAlternateMovieId, "Movie - 4K", versionKey, _primaryMovieId)); + + // One group whose primary is not in the result set, so the lowest id represents it. + var orphanKey = _orphanPrimaryId.ToString("N"); + context.BaseItems.Add(CreateMovie(_firstOrphanId, "Orphan - 1080p", orphanKey, _orphanPrimaryId)); + context.BaseItems.Add(CreateMovie(_secondOrphanId, "Orphan - 4K", orphanKey, _orphanPrimaryId)); + + context.BaseItems.Add(CreateMovie(_standaloneMovieId, "Standalone", _standaloneMovieId.ToString("N"), null)); + + // Two genre entities sharing a presentation key, credited on one movie so the by-name + // item-value join sees them. + var genreKey = _firstGenreId.ToString("N"); + context.BaseItems.Add(CreateGenre(_firstGenreId, "Action", genreKey)); + context.BaseItems.Add(CreateGenre(_secondGenreId, "Action", genreKey)); + + var genreMovie = CreateMovie(_genreMovieId, "Genre Movie", _genreMovieId.ToString("N"), null); + genreMovie.CleanName = "genre movie"; + context.BaseItems.Add(genreMovie); + + var genreValue = new ItemValue + { + ItemValueId = _genreValueId, + Type = ItemValueType.Genre, + Value = "Action", + CleanValue = "action" + }; + context.ItemValues.Add(genreValue); + context.ItemValuesMap.Add(new ItemValueMap + { + ItemId = _genreMovieId, + ItemValueId = _genreValueId, + Item = genreMovie, + ItemValue = genreValue + }); + + await context.SaveChangesAsync(TestContext.Current.CancellationToken).ConfigureAwait(false); + } + } + + private BaseItemEntity CreateMovie(Guid id, string name, string presentationKey, Guid? primaryVersionId) + { + return new BaseItemEntity + { + Id = id, + Type = _itemTypeLookup.BaseItemKindNames[BaseItemKind.Movie], + Name = name, + PresentationUniqueKey = presentationKey, + PrimaryVersionId = primaryVersionId, + MediaType = "Video", + IsMovie = true, + IsFolder = false, + IsVirtualItem = false + }; + } + + private BaseItemEntity CreateGenre(Guid id, string name, string presentationKey) + { + return new BaseItemEntity + { + Id = id, + Type = _itemTypeLookup.BaseItemKindNames[BaseItemKind.Genre], + Name = name, + CleanName = "action", + PresentationUniqueKey = presentationKey, + IsFolder = false, + IsVirtualItem = false + }; + } +}