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/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/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/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/Jellyfin.Server/CoreAppHost.cs b/Jellyfin.Server/CoreAppHost.cs
index 8e037bf742..ee6c80ba10 100644
--- a/Jellyfin.Server/CoreAppHost.cs
+++ b/Jellyfin.Server/CoreAppHost.cs
@@ -3,6 +3,7 @@ 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;
using Jellyfin.Database.Implementations;
@@ -34,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
{
@@ -106,29 +106,7 @@ 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();
- }
+ serviceCollection.AddTranscodeSessionStore(_startupConfig, Logger);
// Scan-leader lease: gates periodic library-mutating scheduled tasks to a single leader
// instance. Active by default once 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/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/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 0528338071..dfd4f7fb7f 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/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);
+ }
+}
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