Compare commits

..

22 Commits

Author SHA1 Message Date
unkin-agent 33a5fbce9b retry building the quick connect store in the startup probe, not just reading it
ci/woodpecker/pr/ci Pipeline was successful
ci/woodpecker/push/ci Pipeline was successful
2026-09-26 22:56:28 +10:00
unkin-agent 6c7f76fd26 retry the quick connect startup probe, and exit non-zero when a start fails
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
2026-09-26 22:14:14 +10:00
unkin-agent 624d528d28 install libfontconfig1 for the docker test step
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
The startup tests build a real app host, which probes the Skia encoder, and
loading libSkiaSharp needs fontconfig.
2026-09-26 19:30:43 +10:00
unkin-agent a99ca458bf fail startup when the quick connect store is unreachable
ci/woodpecker/push/ci Pipeline failed
ci/woodpecker/pr/ci Pipeline failed
Read the configured quick connect store during InitializeServices and abort
startup, with a log line naming valkey, instead of coming up and failing on
the first request that needs the store.
2026-09-26 19:09:24 +10:00
benvin 2a09108d4f Merge pull request 'fix(ha): share quick connect state between instances' (#34) from benvin/quickconnect-shared-state into main
ci/woodpecker/push/ci Pipeline was successful
Reviewed-on: #34
2026-09-26 16:39:55 +10:00
unkin-agent e1ca272d6c fail quick connect closed on an unreachable valkey and restore the idempotent exchange
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
2026-09-26 14:28:45 +10:00
unkin-agent 6c54a02240 make quick connect authorize atomic and stop the fallback shadowing redis
ci/woodpecker/pr/ci Pipeline was successful
ci/woodpecker/push/ci Pipeline was successful
2026-09-24 23:38:35 +10:00
unkin-agent ba4d487c65 share quick connect state between instances
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
Move the pending requests and authorized secrets out of process so the
initiate, authorize and exchange legs can land on different replicas.
2026-09-24 23:02:56 +10:00
benvin 6691b785c3 Merge pull request #32 from benvin/userdata-cache
ci/woodpecker/push/ci Pipeline was successful
fix(userdata): read user data through to the database
2026-09-21 23:13:37 +10:00
benvin 2adb13f50f Merge pull request #29 from benvin/config-propagation
ci/woodpecker/push/ci Pipeline was successful
fix(ha): propagate shared-config and library-visibility changes between replicas
2026-09-21 21:44:34 +10:00
benvin 1c98f4a074 Merge pull request #31 from benvin/pg-tests-in-ci
ci/woodpecker/push/ci Pipeline was successful
test(db): run the PostgreSQL provider tests in CI
2026-09-21 07:45:29 +10:00
benvin 393994a454 Merge pull request #30 from benvin/nextup-datetime-kind
ci/woodpecker/push/ci Pipeline was canceled
fix(api): normalise the Next Up cutoff to UTC
2026-09-21 07:42:33 +10:00
benvin ad50c4e433 Merge pull request #28 from benvin/gated-task-keys
ci/woodpecker/push/ci Pipeline was canceled
fix(ha): gate the remaining timer-driven scheduled tasks
2026-09-21 07:40:32 +10:00
unkin-agent 483c739fb1 report a peer's port change locally without rewriting it
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
Keep IConfigurationManager source-compatible for plugin implementers.
2026-09-21 00:46:16 +10:00
unkin-agent 1c59e6afcb perf(nextup): batch the user data reads the next episode selection makes
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
Next Up walked every series in an unpaginated loop and read user data one
episode at a time, so the Home screen row cost two queries per series once the
cache was gone.

- read the played state of every candidate episode in one query
- read the versions the resume check and the last played date need in one query each
- default PrefetchedUserData on IUserBaseItemComparer so plugin comparers still compile
- cover the bounded query count and the plugin comparer with tests
- share one database across the replica tests
2026-09-21 00:36:39 +10:00
unkin-agent d39ec60e2c test(db): link PostgreSqlTestServer instead of referencing Jellyfin.Server.Tests
ci/woodpecker/pr/ci Pipeline was successful
ci/woodpecker/push/ci Pipeline was successful
2026-09-21 00:29:11 +10:00
unkin-agent a9d6c749fb keep an applied invalidation from inducing a write that publishes back
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
2026-09-21 00:24:01 +10:00
unkin-agent 44b62dcc64 test(ha): pin the default gated task key set
ci/woodpecker/pr/ci Pipeline was successful
ci/woodpecker/push/ci Pipeline was successful
2026-09-20 23:55:17 +10:00
unkin-agent a7919b9bac fix(userdata): read user data through to the database
ci/woodpecker/pr/ci Pipeline was successful
ci/woodpecker/push/ci Pipeline was successful
The user data cache and the item's in-memory rows were both filled once per
replica and never invalidated, so a pod serving a playback tick read its own
stale resume position and wrote it back over the position another pod had just
saved, silently losing resume points, played state, favourites and ratings.

- drop the user item data cache
- read single and batched user data from the database on every query
- prefetch user data for in-memory sorts and filters so each stays one query
- cover the lost update against real PostgreSQL with two manager instances
2026-09-20 23:51:29 +10:00
unkin-agent a025655b4d test(db): run the PostgreSQL provider tests in CI
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
Point Jellyfin.Database.Tests.PostgreSQL at the server postgres-migration-chain
already runs, give every test its own database, and fix the two tests that only
ever failed against real PostgreSQL.
2026-09-20 23:48:43 +10:00
unkin-agent b662ffa48f propagate shared-config and library-option changes between instances
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
Add a Redis pub/sub invalidation bus, publish from the configuration and library-option write paths, and drop the matching local cache entry on the instances that did not write. No-op without a Redis connection string, and fails open when it is unreachable.
2026-09-20 23:47:21 +10:00
unkin-agent 1965c68a76 fix(ha): gate the remaining timer-driven scheduled tasks
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
Add the eight provider, live TV and plugin-update task keys to
ScanLeaderOptions.GatedTaskKeys and widen the test's task-key discovery to
every assembly that declares an IScheduledTask.
2026-09-20 23:36:34 +10:00
74 changed files with 4976 additions and 395 deletions
+11 -4
View File
@@ -45,9 +45,9 @@ steps:
# testcontainers, and a postgres service container deadlocks the step because the backend mounts
# the ReadWriteOnce workspace volume into service pods and schedules them on another node.
# Its data directory lives on the step's ephemeral storage, not on the workspace volume.
# Scoped to Jellyfin.Server.Tests: the three classes in Jellyfin.Database.Tests.PostgreSQL still
# start their own container, and PostgreSqlProviderTests already fails on main - on an EF 10
# scalar query and on its own data - which a third test in the class then inherits.
# Both projects attach to it through JELLYFIN_TEST_POSTGRES and give every test a database of
# its own, so nothing here depends on a docker daemon.
# Valkey runs in the step for the same reason, attached through JELLYFIN_TEST_REDIS.
- name: postgres-migration-chain
image: mcr.microsoft.com/dotnet/sdk:10.0
depends_on:
@@ -56,15 +56,22 @@ steps:
DOTNET_CLI_TELEMETRY_OPTOUT: "1"
DOTNET_NOLOGO: "1"
JELLYFIN_TEST_POSTGRES: "Host=127.0.0.1;Port=5432;Database=postgres;Username=postgres"
JELLYFIN_TEST_REDIS: "127.0.0.1:6379"
commands:
- apt-get -o Acquire::Retries=3 update || apt-get -o Acquire::Retries=3 update
- DEBIAN_FRONTEND=noninteractive apt-get install -y --no-install-recommends postgresql
# libfontconfig1 is needed here too: the startup tests build a real app host, which probes
# the Skia encoder, and loading libSkiaSharp pulls fontconfig in.
- DEBIAN_FRONTEND=noninteractive apt-get install -y --no-install-recommends postgresql valkey-server libfontconfig1
- install -d -o postgres -g postgres /tmp/pgdata /tmp/pgrun
- PGBIN=$(ls -d /usr/lib/postgresql/*/bin | tail -1)
- su postgres -c "$PGBIN/initdb -D /tmp/pgdata -A trust -U postgres"
- su postgres -c "$PGBIN/pg_ctl -D /tmp/pgdata -o \"-c listen_addresses=127.0.0.1 -k /tmp/pgrun\" -l /tmp/pg.log -w start"
- valkey-server --daemonize yes --bind 127.0.0.1 --port 6379 --save ''
- valkey-cli ping
- dotnet build tests/Jellyfin.Server.Tests/Jellyfin.Server.Tests.csproj -c Release
- dotnet build tests/Jellyfin.Database.Tests.PostgreSQL/Jellyfin.Database.Tests.PostgreSQL.csproj -c Release
- dotnet test tests/Jellyfin.Server.Tests/Jellyfin.Server.Tests.csproj -c Release --no-build --verbosity minimal --filter "Category=RequiresDocker"
- dotnet test tests/Jellyfin.Database.Tests.PostgreSQL/Jellyfin.Database.Tests.PostgreSQL.csproj -c Release --no-build --verbosity minimal --filter "Category=RequiresDocker"
backend_options:
kubernetes:
serviceAccountName: jellyfin-ha-src
@@ -87,6 +87,12 @@ namespace Emby.Server.Implementations.AppBase
/// <value>The application paths.</value>
public IApplicationPaths CommonApplicationPaths { get; private set; }
/// <summary>
/// Gets or sets the bus announcing configuration writes to the other instances sharing this
/// configuration directory. Defaults to a no-op, which is the single-instance behaviour.
/// </summary>
public IConfigurationInvalidationBus InvalidationBus { get; set; } = NullConfigurationInvalidationBus.Instance;
/// <summary>
/// Gets or sets the system configuration.
/// </summary>
@@ -169,6 +175,8 @@ namespace Emby.Server.Implementations.AppBase
}
OnConfigurationUpdated();
InvalidationBus.PublishLocalWrite(ConfigurationInvalidationScope.SystemConfiguration, null);
}
/// <summary>
@@ -350,6 +358,29 @@ namespace Emby.Server.Implementations.AppBase
}
OnNamedConfigurationUpdated(key, configuration);
InvalidationBus.PublishLocalWrite(ConfigurationInvalidationScope.NamedConfiguration, key);
}
/// <inheritdoc />
public void InvalidateCachedConfiguration(string? key)
{
if (string.IsNullOrEmpty(key))
{
lock (_configurationSyncLock)
{
_configuration = null;
}
// Reloads the system configuration off the shared file as a side effect of re-deriving
// the cache path from it, then tells the in-process consumers to re-read it.
OnConfigurationUpdated();
return;
}
_configurations.TryRemove(key, out _);
OnNamedConfigurationUpdated(key, GetConfiguration(key));
}
/// <summary>
+130 -15
View File
@@ -52,6 +52,7 @@ using Jellyfin.Server.Implementations.SystemBackupService;
using MediaBrowser.Common;
using MediaBrowser.Common.Configuration;
using MediaBrowser.Common.Events;
using MediaBrowser.Common.Extensions;
using MediaBrowser.Common.Net;
using MediaBrowser.Common.Plugins;
using MediaBrowser.Common.Updates;
@@ -124,6 +125,21 @@ namespace Emby.Server.Implementations
/// </summary>
public abstract class ApplicationHost : IServerApplicationHost, IDisposable
{
/// <summary>
/// The secret the startup read of the quick connect store looks for. No flow ever mints it, so the
/// read is always a miss and only its reachability is being asked about.
/// </summary>
private const string StartupProbeSecret = "startup-probe";
/// <summary>
/// How long the startup read of the quick connect store is retried before the store counts as
/// unreachable. Long enough to ride out valkey restarting alongside this instance, short enough
/// that a store which is really gone is reported inside one liveness cycle.
/// </summary>
private static readonly TimeSpan _quickConnectProbeDeadline = TimeSpan.FromSeconds(30);
private static readonly TimeSpan _quickConnectProbeRetryDelay = TimeSpan.FromSeconds(1);
/// <summary>
/// The disposable parts.
/// </summary>
@@ -644,6 +660,8 @@ namespace Emby.Server.Implementations
/// <returns>A task representing the service initialization operation.</returns>
public async Task InitializeServices(IConfiguration startupConfig)
{
await ProbeQuickConnectStoreAsync().ConfigureAwait(false);
var localizationManager = (LocalizationManager)Resolve<ILocalizationManager>();
await localizationManager.LoadAll().ConfigureAwait(false);
@@ -652,6 +670,66 @@ namespace Emby.Server.Implementations
FindParts();
}
/// <summary>
/// Builds and reads the quick connect store here so a store that cannot be reached stops startup,
/// rather than being discovered on the first request that needs it. Both halves sit inside the
/// retry: a store built with <c>abortConnect=false</c> constructs without touching the network and
/// only the read settles it, while one built without that option connects eagerly and fails at the
/// resolve. Retried until <see cref="_quickConnectProbeDeadline"/> so a starting instance rides out
/// the blip a running one already tolerates.
/// </summary>
private async Task ProbeQuickConnectStoreAsync()
{
var startTimestamp = Stopwatch.GetTimestamp();
var reported = false;
while (true)
{
try
{
// A throwing singleton factory is not cached, so the resolve is retried along with the read.
var store = Resolve<IQuickConnectStore>();
await store.GetRequestBySecretAsync(StartupProbeSecret).ConfigureAwait(false);
return;
}
catch (Exception ex)
{
var elapsed = Stopwatch.GetElapsedTime(startTimestamp);
if (elapsed + _quickConnectProbeRetryDelay < _quickConnectProbeDeadline)
{
if (!reported)
{
reported = true;
Logger.LogWarning(
ex,
"Quick connect store is not reachable yet, retrying for up to {Seconds}s.",
(int)_quickConnectProbeDeadline.TotalSeconds);
}
else
{
Logger.LogDebug(ex, "Quick connect store is still not reachable, retrying.");
}
await Task.Delay(_quickConnectProbeRetryDelay).ConfigureAwait(false);
continue;
}
Logger.LogCritical(
ex,
"Quick connect is configured against the shared valkey/Redis store at {Key} and it is UNREACHABLE after {Seconds}s, so the server will not start. Bring valkey up, or clear that setting to keep quick connect state on this instance alone.",
TranscodeStoreOptions.RedisConnectionStringKey,
(int)elapsed.TotalSeconds);
if (ex is ServiceUnavailableException)
{
throw;
}
throw new ServiceUnavailableException("Quick connect store is unreachable.", ex);
}
}
}
private X509Certificate2 GetCertificate(string path, string password)
{
if (string.IsNullOrWhiteSpace(path))
@@ -706,6 +784,8 @@ namespace Emby.Server.Implementations
BaseItem.UserDataManager = Resolve<IUserDataManager>();
CollectionFolder.XmlSerializer = _xmlSerializer;
CollectionFolder.ApplicationHost = this;
CollectionFolder.InvalidationBus = Resolve<IConfigurationInvalidationBus>();
ConfigurationManager.InvalidationBus = CollectionFolder.InvalidationBus;
Folder.UserViewManager = Resolve<IUserViewManager>();
Folder.CollectionManager = Resolve<ICollectionManager>();
Folder.LimitedConcurrencyLibraryScheduler = Resolve<ILimitedConcurrencyLibraryScheduler>();
@@ -790,6 +870,43 @@ namespace Emby.Server.Implementations
}
}
/// <summary>
/// Works out what a configuration update means for the ports this process bound at startup.
/// </summary>
/// <param name="boundHttpPort">The HTTP port this process is bound to.</param>
/// <param name="boundHttpsPort">The HTTPS port this process is bound to.</param>
/// <param name="configuredHttpPort">The HTTP port the shared configuration now carries.</param>
/// <param name="configuredHttpsPort">The HTTPS port the shared configuration now carries.</param>
/// <param name="isPortAuthorized">Whether the shared configuration still marks the port as authorized.</param>
/// <param name="isApplyingRemoteInvalidation">Whether this update is another instance's write being applied.</param>
/// <returns>What the update requires of this instance.</returns>
internal static PortChangeOutcome EvaluatePortChange(
int boundHttpPort,
int boundHttpsPort,
int configuredHttpPort,
int configuredHttpsPort,
bool isPortAuthorized,
bool isApplyingRemoteInvalidation)
{
// Nothing is bound yet, so nothing has gone stale.
if (boundHttpPort == 0 || boundHttpsPort == 0)
{
return default;
}
if (configuredHttpPort == boundHttpPort && configuredHttpsPort == boundHttpsPort)
{
return default;
}
// Whoever wrote the change, this process is still listening on a port the configuration no
// longer names, so the pending restart is reported either way. The authorization flag belongs
// to the instance that made the change: it cleared the flag along with the port, and clearing
// it again here would write shared configuration on that instance's behalf and announce it a
// second time.
return new PortChangeOutcome(true, isPortAuthorized && !isApplyingRemoteInvalidation);
}
/// <summary>
/// Called when [configuration updated].
/// </summary>
@@ -797,26 +914,24 @@ namespace Emby.Server.Implementations
/// <param name="e">The <see cref="EventArgs"/> instance containing the event data.</param>
private void OnConfigurationUpdated(object sender, EventArgs e)
{
var requiresRestart = false;
var networkConfiguration = ConfigurationManager.GetNetworkConfiguration();
// Don't do anything if these haven't been set yet
if (HttpPort != 0 && HttpsPort != 0)
{
// Need to restart if ports have changed
if (networkConfiguration.InternalHttpPort != HttpPort
|| networkConfiguration.InternalHttpsPort != HttpsPort)
{
if (ConfigurationManager.Configuration.IsPortAuthorized)
{
ConfigurationManager.Configuration.IsPortAuthorized = false;
ConfigurationManager.SaveConfiguration();
var portChange = EvaluatePortChange(
HttpPort,
HttpsPort,
networkConfiguration.InternalHttpPort,
networkConfiguration.InternalHttpsPort,
ConfigurationManager.Configuration.IsPortAuthorized,
ConfigurationInvalidationContext.IsApplyingRemoteInvalidation);
requiresRestart = true;
}
}
if (portChange.ClearsPortAuthorization)
{
ConfigurationManager.Configuration.IsPortAuthorized = false;
ConfigurationManager.SaveConfiguration();
}
var requiresRestart = portChange.RequiresRestart;
if (ValidateSslCertificate(networkConfiguration))
{
requiresRestart = true;
@@ -0,0 +1,82 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using MediaBrowser.Common.Configuration;
using MediaBrowser.Controller.Entities;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
namespace Emby.Server.Implementations.Configuration
{
/// <summary>
/// Applies the configuration invalidations published by the other instances sharing this
/// configuration directory, dropping the local cache entry so the next read comes off the shared file.
/// </summary>
public sealed class ConfigurationInvalidationSubscriber : IHostedService
{
private readonly IConfigurationInvalidationBus _bus;
private readonly IConfigurationManager _configurationManager;
private readonly ILogger<ConfigurationInvalidationSubscriber> _logger;
/// <summary>
/// Initializes a new instance of the <see cref="ConfigurationInvalidationSubscriber"/> class.
/// </summary>
/// <param name="bus">The invalidation bus.</param>
/// <param name="configurationManager">The configuration manager holding the cached configuration.</param>
/// <param name="logger">The logger.</param>
public ConfigurationInvalidationSubscriber(
IConfigurationInvalidationBus bus,
IConfigurationManager configurationManager,
ILogger<ConfigurationInvalidationSubscriber> logger)
{
_bus = bus;
_configurationManager = configurationManager;
_logger = logger;
}
/// <inheritdoc />
public Task StartAsync(CancellationToken cancellationToken)
{
_bus.Subscribe(Apply);
return Task.CompletedTask;
}
/// <inheritdoc />
public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask;
private void Apply(ConfigurationInvalidation invalidation)
{
try
{
// Applying re-raises the same update events a local save raises, so the in-process
// consumers re-read. The scope tells those consumers that the write was somebody else's,
// so the ones that answer an update by writing neither repeat it nor publish it back.
using var scope = ConfigurationInvalidationContext.BeginApply();
switch (invalidation.Scope)
{
case ConfigurationInvalidationScope.SystemConfiguration:
_configurationManager.InvalidateCachedConfiguration(null);
break;
case ConfigurationInvalidationScope.NamedConfiguration when !string.IsNullOrEmpty(invalidation.Target):
_configurationManager.InvalidateCachedConfiguration(invalidation.Target);
break;
case ConfigurationInvalidationScope.LibraryOptions when !string.IsNullOrEmpty(invalidation.Target):
CollectionFolder.InvalidateLibraryOptions(invalidation.Target);
break;
case ConfigurationInvalidationScope.AllLibraryOptions:
CollectionFolder.InvalidateAllLibraryOptions();
break;
default:
return;
}
_logger.LogDebug("Applied {Scope} invalidation for {Target} from {OriginId}.", invalidation.Scope, invalidation.Target, invalidation.OriginId);
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Failed to apply a {Scope} invalidation for {Target}.", invalidation.Scope, invalidation.Target);
}
}
}
}
@@ -0,0 +1,96 @@
using System;
using System.Text.Json;
using Jellyfin.Extensions.Json;
using MediaBrowser.Common.Configuration;
using Microsoft.Extensions.Logging;
using StackExchange.Redis;
namespace Emby.Server.Implementations.Configuration
{
/// <summary>
/// A Redis pub/sub <see cref="IConfigurationInvalidationBus"/>. Notices are broadcast on one channel
/// and every instance but the publisher applies them.
/// </summary>
public sealed class RedisConfigurationInvalidationBus : IConfigurationInvalidationBus
{
private const string ChannelName = "jellyfin:configinvalidation";
private static readonly JsonSerializerOptions _jsonOptions = JsonDefaults.Options;
private readonly ISubscriber _subscriber;
private readonly ILogger<RedisConfigurationInvalidationBus> _logger;
private readonly string _originId;
/// <summary>
/// Initializes a new instance of the <see cref="RedisConfigurationInvalidationBus"/> class.
/// </summary>
/// <param name="redis">The Redis connection multiplexer.</param>
/// <param name="logger">The logger.</param>
public RedisConfigurationInvalidationBus(IConnectionMultiplexer redis, ILogger<RedisConfigurationInvalidationBus> logger)
{
ArgumentNullException.ThrowIfNull(redis);
_subscriber = redis.GetSubscriber();
_logger = logger;
_originId = Environment.GetEnvironmentVariable("JELLYFIN_INSTANCE_ID") ?? Environment.MachineName;
}
/// <inheritdoc />
public void Publish(ConfigurationInvalidationScope scope, string? target)
{
var invalidation = new ConfigurationInvalidation
{
Scope = scope,
Target = target,
OriginId = _originId
};
try
{
// Fire and forget: an admin saving configuration must never wait on, or fail because of,
// the bus. The write has already reached the shared directory by this point.
_subscriber.Publish(
RedisChannel.Literal(ChannelName),
JsonSerializer.Serialize(invalidation, _jsonOptions),
CommandFlags.FireAndForget);
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Failed to publish {Scope} invalidation for {Target}; other instances keep their cached copy until they restart.", scope, target);
}
}
/// <inheritdoc />
public void Subscribe(Action<ConfigurationInvalidation> handler)
{
ArgumentNullException.ThrowIfNull(handler);
try
{
_subscriber.Subscribe(RedisChannel.Literal(ChannelName), (_, value) => Dispatch(handler, value));
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Failed to subscribe to configuration invalidations; this instance keeps its cached configuration until it restarts.");
}
}
private void Dispatch(Action<ConfigurationInvalidation> handler, RedisValue value)
{
try
{
var invalidation = JsonSerializer.Deserialize<ConfigurationInvalidation>(value.ToString(), _jsonOptions);
if (invalidation is null || string.Equals(invalidation.OriginId, _originId, StringComparison.Ordinal))
{
return;
}
handler(invalidation);
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Failed to apply a configuration invalidation.");
}
}
}
}
@@ -2332,7 +2332,10 @@ namespace Emby.Server.Implementations.Library
{
IOrderedEnumerable<BaseItem>? orderedItems = null;
foreach (var orderBy in sortBy.Select(o => GetComparer(o, user)).Where(c => c is not null))
var comparers = sortBy.Select(o => GetComparer(o, user)).Where(c => c is not null).ToList();
items = PrefetchUserData(items, user, comparers);
foreach (var orderBy in comparers)
{
if (orderBy is RandomComparer)
{
@@ -2364,14 +2367,14 @@ namespace Emby.Server.Implementations.Library
{
IOrderedEnumerable<BaseItem>? orderedItems = null;
foreach (var (name, sortOrder) in orderBy)
{
var comparer = GetComparer(name, user);
if (comparer is null)
{
continue;
}
var comparers = orderBy
.Select(o => (Comparer: GetComparer(o.OrderBy, user), o.SortOrder))
.Where(c => c.Comparer is not null)
.ToList();
items = PrefetchUserData(items, user, comparers.Select(c => c.Comparer).ToList());
foreach (var (comparer, sortOrder) in comparers)
{
if (comparer is RandomComparer)
{
var randomItems = items.ToArray();
@@ -2397,6 +2400,31 @@ namespace Emby.Server.Implementations.Library
return orderedItems ?? items;
}
// The user comparers read user data per item, so without one batched read up front an
// in-memory sort would issue a database round trip per comparison.
private IEnumerable<BaseItem> PrefetchUserData(IEnumerable<BaseItem> items, User? user, IReadOnlyList<IBaseItemComparer?> comparers)
{
if (user is null)
{
return items;
}
var userComparers = comparers.OfType<IUserBaseItemComparer>().ToList();
if (userComparers.Count == 0)
{
return items;
}
var itemList = items as IReadOnlyList<BaseItem> ?? items.ToList();
var userData = _userDataManager.GetUserDataBatch(itemList, user);
foreach (var comparer in userComparers)
{
comparer.PrefetchedUserData = userData;
}
return itemList;
}
/// <summary>
/// Gets the comparer.
/// </summary>
@@ -2,10 +2,8 @@
using System;
using System.Collections.Generic;
using System.Globalization;
using System.Linq;
using System.Threading;
using BitFaster.Caching.Lru;
using Jellyfin.Database.Implementations;
using Jellyfin.Database.Implementations.Entities;
using MediaBrowser.Controller.Configuration;
@@ -27,7 +25,6 @@ namespace Emby.Server.Implementations.Library
{
private readonly IServerConfigurationManager _config;
private readonly IDbContextFactory<JellyfinDbContext> _repository;
private readonly FastConcurrentLru<string, UserItemData> _cache;
/// <summary>
/// Initializes a new instance of the <see cref="UserDataManager"/> class.
@@ -40,7 +37,6 @@ namespace Emby.Server.Implementations.Library
{
_config = config;
_repository = repository;
_cache = new FastConcurrentLru<string, UserItemData>(Environment.ProcessorCount, _config.Configuration.CacheSize, StringComparer.OrdinalIgnoreCase);
}
/// <inheritdoc />
@@ -77,11 +73,6 @@ namespace Emby.Server.Implementations.Library
dbContext.SaveChanges();
transaction.Commit();
var userId = user.InternalId;
var cacheKey = GetCacheKey(userId, item.Id);
_cache.AddOrUpdate(cacheKey, userData);
item.UserData = dbContext.UserData.Where(e => e.ItemId == item.Id).AsNoTracking().ToArray(); // rehydrate the cached userdata
UserDataSaved?.Invoke(this, new UserDataSaveEventArgs
{
Keys = keys,
@@ -180,64 +171,41 @@ namespace Emby.Server.Implementations.Library
/// <inheritdoc />
public Dictionary<Guid, UserItemData> GetUserDataBatch(IReadOnlyList<BaseItem> items, User user)
{
ArgumentNullException.ThrowIfNull(items);
ArgumentNullException.ThrowIfNull(user);
var result = new Dictionary<Guid, UserItemData>(items.Count);
var itemsNeedingQuery = new List<(BaseItem Item, List<string> Keys)>();
foreach (var item in items)
{
var cacheKey = GetCacheKey(user.InternalId, item.Id);
if (_cache.TryGet(cacheKey, out var cachedData))
{
result[item.Id] = cachedData;
}
else
{
var userDataRow = ResolveUserDataRow(item, item.UserData?.Where(e => e.UserId.Equals(user.Id)));
var userData = userDataRow is not null ? Map(userDataRow) : null;
if (userData is not null)
{
result[item.Id] = userData;
_cache.AddOrUpdate(cacheKey, userData);
}
else
{
var keys = item.GetUserDataKeys();
itemsNeedingQuery.Add((item, keys));
}
}
}
if (itemsNeedingQuery.Count == 0)
if (items.Count == 0)
{
return result;
}
// Build a single query for all missing items. Fetch rows by item alone so rows kept
// under keys from older metadata resolve the same way as the in-memory path.
var allItemIds = itemsNeedingQuery.Select(x => x.Item.Id).ToList();
// Fetch rows by item alone so rows kept under keys from older metadata resolve the same
// way as the single item path.
var itemIds = items.Select(e => e.Id).Distinct().ToList();
using var context = _repository.CreateDbContext();
var userDataArray = context.UserData
var userDataByItem = context.UserData
.AsNoTracking()
.Where(e => e.UserId.Equals(user.Id))
.WhereOneOrMany(allItemIds, e => e.ItemId)
.ToArray();
.WhereOneOrMany(itemIds, e => e.ItemId)
.ToArray()
.GroupBy(e => e.ItemId)
.ToDictionary(g => g.Key, g => g.ToArray());
var userDataByItem = userDataArray.GroupBy(e => e.ItemId).ToDictionary(g => g.Key, g => g.ToArray());
foreach (var (item, keys) in itemsNeedingQuery)
foreach (var item in items)
{
UserItemData userData;
if (userDataByItem.TryGetValue(item.Id, out var itemUserData) && itemUserData.Length > 0)
if (result.ContainsKey(item.Id))
{
userData = Map(ResolveUserDataRow(item, itemUserData)!);
}
else
{
userData = new UserItemData { Key = keys.Count > 0 ? keys[0] : string.Empty };
continue;
}
result[item.Id] = userData;
var cacheKey = GetCacheKey(user.InternalId, item.Id);
_cache.AddOrUpdate(cacheKey, userData);
var row = userDataByItem.TryGetValue(item.Id, out var itemUserData)
? ResolveUserDataRow(item, itemUserData)
: null;
result[item.Id] = row is not null
? Map(row)
: new UserItemData { Key = item.GetUserDataKeys().FirstOrDefault() ?? string.Empty };
}
return result;
@@ -340,20 +308,19 @@ namespace Emby.Server.Implementations.Library
return result;
}
/// <summary>
/// Gets the internal key.
/// </summary>
/// <returns>System.String.</returns>
private static string GetCacheKey(long internalUserId, Guid itemId)
{
return internalUserId.ToString(CultureInfo.InvariantCulture) + "-" + itemId.ToString("N", CultureInfo.InvariantCulture);
}
/// <inheritdoc />
public UserItemData? GetUserData(User user, BaseItem item)
{
ArgumentNullException.ThrowIfNull(user);
var row = ResolveUserDataRow(item, item.UserData?.Where(e => e.UserId.Equals(user.Id)));
ArgumentNullException.ThrowIfNull(item);
using var dbContext = _repository.CreateDbContext();
var rows = dbContext.UserData
.AsNoTracking()
.Where(e => e.ItemId == item.Id && e.UserId == user.Id)
.ToArray();
var row = ResolveUserDataRow(item, rows);
return row is not null ? Map(row) : new UserItemData()
{
Key = item.GetUserDataKeys()[0],
@@ -536,16 +503,6 @@ namespace Emby.Server.Implementations.Library
}
dbContext.SaveChanges();
var cacheKey = GetCacheKey(user.InternalId, item.Id);
if (_cache.TryGet(cacheKey, out var cached))
{
cached.AudioStreamIndex = null;
cached.SubtitleStreamIndex = null;
_cache.AddOrUpdate(cacheKey, cached);
}
item.UserData = dbContext.UserData.Where(e => e.ItemId == item.Id).AsNoTracking().ToArray();
}
}
}
@@ -10,9 +10,14 @@ using StackExchange.Redis;
namespace Emby.Server.Implementations.MediaEncoding;
/// <summary>
/// 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.
/// Reports the round trip to the configured Redis transcode session store once at startup, so the
/// state of cross-pod takeover is visible where the server is started.
/// </summary>
/// <remarks>
/// This reports, it does not gate. <see cref="ApplicationHost.InitializeServices"/> has already read the
/// same connection for quick connect by the time this runs and has stopped startup if it could not be
/// reached, so the error branch here only covers a store that went away in between.
/// </remarks>
public sealed class TranscodeStoreConnectivityProbe : IHostedService
{
private readonly IServiceProvider _serviceProvider;
@@ -34,7 +39,7 @@ public sealed class TranscodeStoreConnectivityProbe : IHostedService
{
try
{
// Resolved here rather than injected: connecting must not be able to abort startup.
// Resolved here rather than injected so a store lost after the quick connect gate is reported.
var redis = _serviceProvider.GetRequiredService<IConnectionMultiplexer>();
var roundTrip = await redis.GetDatabase().PingAsync().ConfigureAwait(false);
@@ -0,0 +1,14 @@
namespace Emby.Server.Implementations
{
/// <summary>
/// What a configuration update carrying different ports requires of the instance reading it.
/// </summary>
/// <param name="RequiresRestart">
/// Whether this process is still bound to a port the shared configuration no longer names, and so has
/// to report a pending restart.
/// </param>
/// <param name="ClearsPortAuthorization">
/// Whether this instance is the one that has to clear the port authorization flag and save it.
/// </param>
internal readonly record struct PortChangeOutcome(bool RequiresRestart, bool ClearsPortAuthorization);
}
@@ -1,7 +1,5 @@
using System;
using System.Collections.Concurrent;
using System.Globalization;
using System.Linq;
using System.Security.Cryptography;
using System.Threading.Tasks;
using MediaBrowser.Common.Extensions;
@@ -30,12 +28,10 @@ namespace Emby.Server.Implementations.QuickConnect
/// </summary>
private const int Timeout = 10;
private readonly ConcurrentDictionary<string, QuickConnectResult> _currentRequests = new();
private readonly ConcurrentDictionary<string, (DateTime Timestamp, AuthenticationResult AuthenticationResult)> _authorizedSecrets = new();
private readonly IServerConfigurationManager _config;
private readonly ILogger<QuickConnectManager> _logger;
private readonly ISessionManager _sessionManager;
private readonly IQuickConnectStore _store;
/// <summary>
/// Initializes a new instance of the <see cref="QuickConnectManager"/> class.
@@ -44,14 +40,17 @@ namespace Emby.Server.Implementations.QuickConnect
/// <param name="config">Configuration.</param>
/// <param name="logger">Logger.</param>
/// <param name="sessionManager">Session Manager.</param>
/// <param name="store">Quick connect store.</param>
public QuickConnectManager(
IServerConfigurationManager config,
ILogger<QuickConnectManager> logger,
ISessionManager sessionManager)
ISessionManager sessionManager,
IQuickConnectStore store)
{
_config = config;
_logger = logger;
_sessionManager = sessionManager;
_store = store;
}
/// <inheritdoc />
@@ -69,7 +68,7 @@ namespace Emby.Server.Implementations.QuickConnect
}
/// <inheritdoc/>
public QuickConnectResult TryConnect(AuthorizationInfo authorizationInfo)
public async Task<QuickConnectResult> TryConnect(AuthorizationInfo authorizationInfo)
{
ArgumentException.ThrowIfNullOrEmpty(authorizationInfo.DeviceId);
ArgumentException.ThrowIfNullOrEmpty(authorizationInfo.Device);
@@ -77,7 +76,6 @@ namespace Emby.Server.Implementations.QuickConnect
ArgumentException.ThrowIfNullOrEmpty(authorizationInfo.Version);
AssertActive();
ExpireRequests();
var secret = GenerateSecureRandom();
var code = GenerateCode();
@@ -90,19 +88,17 @@ namespace Emby.Server.Implementations.QuickConnect
authorizationInfo.Client,
authorizationInfo.Version);
_currentRequests[code] = result;
await _store.SetRequestAsync(result, ExpiryOf(result)).ConfigureAwait(false);
return result;
}
/// <inheritdoc/>
public QuickConnectResult CheckRequestStatus(string secret)
public async Task<QuickConnectResult> CheckRequestStatus(string secret)
{
AssertActive();
ExpireRequests();
string code = _currentRequests.Where(x => x.Value.Secret == secret).Select(x => x.Value.Code).DefaultIfEmpty(string.Empty).First();
if (!_currentRequests.TryGetValue(code, out QuickConnectResult? result))
var result = await _store.GetRequestBySecretAsync(secret).ConfigureAwait(false);
if (result is null)
{
throw new ResourceNotFoundException("Unable to find request with provided secret");
}
@@ -136,21 +132,27 @@ namespace Emby.Server.Implementations.QuickConnect
public async Task<bool> AuthorizeRequest(Guid userId, string code)
{
AssertActive();
ExpireRequests();
if (!_currentRequests.TryGetValue(code, out QuickConnectResult? result))
var result = await _store.GetRequestByCodeAsync(code).ConfigureAwait(false);
if (result is null)
{
throw new ResourceNotFoundException("Unable to find request");
}
if (result.Authenticated)
{
throw new InvalidOperationException("Request is already authorized");
throw new ConflictException("Request is already authorized");
}
// Change the time on the request so it expires one minute into the future. It can't expire immediately as otherwise some clients wouldn't ever see that they have been authenticated.
result.DateAdded = DateTime.UtcNow.Add(TimeSpan.FromMinutes(1));
// The guard above is a read on shared state, so it cannot settle a race between instances; the claim can.
if (!await _store.TryClaimAuthorizationAsync(result.Secret, ExpiryOf(result)).ConfigureAwait(false))
{
throw await RefusedClaimAsync(result.Secret).ConfigureAwait(false);
}
var authenticationResult = await _sessionManager.AuthenticateDirect(new AuthenticationRequest
{
UserId = userId,
@@ -160,9 +162,10 @@ namespace Emby.Server.Implementations.QuickConnect
AppVersion = result.AppVersion
}).ConfigureAwait(false);
_authorizedSecrets[result.Secret] = (DateTime.UtcNow, authenticationResult);
result.Authenticated = true;
_currentRequests[code] = result;
await _store.SetAuthorizationAsync(result.Secret, authenticationResult, DateTime.UtcNow.AddMinutes(Timeout)).ConfigureAwait(false);
await _store.SetRequestAsync(result, ExpiryOf(result)).ConfigureAwait(false);
_logger.LogDebug("Authorizing device with code {Code} to login as user {UserId}", code, userId);
@@ -170,17 +173,33 @@ namespace Emby.Server.Implementations.QuickConnect
}
/// <inheritdoc/>
public AuthenticationResult GetAuthorizedRequest(string secret)
public async Task<AuthenticationResult> GetAuthorizedRequest(string secret)
{
AssertActive();
ExpireRequests();
if (!_authorizedSecrets.TryGetValue(secret, out var result))
var result = await _store.GetAuthorizationAsync(secret).ConfigureAwait(false);
if (result is null)
{
throw new ResourceNotFoundException("Unable to find request");
}
return result.AuthenticationResult;
return result;
}
private static DateTime ExpiryOf(QuickConnectResult request) => request.DateAdded.AddMinutes(Timeout);
/// <summary>
/// Explains a refused claim. The claim outlives a failed mint on purpose, so it can mean either
/// that the request is authorized or that authorizing it did not finish; the two are told apart
/// by re-reading the request rather than reported as the same thing.
/// </summary>
private async Task<ConflictException> RefusedClaimAsync(string secret)
{
var current = await _store.GetRequestBySecretAsync(secret).ConfigureAwait(false);
return current?.Authenticated == true
? new ConflictException("Request is already authorized")
: new ConflictException("Request is being authorized elsewhere, or an earlier attempt to authorize it did not complete. Start quick connect again for a new code.");
}
private string GenerateSecureRandom(int length = 32)
@@ -190,42 +209,5 @@ namespace Emby.Server.Implementations.QuickConnect
return Convert.ToHexString(bytes);
}
/// <summary>
/// Expire quick connect requests that are over the time limit. If <paramref name="expireAll"/> is true, all requests are unconditionally expired.
/// </summary>
/// <param name="expireAll">If true, all requests will be expired.</param>
private void ExpireRequests(bool expireAll = false)
{
// All requests before this timestamp have expired
var minTime = DateTime.UtcNow.AddMinutes(-Timeout);
// Expire stale connection requests
foreach (var (_, currentRequest) in _currentRequests)
{
if (expireAll || currentRequest.DateAdded < minTime)
{
var code = currentRequest.Code;
_logger.LogDebug("Removing expired request {Code}", code);
if (!_currentRequests.TryRemove(code, out _))
{
_logger.LogWarning("Request {Code} already expired", code);
}
}
}
foreach (var (secret, (timestamp, _)) in _authorizedSecrets)
{
if (expireAll || timestamp < minTime)
{
_logger.LogDebug("Removing expired secret {Secret}", secret);
if (!_authorizedSecrets.TryRemove(secret, out _))
{
_logger.LogWarning("Secret {Secret} already expired", secret);
}
}
}
}
}
}
@@ -0,0 +1,164 @@
using System;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using Jellyfin.Extensions.Json;
using MediaBrowser.Common.Extensions;
using MediaBrowser.Controller.Authentication;
using MediaBrowser.Controller.QuickConnect;
using MediaBrowser.Model.QuickConnect;
using Microsoft.Extensions.Logging;
using StackExchange.Redis;
namespace Emby.Server.Implementations.QuickConnect;
/// <summary>
/// A Redis-backed <see cref="IQuickConnectStore"/> that lets the initiate, authorize and exchange legs
/// of a quick connect flow land on different instances. Expiry is the key TTL and an authorization is
/// claimed with a Lua check-and-set, so only one instance can ever mint a given secret's access token.
/// </summary>
/// <remarks>
/// There is no local fallback: a call Redis did not answer is inconclusive, and reporting it as a miss
/// would tell a polling client its secret is invalid. Quick connect is unavailable for as long as Redis
/// is, which password login is not.
/// </remarks>
public sealed class RedisQuickConnectStore : IQuickConnectStore
{
private const string KeyPrefix = "jellyfin:quickconnect:";
/// <summary>
/// Lua script writing the two keys a request is resolvable by in one step, so it can never be
/// reachable by its secret while the code the user is reading off the screen resolves to nothing.
/// </summary>
private const string SetRequestScript = @"
redis.call('SET', KEYS[1], ARGV[1], 'PX', ARGV[3])
redis.call('SET', KEYS[2], ARGV[2], 'PX', ARGV[3])
return 1";
/// <summary>
/// Lua script for the atomic claim of the sole right to authorize a request: the request has to
/// exist and not already be authorized, and the claim marker is taken with <c>SET NX</c>, so of two
/// instances racing on one code exactly one goes on to mint an access token.
/// </summary>
private const string ClaimAuthorizationScript = @"
local raw = redis.call('GET', KEYS[1])
if not raw then return 0 end
if cjson.decode(raw)['Authenticated'] then return 0 end
if redis.call('SET', KEYS[2], '1', 'NX', 'PX', ARGV[1]) then return 1 end
return 0";
private readonly IDatabase _db;
private readonly ILogger<RedisQuickConnectStore> _logger;
/// <summary>
/// Initializes a new instance of the <see cref="RedisQuickConnectStore"/> class.
/// </summary>
/// <param name="redis">The Redis connection multiplexer.</param>
/// <param name="logger">The logger.</param>
public RedisQuickConnectStore(IConnectionMultiplexer redis, ILogger<RedisQuickConnectStore> logger)
{
ArgumentNullException.ThrowIfNull(redis);
_db = redis.GetDatabase();
_logger = logger;
}
/// <inheritdoc />
public async Task<QuickConnectResult?> GetRequestBySecretAsync(string secret, CancellationToken cancellationToken = default)
{
var raw = await CallAsync(() => _db.StringGetAsync(RequestKey(secret))).ConfigureAwait(false);
// Deserialization is outside the guard: a malformed stored value is a fault of its own, not Redis
// being unavailable.
return raw.HasValue ? JsonSerializer.Deserialize<QuickConnectResult>(raw.ToString(), JsonDefaults.Options) : null;
}
/// <inheritdoc />
public async Task<QuickConnectResult?> GetRequestByCodeAsync(string code, CancellationToken cancellationToken = default)
{
var secret = await CallAsync(() => _db.StringGetAsync(CodeKey(code))).ConfigureAwait(false);
return secret.HasValue
? await GetRequestBySecretAsync(secret.ToString(), cancellationToken).ConfigureAwait(false)
: null;
}
/// <inheritdoc />
public async Task SetRequestAsync(QuickConnectResult request, DateTime expiresUtc, CancellationToken cancellationToken = default)
{
ArgumentNullException.ThrowIfNull(request);
var ttl = expiresUtc - DateTime.UtcNow;
if (ttl <= TimeSpan.Zero)
{
return;
}
var json = JsonSerializer.Serialize(request, JsonDefaults.Options);
await CallAsync(() => _db.ScriptEvaluateAsync(
SetRequestScript,
keys: new RedisKey[] { RequestKey(request.Secret), CodeKey(request.Code) },
values: new RedisValue[] { json, request.Secret, (long)ttl.TotalMilliseconds })).ConfigureAwait(false);
}
/// <inheritdoc />
public async Task<bool> TryClaimAuthorizationAsync(string secret, DateTime expiresUtc, CancellationToken cancellationToken = default)
{
var ttl = expiresUtc - DateTime.UtcNow;
if (ttl <= TimeSpan.Zero)
{
return false;
}
var claimed = (long?)await CallAsync(() => _db.ScriptEvaluateAsync(
ClaimAuthorizationScript,
keys: new RedisKey[] { RequestKey(secret), ClaimKey(secret) },
values: new RedisValue[] { (long)ttl.TotalMilliseconds })).ConfigureAwait(false);
return claimed == 1;
}
/// <inheritdoc />
public async Task SetAuthorizationAsync(string secret, AuthenticationResult authenticationResult, DateTime expiresUtc, CancellationToken cancellationToken = default)
{
var ttl = expiresUtc - DateTime.UtcNow;
if (ttl <= TimeSpan.Zero)
{
return;
}
var json = JsonSerializer.Serialize(authenticationResult, JsonDefaults.Options);
await CallAsync(() => _db.StringSetAsync(AuthorizationKey(secret), json, ttl)).ConfigureAwait(false);
}
/// <inheritdoc />
public async Task<AuthenticationResult?> GetAuthorizationAsync(string secret, CancellationToken cancellationToken = default)
{
var raw = await CallAsync(() => _db.StringGetAsync(AuthorizationKey(secret))).ConfigureAwait(false);
return raw.HasValue
? JsonSerializer.Deserialize<AuthenticationResult>(raw.ToString(), JsonDefaults.Options)
: null;
}
private static string RequestKey(string secret) => KeyPrefix + "request:" + secret;
private static string CodeKey(string code) => KeyPrefix + "code:" + code;
private static string ClaimKey(string secret) => KeyPrefix + "claim:" + secret;
private static string AuthorizationKey(string secret) => KeyPrefix + "auth:" + secret;
private async Task<T> CallAsync<T>(Func<Task<T>> call)
{
try
{
return await call().ConfigureAwait(false);
}
catch (Exception exception) when (exception is RedisException or RedisCommandException or TimeoutException)
{
_logger.LogError(exception, "Quick connect state could not be reached in Redis.");
throw new ServiceUnavailableException("Quick connect is temporarily unavailable.", exception);
}
}
}
@@ -13,6 +13,11 @@ namespace Emby.Server.Implementations.ScheduledTasks;
/// TTL lease on a shared key. The lease value is this instance's pod identity; a leader that keeps
/// renewing retains the lease, and any instance can claim it once the previous leader's lease expires.
/// </summary>
/// <remarks>
/// This fails open on an unreachable Redis while quick connect's startup read fails closed on the same
/// connection. They do not compete: the startup read decides whether the instance runs at all, and this
/// only decides what a running instance does about a store that went away afterwards.
/// </remarks>
public sealed class RedisScanLeaderLease : IScanLeaderLease
{
private const string LeaderKey = "jellyfin:scanleader";
@@ -1,6 +1,7 @@
#nullable disable
using System;
using System.Collections.Generic;
using Jellyfin.Data.Enums;
using Jellyfin.Database.Implementations.Entities;
using MediaBrowser.Controller.Entities;
@@ -27,6 +28,12 @@ namespace Emby.Server.Implementations.Sorting
/// <value>The user manager.</value>
public IUserManager UserManager { get; set; }
/// <summary>
/// Gets or sets the prefetched user data.
/// </summary>
/// <value>The prefetched user data.</value>
public IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
/// <summary>
/// Gets or sets the user data manager.
/// </summary>
@@ -57,7 +64,7 @@ namespace Emby.Server.Implementations.Sorting
/// <returns>DateTime.</returns>
private DateTime GetDate(BaseItem x)
{
var userdata = UserDataManager.GetUserData(User, x);
var userdata = this.GetUserData(x);
if (userdata is not null && userdata.LastPlayedDate.HasValue)
{
@@ -1,6 +1,8 @@
#nullable disable
#pragma warning disable CS1591
using System;
using System.Collections.Generic;
using Jellyfin.Data.Enums;
using Jellyfin.Database.Implementations.Entities;
using MediaBrowser.Controller.Entities;
@@ -35,6 +37,12 @@ namespace Emby.Server.Implementations.Sorting
/// <value>The user manager.</value>
public IUserManager UserManager { get; set; }
/// <summary>
/// Gets or sets the prefetched user data.
/// </summary>
/// <value>The prefetched user data.</value>
public IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
/// <summary>
/// Compares the specified x.
/// </summary>
@@ -53,7 +61,7 @@ namespace Emby.Server.Implementations.Sorting
/// <returns>DateTime.</returns>
private int GetValue(BaseItem x)
{
return x.IsFavoriteOrLiked(User, userItemData: null) ? 0 : 1;
return x.IsFavoriteOrLiked(User, this.GetUserData(x)) ? 0 : 1;
}
}
}
@@ -2,6 +2,8 @@
#pragma warning disable CS1591
using System;
using System.Collections.Generic;
using Jellyfin.Data.Enums;
using Jellyfin.Database.Implementations.Entities;
using MediaBrowser.Controller.Entities;
@@ -36,6 +38,12 @@ namespace Emby.Server.Implementations.Sorting
/// <value>The user manager.</value>
public IUserManager UserManager { get; set; }
/// <summary>
/// Gets or sets the prefetched user data.
/// </summary>
/// <value>The prefetched user data.</value>
public IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
/// <summary>
/// Compares the specified x.
/// </summary>
@@ -54,7 +62,7 @@ namespace Emby.Server.Implementations.Sorting
/// <returns>DateTime.</returns>
private int GetValue(BaseItem x)
{
return x.IsPlayed(User, userItemData: null) ? 0 : 1;
return x.IsPlayed(User, this.GetUserData(x)) ? 0 : 1;
}
}
}
@@ -2,6 +2,8 @@
#pragma warning disable CS1591
using System;
using System.Collections.Generic;
using Jellyfin.Data.Enums;
using Jellyfin.Database.Implementations.Entities;
using MediaBrowser.Controller.Entities;
@@ -36,6 +38,12 @@ namespace Emby.Server.Implementations.Sorting
/// <value>The user manager.</value>
public IUserManager UserManager { get; set; }
/// <summary>
/// Gets or sets the prefetched user data.
/// </summary>
/// <value>The prefetched user data.</value>
public IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
/// <summary>
/// Compares the specified x.
/// </summary>
@@ -54,7 +62,7 @@ namespace Emby.Server.Implementations.Sorting
/// <returns>DateTime.</returns>
private int GetValue(BaseItem x)
{
return x.IsUnplayed(User, userItemData: null) ? 0 : 1;
return x.IsUnplayed(User, this.GetUserData(x)) ? 0 : 1;
}
}
}
@@ -1,5 +1,7 @@
#nullable disable
using System;
using System.Collections.Generic;
using Jellyfin.Data.Enums;
using Jellyfin.Database.Implementations.Entities;
using MediaBrowser.Controller.Entities;
@@ -38,6 +40,12 @@ namespace Emby.Server.Implementations.Sorting
/// <value>The user manager.</value>
public IUserManager UserManager { get; set; }
/// <summary>
/// Gets or sets the prefetched user data.
/// </summary>
/// <value>The prefetched user data.</value>
public IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData { get; set; }
/// <summary>
/// Compares the specified x.
/// </summary>
@@ -56,7 +64,7 @@ namespace Emby.Server.Implementations.Sorting
/// <returns>DateTime.</returns>
private int GetValue(BaseItem x)
{
var userdata = UserDataManager.GetUserData(User, x);
var userdata = this.GetUserData(x);
return userdata is null ? 0 : userdata.PlayCount;
}
+127 -62
View File
@@ -13,6 +13,7 @@ using MediaBrowser.Controller.Configuration;
using MediaBrowser.Controller.Dto;
using MediaBrowser.Controller.Entities;
using MediaBrowser.Controller.Library;
using MediaBrowser.Controller.Persistence;
using MediaBrowser.Controller.TV;
using MediaBrowser.Model.Querying;
using Episode = MediaBrowser.Controller.Entities.TV.Episode;
@@ -124,53 +125,100 @@ namespace Emby.Server.Implementations.TV
var batchResult = _libraryManager.GetNextUpEpisodesBatch(query, seriesKeys, includeSpecials, includeRewatching);
var nextUpList = new List<(DateTime LastWatchedDate, Episode Episode)>();
var results = new List<NextUpEpisodeBatchResult>(seriesKeys.Count);
foreach (var seriesKey in seriesKeys)
{
if (!batchResult.TryGetValue(seriesKey, out var result))
if (batchResult.TryGetValue(seriesKey, out var result))
{
continue;
results.Add(result);
}
}
var nextEpisode = DetermineNextEpisode(result, user, includeSpecials, request.EnableResumable, false);
// The selection below tests the played state of every episode it considers, so read the whole
// series batch in one query rather than one query per series.
var selectionCandidates = new List<BaseItem>();
foreach (var result in results)
{
AddCandidate(selectionCandidates, result.NextUp);
AddCandidate(selectionCandidates, result.LastWatched);
AddCandidate(selectionCandidates, result.NextPlayedForRewatching);
AddCandidate(selectionCandidates, result.LastWatchedForRewatching);
if (result.Specials is not null)
{
selectionCandidates.AddRange(result.Specials);
}
}
var selectionUserData = _userDataManager.GetUserDataBatch(selectionCandidates, user);
var candidates = new List<NextUpCandidate>();
foreach (var result in results)
{
var nextEpisode = SelectNextEpisode(result, user, includeSpecials, includePlayed: false, selectionUserData);
if (nextEpisode is not null)
{
// The last played date and the version that was actually played live on the version item's user data
// The played state propagated to the sibling versions carries no date
var (playedVersion, lastPlayedDate) = GetMostRecentlyPlayedVersion(result.LastWatched, user);
nextEpisode = GetPreferredVersion(nextEpisode, result.LastWatched, playedVersion);
DateTime lastWatchedDate = DateTime.MinValue;
if (result.LastWatched is not null)
{
lastWatchedDate = lastPlayedDate ?? DateTime.MinValue.AddDays(1);
}
nextUpList.Add((lastWatchedDate, nextEpisode));
candidates.Add(new NextUpCandidate(nextEpisode, result.LastWatched, !request.EnableResumable));
}
if (includeRewatching)
{
var nextPlayedEpisode = DetermineNextEpisodeForRewatching(result, user, includeSpecials);
var nextPlayedEpisode = SelectNextEpisode(result, user, includeSpecials, includePlayed: true, selectionUserData);
if (nextPlayedEpisode is not null)
{
var (playedVersion, lastPlayedDate) = GetMostRecentlyPlayedVersion(result.LastWatchedForRewatching, user);
nextPlayedEpisode = GetPreferredVersion(nextPlayedEpisode, result.LastWatchedForRewatching, playedVersion);
DateTime rewatchLastWatchedDate = DateTime.MinValue;
if (result.LastWatchedForRewatching is not null)
{
rewatchLastWatchedDate = lastPlayedDate ?? DateTime.MinValue.AddDays(1);
}
nextUpList.Add((rewatchLastWatchedDate, nextPlayedEpisode));
// A rewatch suggestion is dropped once it has been resumed, whatever the request asked for.
candidates.Add(new NextUpCandidate(nextPlayedEpisode, result.LastWatchedForRewatching, true));
}
}
}
// The resume progress may live on an alternate version, so read every version in one query.
var episodeVersions = new List<BaseItem>();
foreach (var candidate in candidates)
{
if (candidate.DropWhenResumed)
{
candidate.EpisodeVersions = candidate.Episode.GetAllVersions();
episodeVersions.AddRange(candidate.EpisodeVersions);
}
}
if (episodeVersions.Count > 0)
{
var resumeUserData = _userDataManager.GetUserDataBatch(episodeVersions, user);
candidates.RemoveAll(candidate => candidate.EpisodeVersions
.Any(version => GetUserData(user, version, resumeUserData)?.PlaybackPositionTicks > 0));
}
// The last played date and the version that was actually played live on the version item's user data
// The played state propagated to the sibling versions carries no date
var lastWatchedVersions = new List<BaseItem>();
foreach (var candidate in candidates)
{
if (candidate.LastWatched is Video lastWatchedVideo)
{
candidate.LastWatchedVersions = lastWatchedVideo.GetAllVersions();
lastWatchedVersions.AddRange(candidate.LastWatchedVersions);
}
}
var lastWatchedUserData = _userDataManager.GetUserDataBatch(lastWatchedVersions, user);
var nextUpList = new List<(DateTime LastWatchedDate, Episode Episode)>(candidates.Count);
foreach (var candidate in candidates)
{
var (playedVersion, lastPlayedDate) = GetMostRecentlyPlayedVersion(candidate.LastWatchedVersions, user, lastWatchedUserData);
var nextEpisode = GetPreferredVersion(candidate.Episode, candidate.LastWatched, playedVersion);
DateTime lastWatchedDate = DateTime.MinValue;
if (candidate.LastWatched is not null)
{
lastWatchedDate = lastPlayedDate ?? DateTime.MinValue.AddDays(1);
}
nextUpList.Add((lastWatchedDate, nextEpisode));
}
var sortedEpisodes = nextUpList
.OrderByDescending(x => x.LastWatchedDate)
.Select(x => (BaseItem)x.Episode);
@@ -178,12 +226,25 @@ namespace Emby.Server.Implementations.TV
return GetResult(sortedEpisodes, request);
}
private Episode? DetermineNextEpisode(
MediaBrowser.Controller.Persistence.NextUpEpisodeBatchResult result,
private static void AddCandidate(List<BaseItem> candidates, BaseItem? item)
{
if (item is not null)
{
candidates.Add(item);
}
}
private UserItemData? GetUserData(User user, BaseItem item, IReadOnlyDictionary<Guid, UserItemData> prefetchedUserData)
=> prefetchedUserData.TryGetValue(item.Id, out var userData)
? userData
: _userDataManager.GetUserData(user, item);
private Episode? SelectNextEpisode(
NextUpEpisodeBatchResult result,
User user,
bool includeSpecials,
bool includeResumable,
bool includePlayed)
bool includePlayed,
IReadOnlyDictionary<Guid, UserItemData> prefetchedUserData)
{
var nextEpisode = (includePlayed ? result.NextPlayedForRewatching : result.NextUp) as Episode;
var lastWatchedEpisode = (includePlayed ? result.LastWatchedForRewatching : result.LastWatched) as Episode;
@@ -217,60 +278,41 @@ namespace Emby.Server.Implementations.TV
if (!includePlayed)
{
sortedEpisodes = sortedEpisodes.Where(episode => _userDataManager.GetUserData(user, episode) is not { Played: true });
sortedEpisodes = sortedEpisodes.Where(episode => GetUserData(user, episode, prefetchedUserData) is not { Played: true });
}
nextEpisode = sortedEpisodes.FirstOrDefault();
}
}
if (nextEpisode is not null && !includeResumable)
{
// The resume progress may live on an alternate version
foreach (var version in nextEpisode.GetAllVersions())
{
if (_userDataManager.GetUserData(user, version)?.PlaybackPositionTicks > 0)
{
return null;
}
}
}
return nextEpisode;
}
private Episode? DetermineNextEpisodeForRewatching(
MediaBrowser.Controller.Persistence.NextUpEpisodeBatchResult result,
User user,
bool includeSpecials)
{
return DetermineNextEpisode(result, user, includeSpecials, includeResumable: false, includePlayed: true);
}
/// <summary>
/// Gets the version of the last watched episode that was actually played, together with its last played date.
/// The version that was played carries the most recent LastPlayedDate.
/// dates.
/// </summary>
/// <param name="lastWatched">The last watched episode (any version).</param>
/// <param name="versions">The versions of the last watched episode.</param>
/// <param name="user">The user.</param>
/// <param name="prefetchedUserData">User data read for every version up front.</param>
/// <returns>The played version and its last played date.</returns>
private (Video? PlayedVersion, DateTime? LastPlayedDate) GetMostRecentlyPlayedVersion(BaseItem? lastWatched, User user)
private (Video? PlayedVersion, DateTime? LastPlayedDate) GetMostRecentlyPlayedVersion(
IReadOnlyList<Video> versions,
User user,
IReadOnlyDictionary<Guid, UserItemData> prefetchedUserData)
{
if (lastWatched is not Video lastWatchedVideo)
if (versions.Count == 0)
{
return (null, null);
}
var versions = lastWatchedVideo.GetAllVersions();
var userDataByVersion = _userDataManager.GetUserDataBatch(versions, user);
var playedVersion = VersionPlaybackSelector.SelectMostRecentlyPlayed(
versions,
version => userDataByVersion.GetValueOrDefault(version.Id),
version => GetUserData(user, version, prefetchedUserData),
data => data.LastPlayedDate.HasValue);
return (playedVersion, playedVersion is null ? null : userDataByVersion[playedVersion.Id].LastPlayedDate);
return (playedVersion, playedVersion is null ? null : GetUserData(user, playedVersion, prefetchedUserData)?.LastPlayedDate);
}
/// <summary>
@@ -346,5 +388,28 @@ namespace Emby.Server.Implementations.TV
totalCount,
items.ToArray());
}
/// <summary>
/// An episode picked for Next Up, together with the versions its user data is read from.
/// </summary>
private sealed class NextUpCandidate
{
public NextUpCandidate(Episode episode, BaseItem? lastWatched, bool dropWhenResumed)
{
Episode = episode;
LastWatched = lastWatched;
DropWhenResumed = dropWhenResumed;
}
public Episode Episode { get; }
public BaseItem? LastWatched { get; }
public bool DropWhenResumed { get; }
public IReadOnlyList<Video> EpisodeVersions { get; set; } = [];
public IReadOnlyList<Video> LastWatchedVersions { get; set; } = [];
}
}
}
@@ -50,16 +50,18 @@ public class QuickConnectController : BaseJellyfinApiController
/// </summary>
/// <response code="200">Quick connect request successfully created.</response>
/// <response code="401">Quick connect is not active on this server.</response>
/// <response code="503">Quick connect state is unavailable.</response>
/// <returns>A <see cref="QuickConnectResult"/> with a secret and code for future use or an error message.</returns>
[HttpPost("Initiate")]
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status401Unauthorized)]
[ProducesResponseType(StatusCodes.Status503ServiceUnavailable)]
public async Task<ActionResult<QuickConnectResult>> InitiateQuickConnect()
{
try
{
var auth = await _authContext.GetAuthorizationInfo(Request).ConfigureAwait(false);
return _quickConnect.TryConnect(auth);
return await _quickConnect.TryConnect(auth).ConfigureAwait(false);
}
catch (AuthenticationException)
{
@@ -73,15 +75,17 @@ public class QuickConnectController : BaseJellyfinApiController
/// <param name="secret">Secret previously returned from the Initiate endpoint.</param>
/// <response code="200">Quick connect result returned.</response>
/// <response code="404">Unknown quick connect secret.</response>
/// <response code="503">Quick connect state is unavailable.</response>
/// <returns>An updated <see cref="QuickConnectResult"/>.</returns>
[HttpGet("Connect")]
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
public ActionResult<QuickConnectResult> GetQuickConnectState([FromQuery, Required] string secret)
[ProducesResponseType(StatusCodes.Status503ServiceUnavailable)]
public async Task<ActionResult<QuickConnectResult>> GetQuickConnectState([FromQuery, Required] string secret)
{
try
{
return _quickConnect.CheckRequestStatus(secret);
return await _quickConnect.CheckRequestStatus(secret).ConfigureAwait(false);
}
catch (ResourceNotFoundException)
{
@@ -100,11 +104,15 @@ public class QuickConnectController : BaseJellyfinApiController
/// <param name="userId">The user the authorize. Access to the requested user is required.</param>
/// <response code="200">Quick connect result authorized successfully.</response>
/// <response code="403">Unknown user id.</response>
/// <response code="409">Request is already authorized, or authorizing it did not complete.</response>
/// <response code="503">Quick connect state is unavailable.</response>
/// <returns>Boolean indicating if the authorization was successful.</returns>
[HttpPost("Authorize")]
[Authorize]
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(StatusCodes.Status409Conflict)]
[ProducesResponseType(StatusCodes.Status503ServiceUnavailable)]
public async Task<ActionResult<bool>> AuthorizeQuickConnect([FromQuery, Required] string code, [FromQuery] Guid? userId = null)
{
userId = RequestHelpers.GetUserId(User, userId);
+6 -2
View File
@@ -241,15 +241,19 @@ public class UserController : BaseJellyfinApiController
/// <param name="request">The <see cref="QuickConnectDto"/> request.</param>
/// <response code="200">User authenticated.</response>
/// <response code="400">Missing token.</response>
/// <response code="404">Unknown or unauthorized quick connect secret.</response>
/// <response code="503">Quick connect state is unavailable.</response>
/// <returns>A <see cref="Task"/> containing an <see cref="AuthenticationRequest"/> with information about the new session.</returns>
[HttpPost("AuthenticateWithQuickConnect")]
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status503ServiceUnavailable)]
[Tags("Authentication")]
public ActionResult<AuthenticationResult> AuthenticateWithQuickConnect([FromBody, Required] QuickConnectDto request)
public async Task<ActionResult<AuthenticationResult>> AuthenticateWithQuickConnect([FromBody, Required] QuickConnectDto request)
{
try
{
return _quickConnectManager.GetAuthorizedRequest(request.Secret);
return await _quickConnectManager.GetAuthorizedRequest(request.Secret).ConfigureAwait(false);
}
catch (SecurityException e)
{
@@ -131,6 +131,8 @@ public class ExceptionMiddleware
FileNotFoundException => StatusCodes.Status404NotFound,
ResourceNotFoundException => StatusCodes.Status404NotFound,
MethodNotAllowedException => StatusCodes.Status405MethodNotAllowed,
ConflictException => StatusCodes.Status409Conflict,
ServiceUnavailableException => StatusCodes.Status503ServiceUnavailable,
_ => StatusCodes.Status500InternalServerError
};
}
+8
View File
@@ -112,6 +112,14 @@ namespace Jellyfin.Server
// instance. Active by default once a Redis connection is configured, no-op otherwise.
serviceCollection.AddScanLeaderLease(_startupConfig, Logger);
// Configuration invalidation bus: propagates shared-configuration and library-option writes
// to the other instances. Redis-backed when configured, no-op otherwise.
serviceCollection.AddConfigurationInvalidationBus(_startupConfig, Logger);
// Quick connect store: shares in-flight quick connect requests so the initiate, authorize and
// exchange legs can land on different instances. Redis-backed when configured, local otherwise.
serviceCollection.AddQuickConnectStore(_startupConfig, Logger);
foreach (var type in GetExportTypes<ILyricProvider>())
{
serviceCollection.AddSingleton(typeof(ILyricProvider), type);
@@ -0,0 +1,74 @@
using System;
using Emby.Server.Implementations.Configuration;
using MediaBrowser.Common.Configuration;
using MediaBrowser.Controller.MediaEncoding;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using StackExchange.Redis;
namespace Jellyfin.Server.Extensions;
/// <summary>
/// Extensions for registering the shared-configuration invalidation bus.
/// </summary>
public static class ConfigurationInvalidationServiceCollectionExtensions
{
/// <summary>
/// Registers the invalidation bus, Redis-backed when a connection string is configured and no-op
/// otherwise, and the subscriber applying the notices other instances publish.
/// </summary>
/// <remarks>
/// The connection string is only set for a multi-instance deployment, which is the only shape where
/// one instance can write the shared configuration directory behind another's back.
/// </remarks>
/// <param name="serviceCollection">The service collection.</param>
/// <param name="configuration">The configuration to read the Redis connection string from.</param>
/// <param name="logger">The logger to report the selected bus on.</param>
/// <returns>The updated service collection.</returns>
public static IServiceCollection AddConfigurationInvalidationBus(
this IServiceCollection serviceCollection,
IConfiguration configuration,
ILogger logger)
{
ArgumentNullException.ThrowIfNull(configuration);
ArgumentNullException.ThrowIfNull(logger);
if (string.IsNullOrEmpty(configuration[TranscodeStoreOptions.RedisConnectionStringKey]))
{
logger.LogInformation(
"Configuration invalidation bus: {Bus}. Shared-configuration and library-visibility changes stay local to the instance that made them; set {Key} to propagate them.",
nameof(NullConfigurationInvalidationBus),
TranscodeStoreOptions.RedisConnectionStringKey);
return serviceCollection.AddSingleton<IConfigurationInvalidationBus>(NullConfigurationInvalidationBus.Instance);
}
logger.LogInformation(
"Configuration invalidation bus: {Bus}. Shared-configuration and library-visibility changes propagate to every instance.",
nameof(RedisConfigurationInvalidationBus));
serviceCollection.AddSingleton<IConfigurationInvalidationBus>(sp =>
{
try
{
return new RedisConfigurationInvalidationBus(
sp.GetRequiredService<IConnectionMultiplexer>(),
sp.GetRequiredService<ILogger<RedisConfigurationInvalidationBus>>());
}
catch (Exception ex)
{
// Fail open: an unreachable Redis degrades to the single-instance behaviour of every
// instance keeping its own cached configuration, rather than aborting startup.
sp.GetRequiredService<ILogger<CoreAppHost>>().LogError(
ex,
"Redis is configured but unavailable, so shared-configuration changes will not propagate between instances. Check {Key}.",
TranscodeStoreOptions.RedisConnectionStringKey);
return NullConfigurationInvalidationBus.Instance;
}
});
return serviceCollection.AddHostedService<ConfigurationInvalidationSubscriber>();
}
}
@@ -0,0 +1,58 @@
using System;
using Emby.Server.Implementations.QuickConnect;
using MediaBrowser.Controller.MediaEncoding;
using MediaBrowser.Controller.QuickConnect;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using StackExchange.Redis;
namespace Jellyfin.Server.Extensions;
/// <summary>
/// Extensions for registering the quick connect store.
/// </summary>
public static class QuickConnectStoreServiceCollectionExtensions
{
/// <summary>
/// Registers the quick connect store, Redis-backed when a connection string is configured and
/// process-local otherwise, and reports the selected store at <see cref="LogLevel.Information"/>.
/// </summary>
/// <remarks>
/// The connection string is only set for a multi-instance deployment, which is the only shape where
/// the initiate, authorize and exchange legs of one flow can land on different instances. Set but
/// unreachable is a misconfigured deployment rather than a single-instance one, so the store is read
/// once during startup and an unreachable one stops the server coming up, rather than quietly handing
/// out a store the other instances cannot see.
/// </remarks>
/// <param name="serviceCollection">The service collection.</param>
/// <param name="configuration">The configuration to read the Redis connection string from.</param>
/// <param name="logger">The logger to report the selected store on.</param>
/// <returns>The updated service collection.</returns>
public static IServiceCollection AddQuickConnectStore(
this IServiceCollection serviceCollection,
IConfiguration configuration,
ILogger logger)
{
ArgumentNullException.ThrowIfNull(configuration);
ArgumentNullException.ThrowIfNull(logger);
if (string.IsNullOrEmpty(configuration[TranscodeStoreOptions.RedisConnectionStringKey]))
{
logger.LogInformation(
"Quick connect store: {Store}. A quick connect flow has to complete against one instance; set {Key} to share it.",
nameof(InMemoryQuickConnectStore),
TranscodeStoreOptions.RedisConnectionStringKey);
return serviceCollection.AddSingleton<IQuickConnectStore, InMemoryQuickConnectStore>();
}
logger.LogInformation(
"Quick connect store: {Store}. Quick connect flows complete across any instance.",
nameof(RedisQuickConnectStore));
return serviceCollection.AddSingleton<IQuickConnectStore>(sp => new RedisQuickConnectStore(
sp.GetRequiredService<IConnectionMultiplexer>(),
sp.GetRequiredService<ILogger<RedisQuickConnectStore>>()));
}
}
@@ -18,6 +18,13 @@ 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 <see cref="LogLevel.Information"/>.
/// </summary>
/// <remarks>
/// The <see cref="IConnectionMultiplexer"/> registered here is the one connection every Redis-backed
/// component shares, so configuring it makes valkey a hard startup dependency: quick connect reads it
/// during <c>InitializeServices</c> and stops the server when it cannot be reached. The softer
/// policies elsewhere - this store's connectivity report, the scan leader's fail-open - are
/// subordinate to that and only govern a store that goes away after that read has passed.
/// </remarks>
/// <param name="serviceCollection">The service collection.</param>
/// <param name="configuration">The configuration to read <c>Jellyfin:TranscodeStore</c> from.</param>
/// <param name="logger">The logger to report the selected store on.</param>
+14 -2
View File
@@ -54,6 +54,13 @@ namespace Jellyfin.Server
/// </summary>
public const string LoggingConfigFileSystem = "logging.json";
/// <summary>
/// How long a failed start keeps the setup server answering before the process exits, so an
/// install without a supervisor shows the failure instead of spinning. Skipped under an
/// orchestrator, where restarting is its job and holding only stretches the crash loop.
/// </summary>
private static readonly TimeSpan _setupServerHoldAfterFailedStart = TimeSpan.FromMinutes(10);
private static readonly SerilogLoggerFactory _loggerFactory = new SerilogLoggerFactory();
private static SetupServer? _setupServer;
private static CoreAppHost? _appHost;
@@ -255,13 +262,14 @@ namespace Jellyfin.Server
catch (Exception ex)
{
_restartOnShutdown = false;
Environment.ExitCode = 1;
_logger.LogCritical(ex, "Error while starting server");
if (_setupServer!.IsAlive && !configurationCompleted)
{
_setupServer!.SoftStop();
if (options.StartupMode is null or Configuration.StartupMode.MediaServer)
if ((options.StartupMode is null or Configuration.StartupMode.MediaServer) && !IsRunningInContainer())
{
await Task.Delay(TimeSpan.FromMinutes(10)).ConfigureAwait(false);
await Task.Delay(_setupServerHoldAfterFailedStart).ConfigureAwait(false);
}
await _setupServer!.StopAsync().ConfigureAwait(false);
@@ -285,6 +293,10 @@ namespace Jellyfin.Server
}
}
// Set by every official .NET base image, including the one this server ships in.
private static bool IsRunningInContainer()
=> string.Equals(Environment.GetEnvironmentVariable("DOTNET_RUNNING_IN_CONTAINER"), "true", StringComparison.OrdinalIgnoreCase);
/// <summary>
/// [Internal]Runs the startup Migrations.
/// </summary>
@@ -0,0 +1,26 @@
namespace MediaBrowser.Common.Configuration
{
/// <summary>
/// A notice that one instance has written shared configuration, so every other instance has to drop
/// its locally cached copy and read the shared file again.
/// </summary>
public sealed class ConfigurationInvalidation
{
/// <summary>
/// Gets or sets the cache this notice refers to.
/// </summary>
public ConfigurationInvalidationScope Scope { get; set; }
/// <summary>
/// Gets or sets what was invalidated within the scope: the configuration key for
/// <see cref="ConfigurationInvalidationScope.NamedConfiguration"/>, the library path for
/// <see cref="ConfigurationInvalidationScope.LibraryOptions"/>, and <c>null</c> otherwise.
/// </summary>
public string? Target { get; set; }
/// <summary>
/// Gets or sets the identity of the instance that published the notice, so it can ignore its own.
/// </summary>
public string? OriginId { get; set; }
}
}
@@ -0,0 +1,30 @@
namespace MediaBrowser.Common.Configuration
{
/// <summary>
/// Extensions for <see cref="IConfigurationInvalidationBus"/>.
/// </summary>
public static class ConfigurationInvalidationBusExtensions
{
/// <summary>
/// Announces a write this instance originated, and stays silent for a write induced by an
/// invalidation another instance published.
/// </summary>
/// <remarks>
/// This is the backstop for the write path: a consumer reacting to an applied invalidation by
/// writing - including a plugin that knows nothing about the bus - cannot turn that write into a
/// notice of its own.
/// </remarks>
/// <param name="bus">The bus.</param>
/// <param name="scope">The cache that was written.</param>
/// <param name="target">The configuration key or library path that was written, if any.</param>
public static void PublishLocalWrite(this IConfigurationInvalidationBus bus, ConfigurationInvalidationScope scope, string? target)
{
if (ConfigurationInvalidationContext.IsApplyingRemoteInvalidation)
{
return;
}
bus.Publish(scope, target);
}
}
}
@@ -0,0 +1,57 @@
using System;
using System.Threading;
namespace MediaBrowser.Common.Configuration
{
/// <summary>
/// Marks the flow of control that is applying an invalidation published by another instance, so the
/// reactions to it can tell a remote write apart from a local one.
/// </summary>
/// <remarks>
/// Applying an invalidation raises the same update events a local save raises, because the in-process
/// consumers have to re-read the configuration either way. Some of those consumers answer an update by
/// writing, and that write must neither repeat what the publishing instance already did nor fan back
/// out over the bus. The flag rides the execution context, so it reaches the queued and asynchronous
/// event handlers as well as the synchronous ones.
/// </remarks>
public static class ConfigurationInvalidationContext
{
private static readonly AsyncLocal<bool> _applyingRemoteInvalidation = new();
/// <summary>
/// Gets a value indicating whether the current flow of control is applying an invalidation
/// published by another instance rather than handling a local save.
/// </summary>
public static bool IsApplyingRemoteInvalidation => _applyingRemoteInvalidation.Value;
/// <summary>
/// Marks the current flow of control as applying a remote invalidation until the returned scope is
/// disposed.
/// </summary>
/// <returns>The scope to dispose once the invalidation has been applied.</returns>
public static IDisposable BeginApply() => new ApplyScope();
private sealed class ApplyScope : IDisposable
{
private readonly bool _previous;
private bool _disposed;
public ApplyScope()
{
_previous = _applyingRemoteInvalidation.Value;
_applyingRemoteInvalidation.Value = true;
}
public void Dispose()
{
if (_disposed)
{
return;
}
_disposed = true;
_applyingRemoteInvalidation.Value = _previous;
}
}
}
}
@@ -0,0 +1,29 @@
namespace MediaBrowser.Common.Configuration
{
/// <summary>
/// Identifies which locally cached copy of the shared configuration a
/// <see cref="ConfigurationInvalidation"/> refers to.
/// </summary>
public enum ConfigurationInvalidationScope
{
/// <summary>
/// The system configuration cached by the configuration manager.
/// </summary>
SystemConfiguration = 0,
/// <summary>
/// A single named configuration, identified by its key.
/// </summary>
NamedConfiguration = 1,
/// <summary>
/// The library options of a single collection folder, identified by its path.
/// </summary>
LibraryOptions = 2,
/// <summary>
/// The library options of every collection folder, for changes to the library structure itself.
/// </summary>
AllLibraryOptions = 3
}
}
@@ -0,0 +1,28 @@
using System;
namespace MediaBrowser.Common.Configuration
{
/// <summary>
/// Carries cache-invalidation notices between the instances that share one configuration directory.
/// </summary>
/// <remarks>
/// The shared directory carries the content; this bus only carries the fact that it changed. Every
/// implementation is expected to fail open: a bus that cannot deliver must not throw into the write
/// path, leaving each instance on its own locally cached copy until it restarts.
/// </remarks>
public interface IConfigurationInvalidationBus
{
/// <summary>
/// Announces that this instance has written shared configuration.
/// </summary>
/// <param name="scope">The cache that was written.</param>
/// <param name="target">The configuration key or library path that was written, if any.</param>
void Publish(ConfigurationInvalidationScope scope, string? target);
/// <summary>
/// Registers the handler invoked for notices published by other instances.
/// </summary>
/// <param name="handler">The handler applying the invalidation locally.</param>
void Subscribe(Action<ConfigurationInvalidation> handler);
}
}
@@ -85,6 +85,20 @@ namespace MediaBrowser.Common.Configuration
/// </summary>
/// <param name="factories">The factories.</param>
void AddParts(IEnumerable<IConfigurationFactory> factories);
/// <summary>
/// Drops the locally cached copy of configuration another instance has written to the shared
/// configuration directory, so the next read reloads it, and raises the local update event.
/// </summary>
/// <remarks>
/// An implementation predating the invalidation bus keeps the default, which reports that it
/// cannot drop its cache rather than quietly leaving it stale. The caller applying a remote
/// notice treats that as a failed apply and logs it.
/// </remarks>
/// <param name="key">The named configuration key, or <c>null</c> for the system configuration.</param>
/// <exception cref="NotSupportedException">The implementation cannot drop its cached configuration.</exception>
void InvalidateCachedConfiguration(string? key)
=> throw new NotSupportedException(GetType().Name + " cannot drop configuration cached from the shared configuration directory, so writes by other instances will not be picked up.");
}
public static class ConfigurationManagerExtensions
@@ -0,0 +1,27 @@
using System;
namespace MediaBrowser.Common.Configuration
{
/// <summary>
/// A no-op <see cref="IConfigurationInvalidationBus"/> used by single-instance installs and whenever
/// the shared bus is unavailable. Every instance keeps its own cached configuration, which is the
/// behaviour of an install that has only one.
/// </summary>
public sealed class NullConfigurationInvalidationBus : IConfigurationInvalidationBus
{
/// <summary>
/// Gets the shared instance.
/// </summary>
public static NullConfigurationInvalidationBus Instance { get; } = new NullConfigurationInvalidationBus();
/// <inheritdoc />
public void Publish(ConfigurationInvalidationScope scope, string? target)
{
}
/// <inheritdoc />
public void Subscribe(Action<ConfigurationInvalidation> handler)
{
}
}
}
@@ -0,0 +1,36 @@
using System;
namespace MediaBrowser.Common.Extensions
{
/// <summary>
/// Thrown when the current state of a resource does not allow the requested operation.
/// </summary>
public class ConflictException : Exception
{
/// <summary>
/// Initializes a new instance of the <see cref="ConflictException" /> class.
/// </summary>
public ConflictException()
{
}
/// <summary>
/// Initializes a new instance of the <see cref="ConflictException" /> class.
/// </summary>
/// <param name="message">The message.</param>
public ConflictException(string message)
: base(message)
{
}
/// <summary>
/// Initializes a new instance of the <see cref="ConflictException" /> class.
/// </summary>
/// <param name="message">The message.</param>
/// <param name="innerException">The exception that caused this one.</param>
public ConflictException(string message, Exception innerException)
: base(message, innerException)
{
}
}
}
@@ -0,0 +1,37 @@
using System;
namespace MediaBrowser.Common.Extensions
{
/// <summary>
/// Thrown when an operation cannot be answered because a backing service is unreachable, rather than
/// because the thing it was asked about does not exist.
/// </summary>
public class ServiceUnavailableException : Exception
{
/// <summary>
/// Initializes a new instance of the <see cref="ServiceUnavailableException" /> class.
/// </summary>
public ServiceUnavailableException()
{
}
/// <summary>
/// Initializes a new instance of the <see cref="ServiceUnavailableException" /> class.
/// </summary>
/// <param name="message">The message.</param>
public ServiceUnavailableException(string message)
: base(message)
{
}
/// <summary>
/// Initializes a new instance of the <see cref="ServiceUnavailableException" /> class.
/// </summary>
/// <param name="message">The message.</param>
/// <param name="innerException">The exception that caused this one.</param>
public ServiceUnavailableException(string message, Exception innerException)
: base(message, innerException)
{
}
}
}
@@ -14,6 +14,7 @@ using System.Threading.Tasks;
using Jellyfin.Data.Enums;
using Jellyfin.Database.Implementations.Entities;
using Jellyfin.Extensions.Json;
using MediaBrowser.Common.Configuration;
using MediaBrowser.Controller.IO;
using MediaBrowser.Controller.Library;
using MediaBrowser.Controller.Providers;
@@ -70,6 +71,12 @@ namespace MediaBrowser.Controller.Entities
public static IServerApplicationHost ApplicationHost { get; set; }
/// <summary>
/// Gets or sets the bus announcing library option writes to the other instances sharing this
/// configuration directory. Defaults to a no-op, which is the single-instance behaviour.
/// </summary>
public static IConfigurationInvalidationBus InvalidationBus { get; set; } = NullConfigurationInvalidationBus.Instance;
[JsonIgnore]
public override bool SupportsPlayedStatus => false;
@@ -188,11 +195,36 @@ namespace MediaBrowser.Controller.Entities
XmlSerializer.SerializeToFile(clone, GetLibraryOptionsPath(path));
LibraryOptionsUpdated?.Invoke(null, new LibraryOptionsUpdatedEventArgs(path, options));
InvalidationBus.PublishLocalWrite(ConfigurationInvalidationScope.LibraryOptions, path);
}
public static void OnCollectionFolderChange()
/// <summary>
/// Drops the cached options of one library so the next read comes off <c>options.xml</c> again.
/// Applied on the instances that did not write, and so does not publish.
/// </summary>
/// <param name="path">The library path.</param>
public static void InvalidateLibraryOptions(string path)
{
_libraryOptions.TryRemove(path, out _);
LibraryOptionsUpdated?.Invoke(null, new LibraryOptionsUpdatedEventArgs(path, GetLibraryOptions(path)));
}
/// <summary>
/// Drops every cached library option set. Applied on the instances that did not write, and so does
/// not publish.
/// </summary>
public static void InvalidateAllLibraryOptions()
=> _libraryOptions.Clear();
public static void OnCollectionFolderChange()
{
InvalidateAllLibraryOptions();
InvalidationBus.PublishLocalWrite(ConfigurationInvalidationScope.AllLibraryOptions, null);
}
public override bool IsSaveLocalMetadataEnabled()
{
return true;
@@ -449,19 +449,26 @@ namespace MediaBrowser.Controller.Entities
IUserDataManager userDataManager,
ILibraryManager libraryManager)
{
var filtered = items.Where(i => Filter(i, user, query, userDataManager, libraryManager));
var itemList = items as IReadOnlyList<BaseItem> ?? items.ToList();
// The user data checks below run per item, so read them all in one query up front.
var userDataBatch = user is not null && RequiresUserData(query)
? userDataManager.GetUserDataBatch(itemList, user)
: null;
var filtered = itemList.Where(i => Filter(i, user, query, userDataManager, libraryManager, userDataBatch));
if (query.IsPlayed.HasValue && user is not null)
{
var itemList = filtered.ToList();
var folderIds = itemList.OfType<Folder>().Select(f => f.Id).ToList();
var filteredList = filtered.ToList();
var folderIds = filteredList.OfType<Folder>().Select(f => f.Id).ToList();
if (folderIds.Count > 0)
{
var counts = libraryManager.GetPlayedAndTotalCountBatch(folderIds, user);
var isPlayedValue = query.IsPlayed.Value;
return itemList.Where(item =>
return filteredList.Where(item =>
{
if (item is Folder)
{
@@ -473,7 +480,7 @@ namespace MediaBrowser.Controller.Entities
});
}
return itemList;
return filteredList;
}
return filtered;
@@ -515,12 +522,29 @@ namespace MediaBrowser.Controller.Entities
itemsArray);
}
private static bool RequiresUserData(InternalItemsQuery query)
=> query.IsLiked.HasValue
|| query.IsFavoriteOrLiked.HasValue
|| query.IsFavorite.HasValue
|| query.IsResumable.HasValue
|| query.IsPlayed.HasValue;
private static UserItemData GetUserData(
IUserDataManager userDataManager,
User user,
BaseItem item,
Dictionary<Guid, UserItemData> userDataBatch)
=> userDataBatch is not null && userDataBatch.TryGetValue(item.Id, out var userData)
? userData
: userDataManager.GetUserData(user, item);
private static bool Filter(
BaseItem item,
User user,
InternalItemsQuery query,
IUserDataManager userDataManager,
ILibraryManager libraryManager)
ILibraryManager libraryManager,
Dictionary<Guid, UserItemData> userDataBatch)
{
if (!string.IsNullOrEmpty(query.NameStartsWith) && !item.SortName.StartsWith(query.NameStartsWith, StringComparison.InvariantCultureIgnoreCase))
{
@@ -568,7 +592,7 @@ namespace MediaBrowser.Controller.Entities
if (query.IsLiked.HasValue)
{
userData = userDataManager.GetUserData(user, item);
userData = GetUserData(userDataManager, user, item, userDataBatch);
if (!userData.Likes.HasValue || userData.Likes != query.IsLiked.Value)
{
return false;
@@ -577,7 +601,7 @@ namespace MediaBrowser.Controller.Entities
if (query.IsFavoriteOrLiked.HasValue)
{
userData ??= userDataManager.GetUserData(user, item);
userData ??= GetUserData(userDataManager, user, item, userDataBatch);
var isFavoriteOrLiked = userData.IsFavorite || (userData.Likes ?? false);
if (isFavoriteOrLiked != query.IsFavoriteOrLiked.Value)
@@ -588,7 +612,7 @@ namespace MediaBrowser.Controller.Entities
if (query.IsFavorite.HasValue)
{
userData ??= userDataManager.GetUserData(user, item);
userData ??= GetUserData(userDataManager, user, item, userDataBatch);
if (userData.IsFavorite != query.IsFavorite.Value)
{
return false;
@@ -597,7 +621,7 @@ namespace MediaBrowser.Controller.Entities
if (query.IsResumable.HasValue)
{
userData ??= userDataManager.GetUserData(user, item);
userData ??= GetUserData(userDataManager, user, item, userDataBatch);
var isResumable = userData.PlaybackPositionTicks > 0;
if (isResumable != query.IsResumable.Value)
@@ -612,7 +636,7 @@ namespace MediaBrowser.Controller.Entities
// Folders are batch-filtered by the collection Filter() overload.
if (!item.IsFolder)
{
userData ??= userDataManager.GetUserData(user, item);
userData ??= GetUserData(userDataManager, user, item, userDataBatch);
if (item.IsPlayed(user, userData) != query.IsPlayed.Value)
{
return false;
@@ -1,5 +1,6 @@
using System;
using System.Threading.Tasks;
using MediaBrowser.Common.Extensions;
using MediaBrowser.Controller.Authentication;
using MediaBrowser.Controller.Net;
using MediaBrowser.Model.QuickConnect;
@@ -21,17 +22,19 @@ namespace MediaBrowser.Controller.QuickConnect
/// </summary>
/// <param name="authorizationInfo">The initiator authorization info.</param>
/// <returns>A quick connect result with tokens to proceed or throws an exception if not active.</returns>
QuickConnectResult TryConnect(AuthorizationInfo authorizationInfo);
Task<QuickConnectResult> TryConnect(AuthorizationInfo authorizationInfo);
/// <summary>
/// Checks the status of an individual request.
/// </summary>
/// <param name="secret">Unique secret identifier of the request.</param>
/// <returns>Quick connect result.</returns>
QuickConnectResult CheckRequestStatus(string secret);
Task<QuickConnectResult> CheckRequestStatus(string secret);
/// <summary>
/// Authorizes a quick connect request to connect as the calling user.
/// Authorizes a quick connect request to connect as the calling user. A request can be authorized
/// once: a second attempt, including one following an attempt that failed part way, throws
/// <see cref="ConflictException"/> and the user has to start quick connect again for a new code.
/// </summary>
/// <param name="userId">User id.</param>
/// <param name="code">Identifying code for the request.</param>
@@ -39,10 +42,11 @@ namespace MediaBrowser.Controller.QuickConnect
Task<bool> AuthorizeRequest(Guid userId, string code);
/// <summary>
/// Gets the authorized request for the secret.
/// Gets the authorized request for the secret. The read does not consume the authorization, so a
/// client that retries the exchange gets the same access token until the authorization expires.
/// </summary>
/// <param name="secret">The secret.</param>
/// <returns>The authentication result.</returns>
AuthenticationResult GetAuthorizedRequest(string secret);
Task<AuthenticationResult> GetAuthorizedRequest(string secret);
}
}
@@ -0,0 +1,77 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using MediaBrowser.Common.Extensions;
using MediaBrowser.Controller.Authentication;
using MediaBrowser.Model.QuickConnect;
namespace MediaBrowser.Controller.QuickConnect;
/// <summary>
/// Holds the state of in-flight quick connect requests. The three legs of a quick connect flow -
/// initiate, authorize and exchange - can each land on a different instance, so the state has to be
/// reachable from all of them.
/// </summary>
/// <remarks>
/// A shared implementation that cannot reach its backend throws <see cref="ServiceUnavailableException"/>
/// rather than reporting a miss, because a miss tells a polling client its secret is invalid.
/// </remarks>
public interface IQuickConnectStore
{
/// <summary>
/// Looks up a pending request by the secret handed to the initiating client.
/// </summary>
/// <param name="secret">The request secret.</param>
/// <param name="cancellationToken">A cancellation token.</param>
/// <returns>The request, or <c>null</c> when it is unknown or has expired.</returns>
Task<QuickConnectResult?> GetRequestBySecretAsync(string secret, CancellationToken cancellationToken = default);
/// <summary>
/// Looks up a pending request by the code shown to the user.
/// </summary>
/// <param name="code">The user facing code.</param>
/// <param name="cancellationToken">A cancellation token.</param>
/// <returns>The request, or <c>null</c> when it is unknown or has expired.</returns>
Task<QuickConnectResult?> GetRequestByCodeAsync(string code, CancellationToken cancellationToken = default);
/// <summary>
/// Stores a new or updated request until <paramref name="expiresUtc"/>, resolvable by both its secret
/// and its code or by neither. A request already past <paramref name="expiresUtc"/> is not stored.
/// </summary>
/// <param name="request">The request to store.</param>
/// <param name="expiresUtc">The instant the request stops being resolvable.</param>
/// <param name="cancellationToken">A cancellation token.</param>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
Task SetRequestAsync(QuickConnectResult request, DateTime expiresUtc, CancellationToken cancellationToken = default);
/// <summary>
/// Atomically claims the sole right to authorize the request behind <paramref name="secret"/>, so
/// that two instances racing on one code cannot both mint an access token. The claim is never
/// released: a mint that failed after writing its token would otherwise be retried into a second one,
/// so a request whose authorization failed has to be started again.
/// </summary>
/// <param name="secret">The request secret.</param>
/// <param name="expiresUtc">The instant the claim lapses, after which the request can be authorized again.</param>
/// <param name="cancellationToken">A cancellation token.</param>
/// <returns><c>true</c> when this caller may go on to authorize the request; <c>false</c> when it is unknown, expired, already authorized or claimed elsewhere.</returns>
Task<bool> TryClaimAuthorizationAsync(string secret, DateTime expiresUtc, CancellationToken cancellationToken = default);
/// <summary>
/// Stores the authentication minted for an authorized request until <paramref name="expiresUtc"/>.
/// </summary>
/// <param name="secret">The request secret the client exchanges.</param>
/// <param name="authenticationResult">The authentication to hand out.</param>
/// <param name="expiresUtc">The instant the authentication stops being exchangeable.</param>
/// <param name="cancellationToken">A cancellation token.</param>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
Task SetAuthorizationAsync(string secret, AuthenticationResult authenticationResult, DateTime expiresUtc, CancellationToken cancellationToken = default);
/// <summary>
/// Reads the authentication for <paramref name="secret"/>. The read does not consume it, so a client
/// that retries an exchange gets the same access token for as long as the authentication lives.
/// </summary>
/// <param name="secret">The request secret.</param>
/// <param name="cancellationToken">A cancellation token.</param>
/// <returns>The authentication, or <c>null</c> when the secret is unknown or has expired.</returns>
Task<AuthenticationResult?> GetAuthorizationAsync(string secret, CancellationToken cancellationToken = default);
}
@@ -0,0 +1,113 @@
using System;
using System.Collections.Concurrent;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using MediaBrowser.Controller.Authentication;
using MediaBrowser.Model.QuickConnect;
namespace MediaBrowser.Controller.QuickConnect;
/// <summary>
/// A process-local <see cref="IQuickConnectStore"/>, and the single-instance default. Nothing it holds
/// is visible to another instance, so a deployment running more than one has to configure a shared
/// store instead.
/// </summary>
public sealed class InMemoryQuickConnectStore : IQuickConnectStore
{
private readonly ConcurrentDictionary<string, Entry<QuickConnectResult>> _requests = new(StringComparer.Ordinal);
private readonly ConcurrentDictionary<string, Entry<AuthenticationResult>> _authorizations = new(StringComparer.Ordinal);
private readonly ConcurrentDictionary<string, DateTime> _authorizationClaims = new(StringComparer.Ordinal);
/// <inheritdoc />
public Task<QuickConnectResult?> GetRequestBySecretAsync(string secret, CancellationToken cancellationToken = default)
{
Expire();
return Task.FromResult(_requests.TryGetValue(secret, out var entry) ? entry.Value : null);
}
/// <inheritdoc />
public Task<QuickConnectResult?> GetRequestByCodeAsync(string code, CancellationToken cancellationToken = default)
{
Expire();
return Task.FromResult(_requests.Values
.Select(entry => entry.Value)
.FirstOrDefault(request => string.Equals(request.Code, code, StringComparison.Ordinal)));
}
/// <inheritdoc />
public Task SetRequestAsync(QuickConnectResult request, DateTime expiresUtc, CancellationToken cancellationToken = default)
{
ArgumentNullException.ThrowIfNull(request);
Expire();
if (expiresUtc > DateTime.UtcNow)
{
_requests[request.Secret] = new Entry<QuickConnectResult>(expiresUtc, request);
}
return Task.CompletedTask;
}
/// <inheritdoc />
public Task<bool> TryClaimAuthorizationAsync(string secret, DateTime expiresUtc, CancellationToken cancellationToken = default)
{
Expire();
if (!_requests.TryGetValue(secret, out var entry) || entry.Value.Authenticated)
{
return Task.FromResult(false);
}
return Task.FromResult(_authorizationClaims.TryAdd(secret, expiresUtc));
}
/// <inheritdoc />
public Task SetAuthorizationAsync(string secret, AuthenticationResult authenticationResult, DateTime expiresUtc, CancellationToken cancellationToken = default)
{
Expire();
if (expiresUtc > DateTime.UtcNow)
{
_authorizations[secret] = new Entry<AuthenticationResult>(expiresUtc, authenticationResult);
}
return Task.CompletedTask;
}
/// <inheritdoc />
public Task<AuthenticationResult?> GetAuthorizationAsync(string secret, CancellationToken cancellationToken = default)
{
Expire();
return Task.FromResult(_authorizations.TryGetValue(secret, out var entry) ? entry.Value : null);
}
private void Expire()
{
var now = DateTime.UtcNow;
foreach (var (secret, entry) in _requests)
{
if (entry.ExpiresUtc <= now)
{
_requests.TryRemove(secret, out _);
}
}
foreach (var (secret, entry) in _authorizations)
{
if (entry.ExpiresUtc <= now)
{
_authorizations.TryRemove(secret, out _);
}
}
foreach (var (secret, expiresUtc) in _authorizationClaims)
{
if (expiresUtc <= now)
{
_authorizationClaims.TryRemove(secret, out _);
}
}
}
private sealed record Entry<T>(DateTime ExpiresUtc, T Value);
}
@@ -43,6 +43,14 @@ public sealed class ScanLeaderOptions
"TaskExtractMediaSegments",
"KeyframeExtraction",
"CleanupUserDataTask",
"OptimizeDatabaseTask"
"OptimizeDatabaseTask",
"DownloadLyrics",
"DownloadSubtitles",
"TmdbRefreshUpcomingEpisodes",
"RefreshTrickplayImages",
"MoveTrickplayImages",
"RefreshInternetChannels",
"RefreshGuide",
"PluginUpdates"
};
}
@@ -1,6 +1,9 @@
#nullable disable
using System;
using System.Collections.Generic;
using Jellyfin.Database.Implementations.Entities;
using MediaBrowser.Controller.Entities;
using MediaBrowser.Controller.Library;
namespace MediaBrowser.Controller.Sorting
@@ -27,5 +30,16 @@ namespace MediaBrowser.Controller.Sorting
/// </summary>
/// <value>The user data repository.</value>
IUserDataManager UserDataManager { get; set; }
/// <summary>
/// Gets or sets user data for the items being sorted, keyed by item id, read once up front.
/// A comparer that does not store it reads its user data one item at a time instead.
/// </summary>
/// <value>The prefetched user data, or <c>null</c> when none was prefetched.</value>
IReadOnlyDictionary<Guid, UserItemData> PrefetchedUserData
{
get => null;
set { }
}
}
}
@@ -0,0 +1,28 @@
#nullable disable
using MediaBrowser.Controller.Entities;
namespace MediaBrowser.Controller.Sorting
{
/// <summary>
/// Helpers shared by the comparers that sort on user data.
/// </summary>
public static class UserBaseItemComparerExtensions
{
/// <summary>
/// Gets the user data for an item, preferring the batch the sort prefetched.
/// </summary>
/// <param name="comparer">The comparer.</param>
/// <param name="item">The item.</param>
/// <returns>The item's user data.</returns>
public static UserItemData GetUserData(this IUserBaseItemComparer comparer, BaseItem item)
{
if (comparer.PrefetchedUserData is not null && comparer.PrefetchedUserData.TryGetValue(item.Id, out var userData))
{
return userData;
}
return comparer.UserDataManager.GetUserData(comparer.User, item);
}
}
}
+23
View File
@@ -83,6 +83,29 @@ Scan-leader gating is off: timer-driven library tasks run on every instance.
`Enabled=true` with no Redis connection string logs a warning, because gating cannot run.
### Valkey is a hard startup dependency
`Jellyfin:TranscodeStore:RedisConnectionString` selects one shared `IConnectionMultiplexer` and every
Redis-backed component hangs off it. Three of them disagree about an unreachable store on purpose, and
the order they run in is what makes that coherent:
| Component | On an unreachable store | When |
|---|---|---|
| Quick connect store (`ApplicationHost.ProbeQuickConnectStoreAsync`) | **Fails closed.** Reads a sentinel secret, retried for 30s, then logs `Critical` and stops startup | `InitializeServices`, before anything is served |
| `TranscodeStoreConnectivityProbe` | Logs `Error` and carries on | `IHostedService` start, after the gate |
| `RedisScanLeaderLease` | Fails open, treats itself as leader | Per scheduled-task tick, long after the gate |
The gate wins because it runs first: with a connection string set, an instance that reaches
`IHostedService` start has already proved the store reachable. The softer policies govern only a store
that goes away *afterwards*, where a running instance degrades rather than dying — quick connect calls
return `503`, transcode takeover stops, every instance scans.
The 30s window is there so a rollout survives valkey restarting alongside the server. Past it the
deployment is misconfigured or broken, the process exits non-zero and the orchestrator reports the real
cause. Upstream's 10-minute hold, which keeps the setup server answering after *any* failed start, is
skipped when `DOTNET_RUNNING_IN_CONTAINER` is set, because there the restart is the supervisor's job and
holding only stretches the crash loop.
### PostgreSQL provider
`src/Jellyfin.Database/Jellyfin.Database.Providers.PostgreSQL/` is an EF Core
@@ -212,7 +212,8 @@ namespace Jellyfin.LiveTv.Channels
if (query.IsFavorite.HasValue)
{
var val = query.IsFavorite.Value;
channels = channels.Where(i => _userDataManager.GetUserData(user, i).IsFavorite == val)
var userData = _userDataManager.GetUserDataBatch(channels, user);
channels = channels.Where(i => userData.TryGetValue(i.Id, out var data) && data.IsFavorite == val)
.ToList();
}
+18 -3
View File
@@ -304,8 +304,17 @@ namespace Jellyfin.LiveTv
if (query.IsAiring ?? false)
{
// Scoring reads the channel's user data per program, so read every channel's in one query.
var channels = programList
.Cast<LiveTvProgram>()
.Select(i => _libraryManager.GetItemById(i.ChannelId))
.OfType<BaseItem>()
.DistinctBy(i => i.Id)
.ToList();
var channelUserData = _userDataManager.GetUserDataBatch(channels, user);
orderedPrograms = orderedPrograms
.ThenByDescending(i => GetRecommendationScore(i, user, true));
.ThenByDescending(i => GetRecommendationScore(i, user, true, channelUserData));
}
IEnumerable<BaseItem> programs = orderedPrograms;
@@ -338,7 +347,11 @@ namespace Jellyfin.LiveTv
_dtoService.GetBaseItemDtos(internalResult.Items, options, query.User)));
}
private int GetRecommendationScore(LiveTvProgram program, User user, bool factorChannelWatchCount)
private int GetRecommendationScore(
LiveTvProgram program,
User user,
bool factorChannelWatchCount,
IReadOnlyDictionary<Guid, UserItemData> channelUserData)
{
var score = 0;
@@ -359,7 +372,9 @@ namespace Jellyfin.LiveTv
return score;
}
var channelUserdata = _userDataManager.GetUserData(user, channel);
var channelUserdata = channelUserData.TryGetValue(channel.Id, out var cached)
? cached
: _userDataManager.GetUserData(user, channel);
if (channelUserdata.Likes.HasValue)
{
@@ -444,6 +444,13 @@ public sealed class RecordingsManager : IRecordingsManager, IDisposable
private async void OnNamedConfigurationUpdated(object? sender, ConfigurationUpdateEventArgs e)
{
// The instance that wrote the change creates the folders; racing it from here would write the
// same virtual folder a second time.
if (ConfigurationInvalidationContext.IsApplyingRemoteInvalidation)
{
return;
}
if (string.Equals(e.Key, "livetv", StringComparison.OrdinalIgnoreCase))
{
await CreateRecordingFolders().ConfigureAwait(false);
@@ -18,6 +18,11 @@
<PackageReference Include="coverlet.collector" />
</ItemGroup>
<ItemGroup>
<!-- Linked, not project-referenced: Jellyfin.Server.Tests drags the whole server into this output. -->
<Compile Include="..\Jellyfin.Server.Tests\Migrations\PostgreSqlTestServer.cs" Link="Migrations\PostgreSqlTestServer.cs" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\src\Jellyfin.Database\Jellyfin.Database.Providers.PostgreSQL\Jellyfin.Database.Providers.PostgreSQL.csproj" />
<ProjectReference Include="..\..\src\Jellyfin.Database\Jellyfin.Database.Implementations\Jellyfin.Database.Implementations.csproj" />
@@ -1,71 +1,68 @@
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using DotNet.Testcontainers.Builders;
using Jellyfin.Database.Implementations;
using Jellyfin.Database.Implementations.DbConfiguration;
using Jellyfin.Database.Implementations.Entities;
using Jellyfin.Database.Implementations.Locking;
using Jellyfin.Database.Providers.PostgreSQL;
using Jellyfin.Server.Tests.Migrations;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging.Abstractions;
using Npgsql;
using Testcontainers.PostgreSql;
using Xunit;
namespace Jellyfin.Database.Tests.PostgreSQL;
/// <summary>
/// Integration tests that verify concurrent access patterns against a real PostgreSQL 16 container.
/// Integration tests that verify concurrent access patterns against a real PostgreSQL server.
/// </summary>
[Xunit.Trait("Category", "RequiresDocker")]
public sealed class PostgreSqlConcurrencyTests : IAsyncLifetime
{
private readonly PostgreSqlContainer _container;
private static int _databaseSequence;
private PostgreSqlTestServer? _server;
private NpgsqlDataSource? _dataSource;
private PostgreSqlDatabaseProvider? _provider;
/// <summary>
/// Initializes a new instance of the <see cref="PostgreSqlConcurrencyTests"/> class.
/// </summary>
public PostgreSqlConcurrencyTests()
{
_container = new PostgreSqlBuilder("postgres:16-alpine")
.WithWaitStrategy(Wait.ForUnixContainer().UntilCommandIsCompleted("pg_isready"))
.Build();
}
/// <summary>
/// Starts the PostgreSQL container and applies migrations before any tests in the class run.
/// Attaches to the test server, hands this test a database of its own and applies migrations to it.
/// </summary>
/// <returns>A <see cref="ValueTask"/> representing the asynchronous operation.</returns>
public async ValueTask InitializeAsync()
{
await _container.StartAsync().ConfigureAwait(false);
_server = await PostgreSqlTestServer.StartAsync().ConfigureAwait(false);
_dataSource = new NpgsqlDataSourceBuilder(_container.GetConnectionString()).Build();
var databaseName = FormattableString.Invariant($"pg_concurrency_{Interlocked.Increment(ref _databaseSequence)}");
var connectionString = await _server.CreateDatabaseAsync(databaseName, TestContext.Current.CancellationToken).ConfigureAwait(false);
_dataSource = new NpgsqlDataSourceBuilder(connectionString).Build();
_provider = new PostgreSqlDatabaseProvider(_dataSource);
// Apply migrations once for the whole test class.
var context = CreateContext();
await using (context.ConfigureAwait(false))
{
await context.Database.MigrateAsync().ConfigureAwait(false);
await context.Database.MigrateAsync(TestContext.Current.CancellationToken).ConfigureAwait(false);
}
}
/// <summary>
/// Stops and removes the PostgreSQL container after all tests in the class have run.
/// Releases the data source and the test server.
/// </summary>
/// <returns>A <see cref="ValueTask"/> representing the asynchronous operation.</returns>
public async ValueTask DisposeAsync()
{
// InitializeAsync can fail before the data source exists; its error must not be masked by an NRE here.
if (_dataSource is not null)
{
await _dataSource.DisposeAsync().ConfigureAwait(false);
}
await _container.DisposeAsync().ConfigureAwait(false);
if (_server is not null)
{
await _server.DisposeAsync().ConfigureAwait(false);
}
}
/// <summary>
@@ -1,62 +1,68 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using DotNet.Testcontainers.Builders;
using Jellyfin.Database.Implementations;
using Jellyfin.Database.Implementations.DbConfiguration;
using Jellyfin.Database.Implementations.Locking;
using Jellyfin.Database.Providers.PostgreSQL;
using Jellyfin.Server.Tests.Migrations;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging.Abstractions;
using Npgsql;
using Testcontainers.PostgreSql;
using Xunit;
namespace Jellyfin.Database.Tests.PostgreSQL;
/// <summary>
/// Integration tests that validate PostgreSQL migrations against a real container.
/// Integration tests that validate PostgreSQL migrations against a real server.
/// </summary>
[Xunit.Trait("Category", "RequiresDocker")]
public sealed class PostgreSqlMigrationTests : IAsyncLifetime
{
private readonly PostgreSqlContainer _container;
private static int _databaseSequence;
private PostgreSqlTestServer? _server;
private NpgsqlDataSource? _dataSource;
/// <summary>
/// Initializes a new instance of the <see cref="PostgreSqlMigrationTests"/> class.
/// </summary>
public PostgreSqlMigrationTests()
{
_container = new PostgreSqlBuilder("postgres:16-alpine")
.WithWaitStrategy(Wait.ForUnixContainer().UntilCommandIsCompleted("pg_isready"))
.Build();
}
/// <summary>
/// Starts the PostgreSQL container before any tests in the class run.
/// Attaches to the test server and hands this test an empty database of its own.
/// </summary>
/// <returns>A <see cref="ValueTask"/> representing the asynchronous operation.</returns>
public async ValueTask InitializeAsync()
{
await _container.StartAsync().ConfigureAwait(false);
_server = await PostgreSqlTestServer.StartAsync().ConfigureAwait(false);
var databaseName = FormattableString.Invariant($"pg_migration_{Interlocked.Increment(ref _databaseSequence)}");
var connectionString = await _server.CreateDatabaseAsync(databaseName, TestContext.Current.CancellationToken).ConfigureAwait(false);
_dataSource = new NpgsqlDataSourceBuilder(connectionString).Build();
}
/// <summary>
/// Stops and removes the PostgreSQL container after all tests in the class have run.
/// Releases the data source and the test server.
/// </summary>
/// <returns>A <see cref="ValueTask"/> representing the asynchronous operation.</returns>
public async ValueTask DisposeAsync()
{
await _container.DisposeAsync().ConfigureAwait(false);
// InitializeAsync can fail before the data source exists; its error must not be masked by an NRE here.
if (_dataSource is not null)
{
await _dataSource.DisposeAsync().ConfigureAwait(false);
}
if (_server is not null)
{
await _server.DisposeAsync().ConfigureAwait(false);
}
}
/// <summary>
/// Verifies that the <c>InitialPostgreSql</c> migration applies cleanly to a fresh PostgreSQL 16 container.
/// Verifies that the <c>InitialPostgreSql</c> migration applies cleanly to a fresh database.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task MigrateAsync_AppliesInitialMigrationCleanly()
{
await using var dataSource = new NpgsqlDataSourceBuilder(_container.GetConnectionString()).Build();
var context = CreateContext(dataSource);
var context = CreateContext(_dataSource!);
await using (context)
{
await context.Database.MigrateAsync(TestContext.Current.CancellationToken);
@@ -73,11 +79,7 @@ public sealed class PostgreSqlMigrationTests : IAsyncLifetime
[Fact]
public void CheckForUnappliedMigrations_PostgreSql()
{
// Use a dummy connection string; HasPendingModelChanges() is a purely in-memory check
// that compares the current compiled model with the migration snapshots — no real DB needed.
const string dummyConnectionString = "Host=localhost;Database=jellyfin;Username=postgres;Password=postgres";
using var dataSource = new NpgsqlDataSourceBuilder(dummyConnectionString).Build();
using var context = CreateContext(dataSource);
using var context = CreateContext(_dataSource!);
Assert.False(
context.Database.HasPendingModelChanges(),
@@ -2,71 +2,67 @@ using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using DotNet.Testcontainers.Builders;
using Jellyfin.Database.Implementations;
using Jellyfin.Database.Implementations.DbConfiguration;
using Jellyfin.Database.Implementations.Entities;
using Jellyfin.Database.Implementations.Locking;
using Jellyfin.Database.Providers.PostgreSQL;
using Jellyfin.Server.Tests.Migrations;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging.Abstractions;
using Npgsql;
using Testcontainers.PostgreSql;
using Xunit;
namespace Jellyfin.Database.Tests.PostgreSQL;
/// <summary>
/// Integration tests for CRUD operations, optimisation, and purge against a real PostgreSQL 16 container.
/// Integration tests for CRUD operations, optimisation, and purge against a real PostgreSQL server.
/// </summary>
[Xunit.Trait("Category", "RequiresDocker")]
public sealed class PostgreSqlProviderTests : IAsyncLifetime
{
private readonly PostgreSqlContainer _container;
private static int _databaseSequence;
private PostgreSqlTestServer? _server;
private NpgsqlDataSource? _dataSource;
private PostgreSqlDatabaseProvider? _provider;
/// <summary>
/// Initializes a new instance of the <see cref="PostgreSqlProviderTests"/> class.
/// </summary>
public PostgreSqlProviderTests()
{
_container = new PostgreSqlBuilder("postgres:16-alpine")
.WithWaitStrategy(Wait.ForUnixContainer().UntilCommandIsCompleted("pg_isready"))
.Build();
}
/// <summary>
/// Starts the PostgreSQL container and applies migrations before any tests in the class run.
/// Attaches to the test server, hands this test a database of its own and applies migrations to it.
/// </summary>
/// <returns>A <see cref="ValueTask"/> representing the asynchronous operation.</returns>
public async ValueTask InitializeAsync()
{
await _container.StartAsync().ConfigureAwait(false);
_server = await PostgreSqlTestServer.StartAsync().ConfigureAwait(false);
_dataSource = new NpgsqlDataSourceBuilder(_container.GetConnectionString()).Build();
var databaseName = FormattableString.Invariant($"pg_provider_{Interlocked.Increment(ref _databaseSequence)}");
var connectionString = await _server.CreateDatabaseAsync(databaseName, TestContext.Current.CancellationToken).ConfigureAwait(false);
_dataSource = new NpgsqlDataSourceBuilder(connectionString).Build();
_provider = new PostgreSqlDatabaseProvider(_dataSource);
// Apply migrations once for the whole test class.
var context = CreateContext();
await using (context.ConfigureAwait(false))
{
await context.Database.MigrateAsync().ConfigureAwait(false);
await context.Database.MigrateAsync(TestContext.Current.CancellationToken).ConfigureAwait(false);
}
}
/// <summary>
/// Stops and removes the PostgreSQL container after all tests in the class have run.
/// Releases the data source and the test server.
/// </summary>
/// <returns>A <see cref="ValueTask"/> representing the asynchronous operation.</returns>
public async ValueTask DisposeAsync()
{
// InitializeAsync can fail before the data source exists; its error must not be masked by an NRE here.
if (_dataSource is not null)
{
await _dataSource.DisposeAsync().ConfigureAwait(false);
}
await _container.DisposeAsync().ConfigureAwait(false);
if (_server is not null)
{
await _server.DisposeAsync().ConfigureAwait(false);
}
}
/// <summary>
@@ -155,11 +151,15 @@ public sealed class PostgreSqlProviderTests : IAsyncLifetime
var ctx = CreateContext();
await using (ctx)
{
var userId = Guid.NewGuid();
// DisplayPreferences.UserId is a foreign key onto Users, which PostgreSQL enforces and SQLite does not.
var user = new User("prefsuser", "Jellyfin.Server.Implementations.Users.DefaultAuthenticationProvider", "Jellyfin.Server.Implementations.Users.DefaultPasswordResetProvider");
ctx.Users.Add(user);
await ctx.SaveChangesAsync(TestContext.Current.CancellationToken);
var itemId = Guid.NewGuid();
// Create
var prefs = new DisplayPreferences(userId, itemId, "TestClient");
var prefs = new DisplayPreferences(user.Id, itemId, "TestClient");
ctx.DisplayPreferences.Add(prefs);
await ctx.SaveChangesAsync(TestContext.Current.CancellationToken);
@@ -278,7 +278,7 @@ public sealed class PostgreSqlProviderTests : IAsyncLifetime
// session_replication_role should be reset to 'origin' (default)
var role = await ctx.Database
.SqlQueryRaw<string>("SELECT current_setting('session_replication_role')")
.SqlQueryRaw<string>("SELECT current_setting('session_replication_role') AS \"Value\"")
.FirstAsync(TestContext.Current.CancellationToken);
Assert.Equal("origin", role);
}
@@ -0,0 +1,81 @@
using Emby.Server.Implementations;
using Xunit;
namespace Jellyfin.Server.Implementations.Tests.Configuration;
/// <summary>
/// The decision <c>ApplicationHost.OnConfigurationUpdated</c> makes about a port change. The ports this
/// process bound are fixed for its lifetime and the pending-restart flag is per-process, so an instance
/// applying another instance's port change still has to notice its own binding went stale - while leaving
/// the authorization write, and the notice that follows it, to the instance that made the change.
/// </summary>
public static class ApplicationHostPortChangeTests
{
/// <summary>
/// The local case, unchanged: clear the authorization flag and report the pending restart.
/// </summary>
[Fact]
public static void LocalPortChange_ClearsAuthorizationAndRequiresRestart()
{
var outcome = ApplicationHost.EvaluatePortChange(8096, 8920, 9096, 8920, true, false);
Assert.True(outcome.RequiresRestart);
Assert.True(outcome.ClearsPortAuthorization);
}
/// <summary>
/// The cross-instance case: the peer wrote the new port and cleared the flag with it, so this
/// instance must not write, but it is still listening on the old port and has to say so.
/// </summary>
[Theory]
[InlineData(true)]
[InlineData(false)]
public static void RemotePortChange_RequiresRestartWithoutWriting(bool isPortAuthorized)
{
var outcome = ApplicationHost.EvaluatePortChange(8096, 8920, 9096, 8920, isPortAuthorized, true);
Assert.True(outcome.RequiresRestart);
Assert.False(outcome.ClearsPortAuthorization);
}
/// <summary>
/// A second update while a port change is already pending must not write the flag again, and the
/// binding is still stale.
/// </summary>
[Fact]
public static void LocalPortChange_WithAuthorizationAlreadyCleared_RequiresRestartWithoutWriting()
{
var outcome = ApplicationHost.EvaluatePortChange(8096, 8920, 9096, 8920, false, false);
Assert.True(outcome.RequiresRestart);
Assert.False(outcome.ClearsPortAuthorization);
}
/// <summary>
/// An update that leaves the ports alone is not a port change, whoever wrote it.
/// </summary>
[Theory]
[InlineData(false)]
[InlineData(true)]
public static void UnchangedPorts_DoNothing(bool isApplyingRemoteInvalidation)
{
var outcome = ApplicationHost.EvaluatePortChange(8096, 8920, 8096, 8920, true, isApplyingRemoteInvalidation);
Assert.False(outcome.RequiresRestart);
Assert.False(outcome.ClearsPortAuthorization);
}
/// <summary>
/// Nothing is decided before the ports have been bound.
/// </summary>
[Theory]
[InlineData(0, 8920)]
[InlineData(8096, 0)]
public static void UnboundPorts_DoNothing(int boundHttpPort, int boundHttpsPort)
{
var outcome = ApplicationHost.EvaluatePortChange(boundHttpPort, boundHttpsPort, 9096, 9920, true, false);
Assert.False(outcome.RequiresRestart);
Assert.False(outcome.ClearsPortAuthorization);
}
}
@@ -0,0 +1,73 @@
using System;
using System.Collections.Generic;
using System.Linq;
using MediaBrowser.Common.Configuration;
namespace Jellyfin.Server.Implementations.Tests.Configuration;
/// <summary>
/// An in-process stand-in for the Redis pub/sub bus: every endpoint connected to one fabric receives
/// what the others publish, and never its own notices.
/// </summary>
internal sealed class FakeInvalidationBusFabric
{
private readonly List<Endpoint> _endpoints = new();
/// <summary>
/// Connects a new instance to the fabric.
/// </summary>
/// <param name="originId">The identity of the connecting instance.</param>
/// <returns>The bus of that instance.</returns>
public IConfigurationInvalidationBus Connect(string originId)
{
var endpoint = new Endpoint(this, originId);
lock (_endpoints)
{
_endpoints.Add(endpoint);
}
return endpoint;
}
private void Broadcast(ConfigurationInvalidation invalidation)
{
Endpoint[] endpoints;
lock (_endpoints)
{
endpoints = _endpoints.ToArray();
}
foreach (var endpoint in endpoints.Where(e => !string.Equals(e.OriginId, invalidation.OriginId, StringComparison.Ordinal)))
{
endpoint.Deliver(invalidation);
}
}
private sealed class Endpoint : IConfigurationInvalidationBus
{
private readonly FakeInvalidationBusFabric _fabric;
private readonly List<Action<ConfigurationInvalidation>> _handlers = new();
public Endpoint(FakeInvalidationBusFabric fabric, string originId)
{
_fabric = fabric;
OriginId = originId;
}
public string OriginId { get; }
public void Publish(ConfigurationInvalidationScope scope, string? target)
=> _fabric.Broadcast(new ConfigurationInvalidation { Scope = scope, Target = target, OriginId = OriginId });
public void Subscribe(Action<ConfigurationInvalidation> handler)
=> _handlers.Add(handler);
public void Deliver(ConfigurationInvalidation invalidation)
{
foreach (var handler in _handlers)
{
handler(invalidation);
}
}
}
}
@@ -0,0 +1,161 @@
using System;
using System.IO;
using System.Threading;
using System.Threading.Tasks;
using Emby.Server.Implementations.Configuration;
using Emby.Server.Implementations.Serialization;
using Jellyfin.Data;
using Jellyfin.Database.Implementations.Entities;
using Jellyfin.Database.Implementations.Enums;
using MediaBrowser.Common.Configuration;
using MediaBrowser.Controller;
using MediaBrowser.Controller.Entities;
using MediaBrowser.Model.Configuration;
using Microsoft.Extensions.Logging.Abstractions;
using Moq;
using Xunit;
namespace Jellyfin.Server.Implementations.Tests.Configuration;
/// <summary>
/// Library options are cached in a process-wide dictionary, so the replica that did not serve the admin's
/// request is the one under test here: the other replica's write reaches the shared library directory, and
/// this one has to stop answering out of its own stale copy.
/// </summary>
/// <remarks>
/// Disabling a library is an access revocation that overrides every per-user check, so a stale replica
/// keeps serving content that is supposed to be hidden from everyone.
/// </remarks>
public sealed class LibraryVisibilityPropagationTests : IDisposable
{
private readonly string _libraryPath;
private readonly MyXmlSerializer _serializer = new MyXmlSerializer();
/// <summary>
/// Initializes a new instance of the <see cref="LibraryVisibilityPropagationTests"/> class.
/// </summary>
public LibraryVisibilityPropagationTests()
{
_libraryPath = Path.Combine(Path.GetTempPath(), "jf-library-prop-" + Guid.NewGuid().ToString("N"));
Directory.CreateDirectory(_libraryPath);
var applicationHost = new Mock<IServerApplicationHost>();
applicationHost.Setup(host => host.ExpandVirtualPath(It.IsAny<string>())).Returns<string>(path => path);
applicationHost.Setup(host => host.ReverseVirtualPath(It.IsAny<string>())).Returns<string>(path => path);
CollectionFolder.XmlSerializer = _serializer;
CollectionFolder.ApplicationHost = applicationHost.Object;
}
/// <inheritdoc />
public void Dispose()
{
CollectionFolder.InvalidationBus = NullConfigurationInvalidationBus.Instance;
CollectionFolder.InvalidateAllLibraryOptions();
try
{
Directory.Delete(_libraryPath, true);
}
catch (IOException)
{
}
}
/// <summary>
/// Disabling a library on one replica has to hide it on every replica. Until it does, the ones that did
/// not serve the request keep the library visible to every user.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task LibraryDisabledOnAnotherInstance_IsNotVisibleHere()
{
var fabric = new FakeInvalidationBusFabric();
var otherInstance = fabric.Connect("pod-a");
await SubscribeThisInstanceAsync(fabric);
WriteSharedOptions(enabled: true);
var user = CreateUser();
var library = new CollectionFolder { Path = _libraryPath, Name = "Movies" };
// This replica answers out of its cache from here on.
Assert.True(library.IsVisible(user));
// The admin disables the library on the other replica: it writes the shared directory and says so.
WriteSharedOptions(enabled: false);
otherInstance.Publish(ConfigurationInvalidationScope.LibraryOptions, _libraryPath);
Assert.False(library.IsVisible(user));
Assert.False(CollectionFolder.GetLibraryOptions(_libraryPath).Enabled);
}
/// <summary>
/// A path remap made on another replica has to reach this one, or it keeps resolving media against a
/// path that is no longer the library's.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task LibraryPathRemappedOnAnotherInstance_IsSeenHere()
{
var fabric = new FakeInvalidationBusFabric();
var otherInstance = fabric.Connect("pod-a");
await SubscribeThisInstanceAsync(fabric);
WriteSharedOptions(enabled: true, mediaPath: "/media/old");
Assert.Equal("/media/old", CollectionFolder.GetLibraryOptions(_libraryPath).PathInfos[0].Path);
WriteSharedOptions(enabled: true, mediaPath: "/media/new");
otherInstance.Publish(ConfigurationInvalidationScope.LibraryOptions, _libraryPath);
Assert.Equal("/media/new", CollectionFolder.GetLibraryOptions(_libraryPath).PathInfos[0].Path);
}
/// <summary>
/// Saving library options here has to tell the other replicas, which is the half of the exchange the
/// tests above take as given.
/// </summary>
[Fact]
public void SaveLibraryOptions_AnnouncesTheLibraryToTheOtherInstances()
{
var fabric = new FakeInvalidationBusFabric();
ConfigurationInvalidation? received = null;
var otherInstance = fabric.Connect("pod-b");
otherInstance.Subscribe(invalidation => received = invalidation);
CollectionFolder.InvalidationBus = fabric.Connect("pod-a");
CollectionFolder.SaveLibraryOptions(_libraryPath, new LibraryOptions { Enabled = false });
Assert.NotNull(received);
Assert.Equal(ConfigurationInvalidationScope.LibraryOptions, received.Scope);
Assert.Equal(_libraryPath, received.Target);
}
private async Task SubscribeThisInstanceAsync(FakeInvalidationBusFabric fabric)
{
var bus = fabric.Connect("pod-b");
CollectionFolder.InvalidationBus = bus;
var subscriber = new ConfigurationInvalidationSubscriber(
bus,
Mock.Of<IConfigurationManager>(),
NullLogger<ConfigurationInvalidationSubscriber>.Instance);
await subscriber.StartAsync(CancellationToken.None);
}
private void WriteSharedOptions(bool enabled, string mediaPath = "/media")
{
// Written the way the other replica writes it, straight onto the shared directory.
var options = new LibraryOptions { Enabled = enabled, PathInfos = [new MediaPathInfo(mediaPath)] };
_serializer.SerializeToFile(options, Path.Combine(_libraryPath, "options.xml"));
}
private static User CreateUser()
{
var user = new User("propagation", "auth", "reset");
user.SetPermission(PermissionKind.EnableAllFolders, true);
return user;
}
}
@@ -0,0 +1,104 @@
using System;
using System.Threading.Tasks;
using Emby.Server.Implementations.Configuration;
using MediaBrowser.Common.Configuration;
using Microsoft.Extensions.Logging.Abstractions;
using StackExchange.Redis;
using Testcontainers.Redis;
using Xunit;
namespace Jellyfin.Server.Implementations.Tests.Configuration;
/// <summary>
/// Round-trips <see cref="RedisConfigurationInvalidationBus"/> through a real Redis, the transport two
/// replicas actually use to tell each other that the shared configuration directory has changed.
/// </summary>
[Trait("Category", "RequiresDocker")]
public sealed class RedisConfigurationInvalidationBusTests : IAsyncLifetime
{
private readonly RedisContainer _container;
private IConnectionMultiplexer? _redis;
/// <summary>
/// Initializes a new instance of the <see cref="RedisConfigurationInvalidationBusTests"/> class.
/// </summary>
public RedisConfigurationInvalidationBusTests()
{
_container = new RedisBuilder("redis:7-alpine").Build();
}
/// <inheritdoc />
public async ValueTask InitializeAsync()
{
await _container.StartAsync();
_redis = await ConnectionMultiplexer.ConnectAsync(_container.GetConnectionString());
}
/// <inheritdoc />
public async ValueTask DisposeAsync()
{
if (_redis is not null)
{
await _redis.DisposeAsync();
}
await _container.DisposeAsync();
}
/// <summary>
/// A notice published by one replica reaches the other, carrying enough to invalidate one entry.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Publish_ReachesTheOtherInstance()
{
var received = new TaskCompletionSource<ConfigurationInvalidation>();
var instanceA = CreateBus("pod-a");
var instanceB = CreateBus("pod-b");
instanceB.Subscribe(invalidation => received.TrySetResult(invalidation));
instanceA.Publish(ConfigurationInvalidationScope.LibraryOptions, "/media/movies");
var invalidation = await received.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken);
Assert.Equal(ConfigurationInvalidationScope.LibraryOptions, invalidation.Scope);
Assert.Equal("/media/movies", invalidation.Target);
Assert.Equal("pod-a", invalidation.OriginId);
}
/// <summary>
/// The publishing replica has already applied the change to its own cache, so it must not act on its
/// own notice and reload what it just wrote.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Publish_IsNotDeliveredToThePublisher()
{
var ownNotice = new TaskCompletionSource<ConfigurationInvalidation>();
var otherNotice = new TaskCompletionSource<ConfigurationInvalidation>();
var instanceA = CreateBus("pod-a");
var instanceB = CreateBus("pod-b");
instanceA.Subscribe(invalidation => ownNotice.TrySetResult(invalidation));
instanceB.Subscribe(invalidation => otherNotice.TrySetResult(invalidation));
instanceA.Publish(ConfigurationInvalidationScope.SystemConfiguration, null);
// Ordering is per channel, so B having the notice means A would have had it too.
await otherNotice.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken);
Assert.False(ownNotice.Task.IsCompleted);
}
private RedisConfigurationInvalidationBus CreateBus(string originId)
{
Environment.SetEnvironmentVariable("JELLYFIN_INSTANCE_ID", originId);
try
{
return new RedisConfigurationInvalidationBus(_redis!, NullLogger<RedisConfigurationInvalidationBus>.Instance);
}
finally
{
Environment.SetEnvironmentVariable("JELLYFIN_INSTANCE_ID", null);
}
}
}
@@ -0,0 +1,205 @@
using System;
using System.IO;
using System.Threading;
using System.Threading.Tasks;
using Emby.Server.Implementations;
using Emby.Server.Implementations.Configuration;
using Emby.Server.Implementations.Serialization;
using MediaBrowser.Common.Configuration;
using MediaBrowser.Common.Net;
using MediaBrowser.Controller.Configuration;
using Microsoft.Extensions.Logging.Abstractions;
using Xunit;
namespace Jellyfin.Server.Implementations.Tests.Configuration;
/// <summary>
/// Applying an invalidation re-raises the same update events a local save raises, and some consumers of
/// those events answer an update by writing - <c>RecordingsManager</c> creating the recording folders for
/// the <c>livetv</c> key, <c>ApplicationHost</c> clearing <c>IsPortAuthorized</c> on a port change. The
/// instance that did not write must not repeat those writes, and no write it is induced into must reach
/// the bus.
/// </summary>
public sealed class RemoteInvalidationApplyTests : IDisposable
{
private readonly string _root;
/// <summary>
/// Initializes a new instance of the <see cref="RemoteInvalidationApplyTests"/> class.
/// </summary>
public RemoteInvalidationApplyTests()
{
_root = Path.Combine(Path.GetTempPath(), "jf-config-apply-" + Guid.NewGuid().ToString("N"));
Directory.CreateDirectory(_root);
}
/// <inheritdoc />
public void Dispose()
{
try
{
Directory.Delete(_root, true);
}
catch (IOException)
{
}
}
/// <summary>
/// A consumer that knows nothing about the bus - a plugin, or anything reached transitively from one -
/// can answer an applied invalidation by writing. That write must not become a notice of its own.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task NamedInvalidation_InducingAWriteOnTheReceiver_DoesNotPublishBack()
{
var fabric = new FakeInvalidationBusFabric();
var instanceA = CreateInstance(fabric, "pod-a");
var instanceB = await CreateSubscribedInstanceAsync(fabric, "pod-b");
var receivedByA = 0;
instanceA.InvalidationBus.Subscribe(_ => Interlocked.Increment(ref receivedByA));
var writesByB = 0;
instanceB.NamedConfigurationUpdated += (_, e) =>
{
// One shot: the induced write raises the event again on this instance.
if (Interlocked.Increment(ref writesByB) > 1)
{
return;
}
var configuration = (NetworkConfiguration)instanceB.GetConfiguration(e.Key);
configuration.PublishedServerUriBySubnet = ["10.0.0.0/8=example"];
instanceB.SaveConfiguration(e.Key, configuration);
};
var updated = instanceA.GetNetworkConfiguration();
updated.EnableRemoteAccess = false;
instanceA.SaveConfiguration(NetworkConfigurationStore.StoreKey, updated);
Assert.Equal(0, receivedByA);
// The point of the bus still holds: B is not left on its stale copy.
Assert.False(instanceB.GetNetworkConfiguration().EnableRemoteAccess);
}
/// <summary>
/// The shape of <c>RecordingsManager</c>: a consumer that answers a named configuration update by
/// writing has to be able to tell that the write was another instance's, and skip it.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task NamedInvalidation_WithAWriteTriggeringConsumer_DoesNotDuplicateTheWrite()
{
var fabric = new FakeInvalidationBusFabric();
var instanceA = CreateInstance(fabric, "pod-a");
var instanceB = await CreateSubscribedInstanceAsync(fabric, "pod-b");
var receivedByA = 0;
instanceA.InvalidationBus.Subscribe(_ => Interlocked.Increment(ref receivedByA));
var writesByB = 0;
instanceB.NamedConfigurationUpdated += (_, e) =>
{
if (ConfigurationInvalidationContext.IsApplyingRemoteInvalidation)
{
return;
}
Interlocked.Increment(ref writesByB);
};
var updated = instanceA.GetNetworkConfiguration();
updated.EnableRemoteAccess = false;
instanceA.SaveConfiguration(NetworkConfigurationStore.StoreKey, updated);
Assert.Equal(0, writesByB);
Assert.Equal(0, receivedByA);
Assert.False(instanceB.GetNetworkConfiguration().EnableRemoteAccess);
// A read-only consumer is still told, which is what the invalidation exists for.
var refreshes = 0;
instanceB.NamedConfigurationUpdated += (_, _) => Interlocked.Increment(ref refreshes);
updated.EnableRemoteAccess = true;
instanceA.SaveConfiguration(NetworkConfigurationStore.StoreKey, updated);
Assert.Equal(1, refreshes);
}
/// <summary>
/// The <c>ApplicationHost.IsPortAuthorized</c> class of consumer: the system configuration event is
/// queued rather than raised inline, so the fix has to survive the hop onto the thread pool.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task SystemInvalidation_InducingAQueuedWriteOnTheReceiver_DoesNotPublishBack()
{
var fabric = new FakeInvalidationBusFabric();
var instanceA = CreateInstance(fabric, "pod-a");
var instanceB = await CreateSubscribedInstanceAsync(fabric, "pod-b");
var receivedByA = 0;
instanceA.InvalidationBus.Subscribe(_ => Interlocked.Increment(ref receivedByA));
var applied = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
var handled = 0;
instanceB.ConfigurationUpdated += (_, _) =>
{
if (Interlocked.Increment(ref handled) > 1)
{
return;
}
instanceB.Configuration.IsPortAuthorized = false;
instanceB.SaveConfiguration();
applied.TrySetResult();
};
instanceA.Configuration.QuickConnectAvailable = false;
instanceA.SaveConfiguration();
await applied.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken);
Assert.Equal(0, receivedByA);
Assert.False(instanceB.Configuration.QuickConnectAvailable);
}
private async Task<ServerConfigurationManager> CreateSubscribedInstanceAsync(FakeInvalidationBusFabric fabric, string originId)
{
var instance = CreateInstance(fabric, originId);
var subscriber = new ConfigurationInvalidationSubscriber(
instance.InvalidationBus,
instance,
NullLogger<ConfigurationInvalidationSubscriber>.Instance);
await subscriber.StartAsync(CancellationToken.None);
return instance;
}
private ServerConfigurationManager CreateInstance(FakeInvalidationBusFabric fabric, string originId)
{
var paths = new ServerApplicationPaths(
Ensure("data"),
Ensure("log"),
Ensure("config"),
Ensure("cache"),
Ensure("web"));
var manager = new ServerConfigurationManager(paths, NullLoggerFactory.Instance, new MyXmlSerializer())
{
InvalidationBus = fabric.Connect(originId)
};
manager.AddParts([new NetworkConfigurationFactory()]);
return manager;
}
private string Ensure(string name)
{
var path = Path.Combine(_root, name);
Directory.CreateDirectory(path);
return path;
}
}
@@ -0,0 +1,147 @@
using System;
using System.IO;
using System.Threading;
using System.Threading.Tasks;
using Emby.Server.Implementations;
using Emby.Server.Implementations.Configuration;
using Emby.Server.Implementations.Serialization;
using MediaBrowser.Common.Configuration;
using MediaBrowser.Common.Net;
using MediaBrowser.Controller.Configuration;
using Microsoft.Extensions.Logging.Abstractions;
using StackExchange.Redis;
using Xunit;
namespace Jellyfin.Server.Implementations.Tests.Configuration;
/// <summary>
/// Two independently constructed <see cref="ServerConfigurationManager"/> instances over one configuration
/// directory are the in-process stand-in for two replicas sharing one <c>/config</c> mount: what either of
/// them writes, the other has to pick up without being restarted.
/// </summary>
public sealed class SharedConfigurationPropagationTests : IDisposable
{
private readonly string _root;
/// <summary>
/// Initializes a new instance of the <see cref="SharedConfigurationPropagationTests"/> class.
/// </summary>
public SharedConfigurationPropagationTests()
{
_root = Path.Combine(Path.GetTempPath(), "jf-config-prop-" + Guid.NewGuid().ToString("N"));
Directory.CreateDirectory(_root);
}
/// <inheritdoc />
public void Dispose()
{
try
{
Directory.Delete(_root, true);
}
catch (IOException)
{
}
}
/// <summary>
/// A system configuration setting tightened on one replica has to hold on every other replica, not
/// only on the one that served the admin's request.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task SystemConfigurationSavedOnOneInstance_IsSeenByAnotherWithoutRestart()
{
var fabric = new FakeInvalidationBusFabric();
var instanceA = CreateInstance(fabric, "pod-a");
var instanceB = await CreateSubscribedInstanceAsync(fabric, "pod-b");
// B has the pre-change configuration in hand before A writes, as a running replica would.
Assert.True(instanceB.Configuration.QuickConnectAvailable);
instanceA.Configuration.QuickConnectAvailable = false;
instanceA.SaveConfiguration();
Assert.False(instanceB.Configuration.QuickConnectAvailable);
}
/// <summary>
/// The same has to hold for the named configurations, which are cached per key and never reloaded.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task NamedConfigurationSavedOnOneInstance_IsSeenByAnotherWithoutRestart()
{
var fabric = new FakeInvalidationBusFabric();
var instanceA = CreateInstance(fabric, "pod-a");
var instanceB = await CreateSubscribedInstanceAsync(fabric, "pod-b");
Assert.True(instanceB.GetNetworkConfiguration().EnableRemoteAccess);
var updated = instanceA.GetNetworkConfiguration();
updated.EnableRemoteAccess = false;
instanceA.SaveConfiguration(NetworkConfigurationStore.StoreKey, updated);
Assert.False(instanceB.GetNetworkConfiguration().EnableRemoteAccess);
}
/// <summary>
/// A replica that cannot reach the bus keeps serving: the admin's save still lands on the shared
/// directory, and the only loss is that the other replicas are not told about it.
/// </summary>
[Fact]
public void SaveConfiguration_WithUnreachableBus_DoesNotThrow()
{
using var multiplexer = ConnectionMultiplexer.Connect("127.0.0.1:1,abortConnect=false,connectTimeout=200,connectRetry=1,syncTimeout=200");
var bus = new RedisConfigurationInvalidationBus(multiplexer, NullLogger<RedisConfigurationInvalidationBus>.Instance);
bus.Subscribe(_ => throw new InvalidOperationException("Nothing can be delivered by an unreachable bus."));
var instance = CreateInstance(new FakeInvalidationBusFabric(), "pod-a");
instance.InvalidationBus = bus;
instance.Configuration.QuickConnectAvailable = false;
instance.SaveConfiguration();
instance.SaveConfiguration(NetworkConfigurationStore.StoreKey, instance.GetNetworkConfiguration());
Assert.False(instance.Configuration.QuickConnectAvailable);
}
private async Task<ServerConfigurationManager> CreateSubscribedInstanceAsync(FakeInvalidationBusFabric fabric, string originId)
{
var instance = CreateInstance(fabric, originId);
var subscriber = new ConfigurationInvalidationSubscriber(
instance.InvalidationBus,
instance,
NullLogger<ConfigurationInvalidationSubscriber>.Instance);
await subscriber.StartAsync(CancellationToken.None);
return instance;
}
private ServerConfigurationManager CreateInstance(FakeInvalidationBusFabric fabric, string originId)
{
// Every instance has its own paths object, all of them pointing at the one shared directory.
var paths = new ServerApplicationPaths(
Ensure("data"),
Ensure("log"),
Ensure("config"),
Ensure("cache"),
Ensure("web"));
var manager = new ServerConfigurationManager(paths, NullLoggerFactory.Instance, new MyXmlSerializer())
{
InvalidationBus = fabric.Connect(originId)
};
manager.AddParts([new NetworkConfigurationFactory()]);
return manager;
}
private string Ensure(string name)
{
var path = Path.Combine(_root, name);
Directory.CreateDirectory(path);
return path;
}
}
@@ -31,6 +31,8 @@
<ItemGroup>
<ProjectReference Include="..\..\Emby.Server.Implementations\Emby.Server.Implementations.csproj" />
<ProjectReference Include="..\..\Jellyfin.Server.Implementations\Jellyfin.Server.Implementations.csproj" />
<ProjectReference Include="..\..\MediaBrowser.Providers\MediaBrowser.Providers.csproj" />
<ProjectReference Include="..\..\src\Jellyfin.LiveTv\Jellyfin.LiveTv.csproj" />
<ProjectReference Include="..\Jellyfin.Server.Integration.Tests\Jellyfin.Server.Integration.Tests.csproj" />
<ProjectReference Include="..\..\src\Jellyfin.Database\Jellyfin.Database.Implementations\Jellyfin.Database.Implementations.csproj" />
</ItemGroup>
@@ -7,6 +7,7 @@ using Emby.Naming.Common;
using Emby.Server.Implementations.Library;
using Emby.Server.Implementations.Sorting;
using Jellyfin.Data.Enums;
using Jellyfin.Database.Implementations.Entities;
using Jellyfin.Database.Implementations.Enums;
using MediaBrowser.Controller.Configuration;
using MediaBrowser.Controller.Entities;
@@ -63,13 +64,49 @@ public class LibraryManagerSortTests
Assert.Equal(new[] { "Alpha", "Mike", "Zulu" }, sorted.Select(i => i.Name));
}
[Fact]
public void Sort_ComparerThatIgnoresPrefetchedUserData_StillSortsFromLiveReads()
{
var alpha = new Audio { Name = "Alpha", SortName = "Alpha", Id = Guid.NewGuid() };
var zulu = new Audio { Name = "Zulu", SortName = "Zulu", Id = Guid.NewGuid() };
var playCounts = new Dictionary<Guid, int> { [alpha.Id] = 1, [zulu.Id] = 9 };
var userDataManager = new Mock<IUserDataManager>();
userDataManager
.Setup(u => u.GetUserData(It.IsAny<User>(), It.IsAny<BaseItem>()))
.Returns<User, BaseItem>((_, item) => new UserItemData { Key = item.Id.ToString("N"), PlayCount = playCounts[item.Id] });
userDataManager
.Setup(u => u.GetUserDataBatch(It.IsAny<IReadOnlyList<BaseItem>>(), It.IsAny<User>()))
.Returns(new Dictionary<Guid, UserItemData>());
var libraryManager = CreateLibraryManager(
new IBaseItemComparer[] { new PluginPlayCountComparer() },
userDataManager);
var sorted = libraryManager.Sort(
new BaseItem[] { alpha, zulu },
new User("sorter", "provider", "provider"),
new[] { (ItemSortBy.PlayCount, SortOrder.Descending) }).ToArray();
Assert.Equal(new[] { "Zulu", "Alpha" }, sorted.Select(i => i.Name));
userDataManager.Verify(u => u.GetUserData(It.IsAny<User>(), It.IsAny<BaseItem>()), Times.AtLeastOnce);
}
private static Folder MakeFolder(string name, DateTime dateLastMediaAdded)
=> new() { Name = name, Id = Guid.NewGuid(), DateLastMediaAdded = dateLastMediaAdded };
private static Emby.Server.Implementations.Library.LibraryManager CreateLibraryManager(IReadOnlyCollection<IBaseItemComparer> comparers)
private static Emby.Server.Implementations.Library.LibraryManager CreateLibraryManager(
IReadOnlyCollection<IBaseItemComparer> comparers,
Mock<IUserDataManager>? userDataManager = null)
{
var fixture = new Fixture().Customize(new AutoMoqCustomization());
fixture.Register(() => new NamingOptions());
if (userDataManager is not null)
{
fixture.Inject(userDataManager.Object);
}
var configMock = fixture.Freeze<Mock<IServerConfigurationManager>>();
configMock.Setup(c => c.ApplicationPaths.ProgramDataPath).Returns("/data");
BaseItem.ConfigurationManager ??= configMock.Object;
@@ -86,4 +123,22 @@ public class LibraryManagerSortTests
fixture.Create<IEnumerable<ILibraryPostScanTask>>()))
.Create();
}
/// <summary>
/// A comparer of the shape a third-party plugin ships: it implements
/// <see cref="IUserBaseItemComparer"/> without ever mentioning PrefetchedUserData.
/// </summary>
public sealed class PluginPlayCountComparer : IUserBaseItemComparer
{
public User User { get; set; } = null!;
public IUserManager UserManager { get; set; } = null!;
public IUserDataManager UserDataManager { get; set; } = null!;
public ItemSortBy Type => ItemSortBy.PlayCount;
public int Compare(BaseItem? x, BaseItem? y)
=> UserDataManager.GetUserData(User, x!)!.PlayCount.CompareTo(UserDataManager.GetUserData(User, y!)!.PlayCount);
}
}
@@ -1,5 +1,4 @@
using System;
using System.Collections.Generic;
using Emby.Server.Implementations.Library;
using Jellyfin.Database.Implementations;
using Jellyfin.Database.Implementations.Entities;
@@ -49,6 +48,12 @@ public sealed class UserDataManagerTests : IDisposable
{
Id = Guid.NewGuid()
};
using (var ctx = CreateDbContext())
{
ctx.Users.Add(_user);
ctx.SaveChanges();
}
}
public void Dispose()
@@ -78,6 +83,23 @@ public sealed class UserDataManagerTests : IDisposable
};
}
private void Seed(AudioBook item, params UserData[] rows)
{
using var ctx = CreateDbContext();
ctx.BaseItems.Add(new BaseItemEntity { Id = item.Id, Type = typeof(AudioBook).FullName! });
ctx.UserData.AddRange(rows);
ctx.SaveChanges();
}
private User CreateOtherUser()
{
var user = new User("other", "auth-provider", "reset-provider") { Id = Guid.NewGuid() };
using var ctx = CreateDbContext();
ctx.Users.Add(user);
ctx.SaveChanges();
return user;
}
private UserData CreateUserDataRow(AudioBook item, string key, long positionTicks)
{
return new UserData
@@ -98,11 +120,10 @@ public sealed class UserDataManagerTests : IDisposable
var currentKey = item.GetUserDataKeys()[0];
// the retired-key row comes first to ensure selection is by key, not row order
item.UserData = new List<UserData>
{
Seed(
item,
CreateUserDataRow(item, "Author-Old Album-0001Old File Name", 111),
CreateUserDataRow(item, currentKey, 222)
};
CreateUserDataRow(item, currentKey, 222));
var userData = _userDataManager.GetUserData(_user, item);
@@ -117,11 +138,10 @@ public sealed class UserDataManagerTests : IDisposable
var item = CreateAudioBook();
var idKey = item.GetUserDataKeys()[1];
item.UserData = new List<UserData>
{
Seed(
item,
CreateUserDataRow(item, "Author-Old Album-0001Old File Name", 111),
CreateUserDataRow(item, idKey, 333)
};
CreateUserDataRow(item, idKey, 333));
var userData = _userDataManager.GetUserData(_user, item);
@@ -135,10 +155,7 @@ public sealed class UserDataManagerTests : IDisposable
{
var item = CreateAudioBook();
item.UserData = new List<UserData>
{
CreateUserDataRow(item, "Author-Old Album-0001Old File Name", 111)
};
Seed(item, CreateUserDataRow(item, "Author-Old Album-0001Old File Name", 111));
var userData = _userDataManager.GetUserData(_user, item);
@@ -150,7 +167,7 @@ public sealed class UserDataManagerTests : IDisposable
public void GetUserData_NoRows_ReturnsDefaultWithPrimaryKey()
{
var item = CreateAudioBook();
item.UserData = new List<UserData>();
Seed(item);
var userData = _userDataManager.GetUserData(_user, item);
@@ -166,13 +183,9 @@ public sealed class UserDataManagerTests : IDisposable
var currentKey = item.GetUserDataKeys()[0];
var otherUserRow = CreateUserDataRow(item, currentKey, 999);
otherUserRow.UserId = Guid.NewGuid();
otherUserRow.UserId = CreateOtherUser().Id;
item.UserData = new List<UserData>
{
otherUserRow,
CreateUserDataRow(item, currentKey, 222)
};
Seed(item, otherUserRow, CreateUserDataRow(item, currentKey, 222));
var userData = _userDataManager.GetUserData(_user, item);
@@ -183,23 +196,15 @@ public sealed class UserDataManagerTests : IDisposable
[Fact]
public void GetUserDataBatch_DatabaseFallback_ResolvesRowsByKeyOrder()
{
// no preloaded navigation data, so the batch takes the database fallback
var fossilItem = CreateAudioBook();
var retiredItem = CreateAudioBook();
using (var ctx = CreateDbContext())
{
ctx.Users.Add(_user);
ctx.BaseItems.Add(new BaseItemEntity { Id = fossilItem.Id, Type = typeof(AudioBook).FullName! });
ctx.BaseItems.Add(new BaseItemEntity { Id = retiredItem.Id, Type = typeof(AudioBook).FullName! });
// the stale id-key row is inserted first so selection by row order would return it
ctx.UserData.AddRange(
CreateUserDataRow(fossilItem, fossilItem.GetUserDataKeys()[1], 111),
CreateUserDataRow(fossilItem, fossilItem.GetUserDataKeys()[0], 222),
CreateUserDataRow(retiredItem, "Author-Old Album-0001Old File Name", 333));
ctx.SaveChanges();
}
// the stale id-key row is inserted first so selection by row order would return it
Seed(
fossilItem,
CreateUserDataRow(fossilItem, fossilItem.GetUserDataKeys()[1], 111),
CreateUserDataRow(fossilItem, fossilItem.GetUserDataKeys()[0], 222));
Seed(retiredItem, CreateUserDataRow(retiredItem, "Author-Old Album-0001Old File Name", 333));
var result = _userDataManager.GetUserDataBatch([fossilItem, retiredItem], _user);
@@ -8,6 +8,7 @@ using MediaBrowser.Common.Extensions;
using MediaBrowser.Controller.Authentication;
using MediaBrowser.Controller.Configuration;
using MediaBrowser.Controller.Net;
using MediaBrowser.Controller.QuickConnect;
using MediaBrowser.Model.Configuration;
using Moq;
using Xunit;
@@ -40,6 +41,8 @@ namespace Jellyfin.Server.Implementations.Tests.QuickConnect
ConfigureMembers = true
}).Inject(configManager.Object);
_fixture.Inject<IQuickConnectStore>(new InMemoryQuickConnectStore());
// User object contains circular references.
_fixture.Behaviors.OfType<ThrowingRecursionBehavior>().ToList()
.ForEach(b => _fixture.Behaviors.Remove(b));
@@ -60,8 +63,8 @@ namespace Jellyfin.Server.Implementations.Tests.QuickConnect
[InlineData("Device", "", "Client", "1.0.0")]
[InlineData("Device", "DeviceId", "", "1.0.0")]
[InlineData("Device", "DeviceId", "Client", "")]
public void TryConnect_InvalidAuthorizationInfo_ThrowsArgumentException(string device, string deviceId, string client, string version)
=> Assert.Throws<ArgumentException>(() => _quickConnectManager.TryConnect(
public async Task TryConnect_InvalidAuthorizationInfo_ThrowsArgumentException(string device, string deviceId, string client, string version)
=> await Assert.ThrowsAsync<ArgumentException>(() => _quickConnectManager.TryConnect(
new AuthorizationInfo
{
Device = device,
@@ -71,17 +74,17 @@ namespace Jellyfin.Server.Implementations.Tests.QuickConnect
}));
[Fact]
public void TryConnect_QuickConnectUnavailable_ThrowsAuthenticationException()
public async Task TryConnect_QuickConnectUnavailable_ThrowsAuthenticationException()
{
_config.QuickConnectAvailable = false;
Assert.Throws<AuthenticationException>(() => _quickConnectManager.TryConnect(_quickConnectAuthInfo));
await Assert.ThrowsAsync<AuthenticationException>(() => _quickConnectManager.TryConnect(_quickConnectAuthInfo));
}
[Fact]
public void CheckRequestStatus_QuickConnectUnavailable_ThrowsAuthenticationException()
public async Task CheckRequestStatus_QuickConnectUnavailable_ThrowsAuthenticationException()
{
_config.QuickConnectAvailable = false;
Assert.Throws<AuthenticationException>(() => _quickConnectManager.CheckRequestStatus(string.Empty));
await Assert.ThrowsAsync<AuthenticationException>(() => _quickConnectManager.CheckRequestStatus(string.Empty));
}
[Fact]
@@ -92,10 +95,10 @@ namespace Jellyfin.Server.Implementations.Tests.QuickConnect
}
[Fact]
public void GetAuthorizedRequest_QuickConnectUnavailable_ThrowsAuthenticationException()
public async Task GetAuthorizedRequest_QuickConnectUnavailable_ThrowsAuthenticationException()
{
_config.QuickConnectAvailable = false;
Assert.Throws<AuthenticationException>(() => _quickConnectManager.GetAuthorizedRequest(string.Empty));
await Assert.ThrowsAsync<AuthenticationException>(() => _quickConnectManager.GetAuthorizedRequest(string.Empty));
}
[Fact]
@@ -106,34 +109,83 @@ namespace Jellyfin.Server.Implementations.Tests.QuickConnect
}
[Fact]
public void CheckRequestStatus_QuickConnectAvailable_Success()
public async Task CheckRequestStatus_QuickConnectAvailable_Success()
{
_config.QuickConnectAvailable = true;
var res1 = _quickConnectManager.TryConnect(_quickConnectAuthInfo);
var res2 = _quickConnectManager.CheckRequestStatus(res1.Secret);
Assert.Equal(res1, res2);
var res1 = await _quickConnectManager.TryConnect(_quickConnectAuthInfo);
var res2 = await _quickConnectManager.CheckRequestStatus(res1.Secret);
Assert.Equal(res1.Secret, res2.Secret);
Assert.Equal(res1.Code, res2.Code);
}
[Fact]
public void CheckRequestStatus_UnknownSecret_ThrowsResourceNotFoundException()
public async Task CheckRequestStatus_UnknownSecret_ThrowsResourceNotFoundException()
{
_config.QuickConnectAvailable = true;
Assert.Throws<ResourceNotFoundException>(() => _quickConnectManager.CheckRequestStatus("Unknown secret"));
await Assert.ThrowsAsync<ResourceNotFoundException>(() => _quickConnectManager.CheckRequestStatus("Unknown secret"));
}
[Fact]
public void GetAuthorizedRequest_UnknownSecret_ThrowsResourceNotFoundException()
public async Task GetAuthorizedRequest_UnknownSecret_ThrowsResourceNotFoundException()
{
_config.QuickConnectAvailable = true;
Assert.Throws<ResourceNotFoundException>(() => _quickConnectManager.GetAuthorizedRequest("Unknown secret"));
await Assert.ThrowsAsync<ResourceNotFoundException>(() => _quickConnectManager.GetAuthorizedRequest("Unknown secret"));
}
[Fact]
public async Task AuthorizeRequest_QuickConnectAvailable_Success()
{
_config.QuickConnectAvailable = true;
var res = _quickConnectManager.TryConnect(_quickConnectAuthInfo);
var res = await _quickConnectManager.TryConnect(_quickConnectAuthInfo);
Assert.True(await _quickConnectManager.AuthorizeRequest(Guid.Empty, res.Code));
}
[Fact]
public async Task AuthorizeRequest_RacedOnOneCode_SucceedsOnce()
{
_config.QuickConnectAvailable = true;
var res = await _quickConnectManager.TryConnect(_quickConnectAuthInfo);
var outcomes = await Task.WhenAll(
Task.Run(() => AuthorizeAsync(res.Code)),
Task.Run(() => AuthorizeAsync(res.Code)));
Assert.Single(outcomes, authorized => authorized);
}
[Fact]
public async Task GetAuthorizedRequest_SecondExchange_ReturnsTheSameResult()
{
_config.QuickConnectAvailable = true;
var res = await _quickConnectManager.TryConnect(_quickConnectAuthInfo);
await _quickConnectManager.AuthorizeRequest(Guid.Empty, res.Code);
var first = await _quickConnectManager.GetAuthorizedRequest(res.Secret);
var second = await _quickConnectManager.GetAuthorizedRequest(res.Secret);
Assert.Same(first, second);
}
[Fact]
public async Task AuthorizeRequest_OfAnAuthorizedRequest_ThrowsConflictException()
{
_config.QuickConnectAvailable = true;
var res = await _quickConnectManager.TryConnect(_quickConnectAuthInfo);
await _quickConnectManager.AuthorizeRequest(Guid.Empty, res.Code);
await Assert.ThrowsAsync<ConflictException>(() => _quickConnectManager.AuthorizeRequest(Guid.Empty, res.Code));
}
private async Task<bool> AuthorizeAsync(string code)
{
try
{
return await _quickConnectManager.AuthorizeRequest(Guid.Empty, code).ConfigureAwait(false);
}
catch (ConflictException)
{
return false;
}
}
}
}
@@ -1,6 +1,8 @@
using System;
using System.Collections.Generic;
using System.IO;
using System.Linq;
using System.Reflection;
using System.Runtime.CompilerServices;
using Emby.Server.Implementations.ScheduledTasks.Tasks;
using MediaBrowser.Controller.ScheduledTasks;
@@ -11,6 +13,14 @@ namespace Jellyfin.Server.Implementations.Tests.ScheduledTasks;
public class ScanLeaderOptionsTests
{
private static readonly Assembly[] _taskAssemblies =
{
typeof(DeleteTranscodeFileTask).Assembly,
typeof(MediaBrowser.Providers.Lyric.LyricScheduledTask).Assembly,
typeof(Jellyfin.LiveTv.Guide.RefreshGuideScheduledTask).Assembly,
typeof(Jellyfin.MediaEncoding.Hls.ScheduledTasks.KeyframeExtractionScheduledTask).Assembly
};
/// <summary>
/// A gated key that matches no registered task silently stops gating anything, so the default
/// set is pinned to the task keys that actually exist in the build.
@@ -18,28 +28,92 @@ public class ScanLeaderOptionsTests
[Fact]
public void DefaultGatedTaskKeys_Should_MatchRegisteredScheduledTasks()
{
var registeredKeys = DiscoverScheduledTaskKeys();
var registeredKeys = DiscoverScheduledTaskKeys(_taskAssemblies);
Assert.NotEmpty(registeredKeys);
Assert.Empty(new ScanLeaderOptions().GatedTaskKeys.Except(registeredKeys, StringComparer.Ordinal));
var unmatched = new ScanLeaderOptions().GatedTaskKeys.Except(registeredKeys, StringComparer.Ordinal).ToList();
Assert.True(
unmatched.Count == 0,
$"Gated keys match no scheduled task: {string.Join(", ", unmatched)}. Known keys: {string.Join(", ", registeredKeys.Order(StringComparer.Ordinal))}");
}
private static HashSet<string> DiscoverScheduledTaskKeys()
/// <summary>
/// A key dropped from the default set silently un-gates that task on every replica, so the whole
/// set is pinned against a hand-maintained expectation rather than read back from the options.
/// </summary>
[Fact]
public void DefaultGatedTaskKeys_Should_BeTheExpectedSet()
{
var keys = new HashSet<string>(StringComparer.Ordinal);
var assemblies = new[]
string[] expected =
{
typeof(DeleteTranscodeFileTask).Assembly,
typeof(Jellyfin.MediaEncoding.Hls.ScheduledTasks.KeyframeExtractionScheduledTask).Assembly
"AudioNormalization",
"CleanupUserDataTask",
"DownloadLyrics",
"DownloadSubtitles",
"KeyframeExtraction",
"MoveTrickplayImages",
"OptimizeDatabaseTask",
"PluginUpdates",
"RefreshChapterImages",
"RefreshGuide",
"RefreshInternetChannels",
"RefreshLibrary",
"RefreshPeople",
"RefreshTrickplayImages",
"TaskExtractMediaSegments",
"TmdbRefreshUpcomingEpisodes"
};
foreach (var type in assemblies.SelectMany(a => a.GetTypes()))
var actual = new ScanLeaderOptions().GatedTaskKeys;
var missing = expected.Except(actual, StringComparer.Ordinal).ToList();
var unexpected = actual.Except(expected, StringComparer.Ordinal).ToList();
Assert.True(
missing.Count == 0 && unexpected.Count == 0,
$"Default gated task keys drifted. Missing: {Describe(missing)}. Unexpected: {Describe(unexpected)}.");
}
/// <summary>
/// The key universe is only as complete as the assemblies it is read from, so a task added to an
/// unscanned assembly must fail here rather than narrow what the previous test can catch.
/// </summary>
[Fact]
public void TaskAssemblies_Should_CoverEveryAssemblyDeclaringScheduledTasks()
{
var scanned = _taskAssemblies.Select(a => a.GetName().Name).ToHashSet(StringComparer.Ordinal);
var missing = new List<string>();
foreach (var path in Directory.EnumerateFiles(AppContext.BaseDirectory, "*.dll"))
{
if (type.IsAbstract || type.IsInterface || !typeof(IScheduledTask).IsAssignableFrom(type))
var name = Path.GetFileNameWithoutExtension(path);
if (scanned.Contains(name)
|| name.EndsWith(".Tests", StringComparison.Ordinal)
|| !(name.StartsWith("Jellyfin.", StringComparison.Ordinal)
|| name.StartsWith("Emby.", StringComparison.Ordinal)
|| name.StartsWith("MediaBrowser.", StringComparison.Ordinal)))
{
continue;
}
if (GetScheduledTaskTypes(Assembly.LoadFrom(path)).Any())
{
missing.Add(name);
}
}
Assert.True(missing.Count == 0, $"Assemblies declaring scheduled tasks but not scanned: {string.Join(", ", missing)}");
}
private static string Describe(IReadOnlyCollection<string> keys)
=> keys.Count == 0 ? "none" : string.Join(", ", keys.Order(StringComparer.Ordinal));
private static HashSet<string> DiscoverScheduledTaskKeys(IEnumerable<Assembly> assemblies)
{
var keys = new HashSet<string>(StringComparer.Ordinal);
foreach (var type in assemblies.SelectMany(GetScheduledTaskTypes))
{
// Task keys are constant expressions, so an uninitialised instance is enough to read
// them without standing up each task's dependency graph.
var task = (IScheduledTask)RuntimeHelpers.GetUninitializedObject(type);
@@ -48,4 +122,21 @@ public class ScanLeaderOptionsTests
return keys;
}
private static IEnumerable<Type> GetScheduledTaskTypes(Assembly assembly)
{
Type?[] types;
try
{
types = assembly.GetTypes();
}
catch (ReflectionTypeLoadException ex)
{
types = ex.Types;
}
return types
.Where(t => t is not null && !t.IsAbstract && !t.IsInterface && typeof(IScheduledTask).IsAssignableFrom(t))
.Select(t => t!);
}
}
@@ -0,0 +1,94 @@
using System;
using System.Collections.Generic;
using System.Globalization;
using System.Linq;
using Emby.Server.Implementations.TV;
using Jellyfin.Database.Implementations.Entities;
using MediaBrowser.Controller.Configuration;
using MediaBrowser.Controller.Dto;
using MediaBrowser.Controller.Entities;
using MediaBrowser.Controller.Entities.TV;
using MediaBrowser.Controller.Library;
using MediaBrowser.Controller.Persistence;
using MediaBrowser.Model.Configuration;
using MediaBrowser.Model.Querying;
using Moq;
using Xunit;
namespace Jellyfin.Server.Implementations.Tests.TV;
public class TVSeriesManagerNextUpTests
{
[Theory]
[InlineData(1)]
[InlineData(25)]
[InlineData(200)]
public void GetNextUp_ReadsUserDataInABoundedNumberOfQueries(int seriesCount)
{
var user = new User("next-up", "provider", "provider");
var libraryManager = new Mock<ILibraryManager>();
var userDataManager = new Mock<IUserDataManager>();
var seriesKeys = Enumerable.Range(0, seriesCount)
.Select(i => i.ToString(CultureInfo.InvariantCulture))
.ToList();
var batch = seriesKeys.ToDictionary(
key => key,
key => new NextUpEpisodeBatchResult
{
NextUp = new Episode { Id = Guid.NewGuid(), Name = "Next " + key },
LastWatched = new Episode { Id = Guid.NewGuid(), Name = "Watched " + key }
});
libraryManager
.Setup(l => l.GetNextUpSeriesKeys(It.IsAny<InternalItemsQuery>(), It.IsAny<IReadOnlyCollection<BaseItem>>(), It.IsAny<DateTime>()))
.Returns(seriesKeys);
libraryManager
.Setup(l => l.GetNextUpEpisodesBatch(It.IsAny<InternalItemsQuery>(), It.IsAny<IReadOnlyList<string>>(), It.IsAny<bool>(), It.IsAny<bool>()))
.Returns(batch);
libraryManager.Setup(l => l.GetLinkedAlternateVersions(It.IsAny<Video>())).Returns([]);
libraryManager.Setup(l => l.GetLocalAlternateVersionIds(It.IsAny<Video>())).Returns([]);
var batchReads = 0;
userDataManager
.Setup(u => u.GetUserDataBatch(It.IsAny<IReadOnlyList<BaseItem>>(), It.IsAny<User>()))
.Returns<IReadOnlyList<BaseItem>, User>((items, _) =>
{
batchReads++;
return items.DistinctBy(i => i.Id).ToDictionary(
i => i.Id,
i => new UserItemData { Key = i.Id.ToString("N", CultureInfo.InvariantCulture) });
});
var previousLibraryManager = BaseItem.LibraryManager;
BaseItem.LibraryManager = libraryManager.Object;
try
{
var manager = new TVSeriesManager(userDataManager.Object, libraryManager.Object, CreateConfigurationManager());
var result = manager.GetNextUp(
new NextUpQuery { User = user, EnableTotalRecordCount = true },
[],
new DtoOptions(false));
Assert.Equal(seriesCount, result.TotalRecordCount);
// Selection, the resume check and the last played date: three reads whatever the library holds.
Assert.Equal(3, batchReads);
userDataManager.Verify(u => u.GetUserData(It.IsAny<User>(), It.IsAny<BaseItem>()), Times.Never);
}
finally
{
BaseItem.LibraryManager = previousLibraryManager;
}
}
private static IServerConfigurationManager CreateConfigurationManager()
{
var configurationManager = new Mock<IServerConfigurationManager>();
configurationManager.SetupGet(c => c.Configuration).Returns(new ServerConfiguration());
return configurationManager.Object;
}
}
@@ -13,6 +13,7 @@ namespace Jellyfin.Server.Tests.HighAvailability;
/// 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.
/// </summary>
[Collection("JellyfinSectionConfiguration")]
public sealed class JellyfinSectionConfigurationTests : IDisposable
{
private const string RedisKey = "Jellyfin:TranscodeStore:RedisConnectionString";
@@ -0,0 +1,223 @@
using System;
using System.Collections.Concurrent;
using System.Globalization;
using System.IO;
using System.Net;
using System.Net.Sockets;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using StackExchange.Redis;
namespace Jellyfin.Server.Tests.HighAvailability;
/// <summary>
/// A loopback TCP proxy in front of a Redis server. Cutting it drops every connection through it and
/// refuses new ones, so a test can take Redis away from one instance mid-flow - and give it back - the
/// way a restarted valkey does, and watch what a real StackExchange.Redis client makes of it.
/// </summary>
public sealed class RedisFaultProxy : IAsyncDisposable
{
private readonly ConcurrentDictionary<TcpClient, byte> _live = new();
private readonly CancellationTokenSource _cts = new();
private readonly TcpListener _listener;
private readonly string _targetHost;
private readonly int _targetPort;
private readonly int _port;
private volatile bool _cut;
private volatile byte[]? _cutAfterMarker;
private RedisFaultProxy(TcpListener listener, int port, string targetHost, int targetPort)
{
_listener = listener;
_port = port;
_targetHost = targetHost;
_targetPort = targetPort;
}
/// <summary>
/// Gets a connection string pointing at the proxy. The timeouts are short so a cut surfaces as a
/// failure in seconds rather than in the library's minute-scale defaults.
/// </summary>
public string ConnectionString => string.Create(
CultureInfo.InvariantCulture,
$"127.0.0.1:{_port},abortConnect=false,connectTimeout=500,syncTimeout=2000,connectRetry=1");
/// <summary>
/// Starts a proxy in front of the server named by <paramref name="target"/>.
/// </summary>
/// <param name="target">The connection string of the server to forward to.</param>
/// <returns>The running proxy.</returns>
public static RedisFaultProxy Start(string target)
{
var endpoint = ConfigurationOptions.Parse(target).EndPoints[0];
var (host, port) = endpoint switch
{
DnsEndPoint dns => (dns.Host, dns.Port),
IPEndPoint ip => (ip.Address.ToString(), ip.Port),
_ => throw new NotSupportedException("Unsupported endpoint " + endpoint)
};
var listener = new TcpListener(IPAddress.Loopback, 0);
listener.Start();
var proxy = new RedisFaultProxy(listener, ((IPEndPoint)listener.LocalEndpoint).Port, host, port);
_ = Task.Run(proxy.AcceptAsync);
return proxy;
}
/// <summary>
/// Takes Redis away from everything connected through the proxy.
/// </summary>
public void Cut()
{
_cut = true;
DropLiveConnections();
}
/// <summary>
/// Arms a cut for the moment after a command containing <paramref name="marker"/> has been forwarded
/// and answered, so a test can take Redis away between two round trips of one operation rather than
/// only before or after all of them.
/// </summary>
/// <param name="marker">Text that identifies the command to cut after.</param>
public void CutAfterForwarding(string marker) => _cutAfterMarker = Encoding.UTF8.GetBytes(marker);
/// <summary>
/// Lets connections through again. Clients reconnect on their own schedule, so callers have to wait
/// for the connection to come back rather than assume it already has.
/// </summary>
public void Restore()
{
_cutAfterMarker = null;
_cut = false;
}
/// <inheritdoc/>
public async ValueTask DisposeAsync()
{
_cut = true;
await _cts.CancelAsync().ConfigureAwait(false);
_listener.Stop();
DropLiveConnections();
_cts.Dispose();
}
private void DropLiveConnections()
{
foreach (var client in _live.Keys)
{
if (_live.TryRemove(client, out _))
{
client.Dispose();
}
}
}
private async Task AcceptAsync()
{
while (!_cts.IsCancellationRequested)
{
TcpClient client;
try
{
client = await _listener.AcceptTcpClientAsync(_cts.Token).ConfigureAwait(false);
}
catch (Exception exception) when (exception is OperationCanceledException or SocketException or ObjectDisposedException)
{
return;
}
if (_cut)
{
client.Dispose();
continue;
}
_ = Task.Run(() => ForwardAsync(client));
}
}
private async Task ForwardAsync(TcpClient client)
{
TcpClient? upstream = null;
try
{
upstream = new TcpClient();
await upstream.ConnectAsync(_targetHost, _targetPort, _cts.Token).ConfigureAwait(false);
_live[client] = 0;
_live[upstream] = 0;
// Registered first, then rechecked: a cut concurrent with this connect would otherwise drop
// the live connections before this pair joined them and leave it running through the outage.
if (_cut)
{
return;
}
var clientStream = client.GetStream();
var upstreamStream = upstream.GetStream();
await Task.WhenAny(
CopyFromClientAsync(clientStream, upstreamStream),
CopyAsync(upstreamStream, clientStream)).ConfigureAwait(false);
}
catch (Exception exception) when (exception is IOException or SocketException or OperationCanceledException or ObjectDisposedException)
{
}
finally
{
_live.TryRemove(client, out _);
client.Dispose();
if (upstream is not null)
{
_live.TryRemove(upstream, out _);
upstream.Dispose();
}
}
}
private async Task CopyFromClientAsync(NetworkStream from, NetworkStream to)
{
var buffer = new byte[16 * 1024];
try
{
while (true)
{
var read = await from.ReadAsync(buffer, _cts.Token).ConfigureAwait(false);
if (read == 0)
{
return;
}
await to.WriteAsync(buffer.AsMemory(0, read), _cts.Token).ConfigureAwait(false);
var marker = _cutAfterMarker;
if (marker is not null && buffer.AsSpan(0, read).IndexOf(marker) >= 0)
{
_cutAfterMarker = null;
// Long enough for the server to have applied the command that was just forwarded.
await Task.Delay(TimeSpan.FromMilliseconds(250), _cts.Token).ConfigureAwait(false);
Cut();
return;
}
}
}
catch (Exception exception) when (exception is IOException or SocketException or OperationCanceledException or ObjectDisposedException)
{
}
}
private async Task CopyAsync(NetworkStream from, NetworkStream to)
{
try
{
await from.CopyToAsync(to, _cts.Token).ConfigureAwait(false);
}
catch (Exception exception) when (exception is IOException or SocketException or OperationCanceledException or ObjectDisposedException)
{
}
}
}
@@ -0,0 +1,91 @@
using System;
using System.Threading.Tasks;
using StackExchange.Redis;
using Testcontainers.Redis;
namespace Jellyfin.Server.Tests.HighAvailability;
/// <summary>
/// Hands out a Redis server for the tests that need one. A server named by <c>JELLYFIN_TEST_REDIS</c> is
/// used as is, so CI can run one beside the step instead of a docker daemon of its own; without it a
/// container is started through testcontainers.
/// </summary>
public sealed class RedisTestServer : IAsyncDisposable
{
/// <summary>
/// The connection string of an already running server.
/// </summary>
public const string ConnectionStringVariable = "JELLYFIN_TEST_REDIS";
private readonly RedisContainer? _container;
private RedisTestServer(RedisContainer? container, string connectionString)
{
_container = container;
ConnectionString = connectionString;
}
/// <summary>
/// Gets the connection string of the running server.
/// </summary>
public string ConnectionString { get; }
/// <summary>
/// Starts or attaches to a Redis server and waits until it accepts connections.
/// </summary>
/// <returns>The running server.</returns>
public static async Task<RedisTestServer> StartAsync()
{
var provided = Environment.GetEnvironmentVariable(ConnectionStringVariable);
if (!string.IsNullOrWhiteSpace(provided))
{
var attached = new RedisTestServer(null, provided);
await attached.WaitUntilReadyAsync().ConfigureAwait(false);
return attached;
}
var container = new RedisBuilder("redis:7-alpine").Build();
await container.StartAsync().ConfigureAwait(false);
var started = new RedisTestServer(container, container.GetConnectionString());
await started.WaitUntilReadyAsync().ConfigureAwait(false);
return started;
}
/// <summary>
/// Opens a connection of its own, so each in-process stand-in for a replica talks to the server the
/// way a separate pod would.
/// </summary>
/// <returns>A new multiplexer.</returns>
public async Task<IConnectionMultiplexer> ConnectAsync()
=> await ConnectionMultiplexer.ConnectAsync(ConnectionString).ConfigureAwait(false);
/// <inheritdoc/>
public async ValueTask DisposeAsync()
{
if (_container is not null)
{
await _container.DisposeAsync().ConfigureAwait(false);
}
}
private async Task WaitUntilReadyAsync()
{
for (var attempt = 1; ; attempt++)
{
try
{
var connection = await ConnectionMultiplexer.ConnectAsync(ConnectionString).ConfigureAwait(false);
await using (connection.ConfigureAwait(false))
{
await connection.GetDatabase().PingAsync().ConfigureAwait(false);
return;
}
}
catch (RedisException) when (attempt < 60)
{
await Task.Delay(TimeSpan.FromSeconds(1)).ConfigureAwait(false);
}
}
}
}
@@ -11,9 +11,10 @@ using Xunit;
namespace Jellyfin.Server.Tests.HighAvailability;
/// <summary>
/// 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.
/// A transcode store that cannot be reached degrades silently: the client is configured not to abort
/// the connection and every call site swallows failures. Quick connect's startup read is what stops a
/// pod coming up against a dead store; this probe only reports it, so both of its outcomes are pinned
/// here - including that it never throws, which is what keeps the two policies from competing.
/// </summary>
public sealed class TranscodeStoreConnectivityProbeTests
{
@@ -12,6 +12,7 @@
<PackageReference Include="Microsoft.NET.Test.Sdk" />
<PackageReference Include="Npgsql" />
<PackageReference Include="Testcontainers.PostgreSql" />
<PackageReference Include="Testcontainers.Redis" />
<PackageReference Include="xunit.v3" />
<PackageReference Include="xunit.runner.visualstudio">
<PrivateAssets>all</PrivateAssets>
@@ -0,0 +1,243 @@
using System;
using System.Collections.Generic;
using System.Globalization;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Emby.Server.Implementations.Library;
using Jellyfin.Database.Implementations;
using Jellyfin.Database.Implementations.DbConfiguration;
using Jellyfin.Database.Implementations.Entities;
using Jellyfin.Database.Implementations.Locking;
using Jellyfin.Database.Providers.PostgreSQL;
using Jellyfin.Server.Tests.Migrations;
using MediaBrowser.Controller.Configuration;
using MediaBrowser.Model.Configuration;
using MediaBrowser.Model.Entities;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging.Abstractions;
using Moq;
using Npgsql;
using Xunit;
using AudioBook = MediaBrowser.Controller.Entities.AudioBook;
namespace Jellyfin.Server.Tests.Library;
/// <summary>
/// Two independently constructed <see cref="UserDataManager"/> instances over one PostgreSQL database are the
/// in-process stand-in for two replicas sharing one database: what either of them writes, the other has to
/// see on its very next read, and a read-modify-write on one must not roll back the other's.
/// </summary>
[Trait("Category", "RequiresDocker")]
public sealed class UserDataManagerReplicaTests : IClassFixture<UserDataManagerReplicaTests.DatabaseFixture>
{
private static readonly long _quarterIn = TimeSpan.FromMinutes(20).Ticks;
private readonly NpgsqlDataSource _dataSource;
public UserDataManagerReplicaTests(DatabaseFixture fixture)
{
_dataSource = fixture.DataSource;
}
/// <summary>
/// A resume position written by the replica serving the playback tick has to be the position the next
/// request reads, whichever replica it lands on - both through the single item read the write path uses
/// and through the batch read the library pages render from.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task ResumePositionWrittenOnOneReplica_IsReadOnAnother()
{
var cancellationToken = TestContext.Current.CancellationToken;
var itemId = Guid.NewGuid();
var user = await CreateUserAndItemAsync(_dataSource, itemId, cancellationToken);
var replicaA = CreateManager(_dataSource);
var replicaB = CreateManager(_dataSource);
var itemOnA = new AudioBook { Id = itemId, Name = "Replica Book" };
var itemOnB = new AudioBook { Id = itemId, Name = "Replica Book" };
var early = replicaA.GetUserData(user, itemOnA)!;
early.PlaybackPositionTicks = TimeSpan.FromMinutes(5).Ticks;
replicaA.SaveUserData(user, itemOnA, early, UserDataSaveReason.PlaybackProgress, cancellationToken);
// Replica B materialised the item before the later tick, so it holds the earlier row in memory.
itemOnB.UserData = await LoadUserDataAsync(_dataSource, itemId, cancellationToken);
var later = replicaA.GetUserData(user, itemOnA)!;
later.PlaybackPositionTicks = _quarterIn;
replicaA.SaveUserData(user, itemOnA, later, UserDataSaveReason.PlaybackProgress, cancellationToken);
Assert.Equal(_quarterIn, replicaB.GetUserData(user, itemOnB)!.PlaybackPositionTicks);
Assert.Equal(_quarterIn, replicaB.GetUserDataBatch([itemOnB], user)[itemId].PlaybackPositionTicks);
}
/// <summary>
/// The playback tick is a read-modify-write of the whole row, so a tick served by one replica must build
/// on the favourite another replica just recorded instead of writing it back out.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task PlaybackTickOnOneReplica_KeepsFavouriteSetOnAnother()
{
var cancellationToken = TestContext.Current.CancellationToken;
var itemId = Guid.NewGuid();
var user = await CreateUserAndItemAsync(_dataSource, itemId, cancellationToken);
var replicaA = CreateManager(_dataSource);
var replicaB = CreateManager(_dataSource);
var itemOnA = new AudioBook { Id = itemId, Name = "Replica Book" };
var itemOnB = new AudioBook { Id = itemId, Name = "Replica Book" };
var seed = replicaA.GetUserData(user, itemOnA)!;
seed.PlaybackPositionTicks = TimeSpan.FromMinutes(5).Ticks;
replicaA.SaveUserData(user, itemOnA, seed, UserDataSaveReason.PlaybackProgress, cancellationToken);
// Replica B is serving the playback session and read the item before the favourite was recorded.
itemOnB.UserData = await LoadUserDataAsync(_dataSource, itemId, cancellationToken);
var favourited = replicaA.GetUserData(user, itemOnA)!;
favourited.IsFavorite = true;
replicaA.SaveUserData(user, itemOnA, favourited, UserDataSaveReason.UpdateUserRating, cancellationToken);
var tick = replicaB.GetUserData(user, itemOnB)!;
tick.PlaybackPositionTicks = _quarterIn;
replicaB.SaveUserData(user, itemOnB, tick, UserDataSaveReason.PlaybackProgress, cancellationToken);
var stored = replicaA.GetUserData(user, itemOnA)!;
Assert.True(stored.IsFavorite);
Assert.Equal(_quarterIn, stored.PlaybackPositionTicks);
}
/// <summary>
/// A tick that lands on the other replica has to carry the position forward from where the session
/// actually is, not from the position that replica happened to have in memory.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task PlaybackTickOnOneReplica_ResumesFromThePositionAnotherWrote()
{
var cancellationToken = TestContext.Current.CancellationToken;
var itemId = Guid.NewGuid();
var user = await CreateUserAndItemAsync(_dataSource, itemId, cancellationToken);
var replicaA = CreateManager(_dataSource);
var replicaB = CreateManager(_dataSource);
var itemOnA = new AudioBook { Id = itemId, Name = "Replica Book" };
var itemOnB = new AudioBook { Id = itemId, Name = "Replica Book" };
var seed = replicaA.GetUserData(user, itemOnA)!;
seed.PlaybackPositionTicks = TimeSpan.FromMinutes(5).Ticks;
replicaA.SaveUserData(user, itemOnA, seed, UserDataSaveReason.PlaybackProgress, cancellationToken);
itemOnB.UserData = await LoadUserDataAsync(_dataSource, itemId, cancellationToken);
// The viewer seeks forward and the tick reporting it lands on replica A.
var seeked = replicaA.GetUserData(user, itemOnA)!;
seeked.PlaybackPositionTicks = _quarterIn;
replicaA.SaveUserData(user, itemOnA, seeked, UserDataSaveReason.PlaybackProgress, cancellationToken);
// The next tick lands on replica B, which adds ten seconds to whatever it reads.
var tick = replicaB.GetUserData(user, itemOnB)!;
tick.PlaybackPositionTicks += TimeSpan.FromSeconds(10).Ticks;
replicaB.SaveUserData(user, itemOnB, tick, UserDataSaveReason.PlaybackProgress, cancellationToken);
var stored = replicaA.GetUserData(user, itemOnA)!;
Assert.Equal(_quarterIn + TimeSpan.FromSeconds(10).Ticks, stored.PlaybackPositionTicks);
}
private static UserDataManager CreateManager(NpgsqlDataSource dataSource)
{
var config = new Mock<IServerConfigurationManager>();
config.SetupGet(c => c.Configuration).Returns(new ServerConfiguration());
return new UserDataManager(config.Object, new DataSourceContextFactory(dataSource));
}
private static async Task<ICollection<UserData>> LoadUserDataAsync(NpgsqlDataSource dataSource, Guid itemId, CancellationToken cancellationToken)
{
var context = CreateContext(dataSource);
await using (context.ConfigureAwait(false))
{
return await context.UserData
.AsNoTracking()
.Where(e => e.ItemId.Equals(itemId))
.ToArrayAsync(cancellationToken)
.ConfigureAwait(false);
}
}
private static async Task<User> CreateUserAndItemAsync(NpgsqlDataSource dataSource, Guid itemId, CancellationToken cancellationToken)
{
var context = CreateContext(dataSource);
await using (context.ConfigureAwait(false))
{
var user = new User("replica-user-" + itemId.ToString("N", CultureInfo.InvariantCulture), "provider", "provider");
context.Users.Add(user);
context.BaseItems.Add(new BaseItemEntity { Id = itemId, Type = typeof(AudioBook).FullName! });
await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false);
return user;
}
}
private static JellyfinDbContext CreateContext(NpgsqlDataSource dataSource)
{
var optionsBuilder = new DbContextOptionsBuilder<JellyfinDbContext>();
var provider = new PostgreSqlDatabaseProvider(dataSource);
provider.Initialise(optionsBuilder, new DatabaseConfigurationOptions { DatabaseType = "PostgreSQL" });
return new JellyfinDbContext(
optionsBuilder.Options,
NullLogger<JellyfinDbContext>.Instance,
provider,
new NoLockBehavior(NullLogger<NoLockBehavior>.Instance));
}
/// <summary>
/// Hands every <see cref="UserDataManager"/> its own context over the one shared database, the way the
/// pooled factory does in the server.
/// </summary>
private sealed class DataSourceContextFactory : IDbContextFactory<JellyfinDbContext>
{
private readonly NpgsqlDataSource _dataSource;
public DataSourceContextFactory(NpgsqlDataSource dataSource)
{
_dataSource = dataSource;
}
public JellyfinDbContext CreateDbContext() => CreateContext(_dataSource);
}
/// <summary>
/// Builds the schema once for the whole class. Every test keeps to its own user and item, so one
/// database serves all of them and the shared server is spared three schema builds.
/// </summary>
public sealed class DatabaseFixture : IAsyncLifetime
{
private PostgreSqlTestServer _server = null!;
public NpgsqlDataSource DataSource { get; private set; } = null!;
/// <inheritdoc/>
public async ValueTask InitializeAsync()
{
_server = await PostgreSqlTestServer.StartAsync().ConfigureAwait(false);
var connectionString = await _server.CreateDatabaseAsync("userdata_replica", CancellationToken.None).ConfigureAwait(false);
DataSource = new NpgsqlDataSourceBuilder(connectionString).Build();
var context = CreateContext(DataSource);
await using (context.ConfigureAwait(false))
{
await context.Database.EnsureCreatedAsync(CancellationToken.None).ConfigureAwait(false);
}
}
/// <inheritdoc/>
public async ValueTask DisposeAsync()
{
await DataSource.DisposeAsync().ConfigureAwait(false);
await _server.DisposeAsync().ConfigureAwait(false);
}
}
}
@@ -0,0 +1,374 @@
using System;
using System.Collections.Generic;
using System.Globalization;
using System.Threading;
using System.Threading.Tasks;
using Emby.Server.Implementations.QuickConnect;
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.HighAvailability;
using Jellyfin.Server.Tests.Migrations;
using MediaBrowser.Common.Extensions;
using MediaBrowser.Controller.Authentication;
using MediaBrowser.Controller.Configuration;
using MediaBrowser.Controller.Devices;
using MediaBrowser.Controller.Library;
using MediaBrowser.Controller.Net;
using MediaBrowser.Controller.QuickConnect;
using MediaBrowser.Controller.Session;
using MediaBrowser.Model.Configuration;
using MediaBrowser.Model.Dto;
using MediaBrowser.Model.QuickConnect;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging.Abstractions;
using Moq;
using Npgsql;
using StackExchange.Redis;
using Xunit;
namespace Jellyfin.Server.Tests.QuickConnect;
/// <summary>
/// Three independently constructed <see cref="QuickConnectManager"/> instances over one PostgreSQL
/// database and one Redis are the in-process stand-in for three replicas without sticky sessions: the
/// initiate, authorize and exchange legs of one flow each land on a different one.
/// </summary>
[Trait("Category", "RequiresDocker")]
public sealed class QuickConnectReplicaTests : IAsyncLifetime
{
private static readonly AuthorizationInfo _authorizationInfo = new AuthorizationInfo
{
Device = "Living Room TV",
DeviceId = "device-1",
Client = "Jellyfin Web",
Version = "1.0.0"
};
private readonly List<IConnectionMultiplexer> _connections = new();
private PostgreSqlTestServer _postgres = null!;
private RedisTestServer _redis = null!;
/// <inheritdoc/>
public async ValueTask InitializeAsync()
{
_postgres = await PostgreSqlTestServer.StartAsync().ConfigureAwait(false);
_redis = await RedisTestServer.StartAsync().ConfigureAwait(false);
}
/// <inheritdoc/>
public async ValueTask DisposeAsync()
{
foreach (var connection in _connections)
{
await connection.DisposeAsync().ConfigureAwait(false);
}
await _redis.DisposeAsync().ConfigureAwait(false);
await _postgres.DisposeAsync().ConfigureAwait(false);
}
/// <summary>
/// The three legs of a quick connect flow land on three different replicas, and the token the third
/// one hands out is the one the second one minted into the shared database.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task InitiateAuthorizeExchange_AcrossThreeReplicas_Succeeds()
{
var cancellationToken = TestContext.Current.CancellationToken;
var connectionString = await _postgres.CreateDatabaseAsync("quickconnect_replica_flow", cancellationToken);
await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build();
var user = await CreateSchemaWithUserAsync(dataSource, cancellationToken);
var replicaA = await CreateReplicaAsync(dataSource, user);
var replicaB = await CreateReplicaAsync(dataSource, user);
var replicaC = await CreateReplicaAsync(dataSource, user);
var initiated = await replicaA.Manager.TryConnect(_authorizationInfo);
// The code is shown to the user on whichever replica serves the dashboard.
Assert.True(await replicaB.Manager.AuthorizeRequest(user.Id, initiated.Code));
var polled = await replicaC.Manager.CheckRequestStatus(initiated.Secret);
Assert.True(polled.Authenticated);
Assert.Equal(initiated.Code, polled.Code);
Assert.Equal(_authorizationInfo.DeviceId, polled.DeviceId);
var exchanged = await replicaC.Manager.GetAuthorizedRequest(initiated.Secret);
Assert.False(string.IsNullOrEmpty(exchanged.AccessToken));
Assert.Equal(user.Id, exchanged.User.Id);
var devices = await replicaA.Devices.GetDevices(new DeviceQuery { AccessToken = exchanged.AccessToken });
Assert.Equal(user.Id, Assert.Single(devices.Items).UserId);
}
/// <summary>
/// Exchanging a secret does not spend it: a client that retries, or whose retry lands on another
/// replica, gets the same access token back rather than a 404, and the device is minted once.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Exchange_RepeatedOnTwoReplicas_ReturnsTheSameToken()
{
var cancellationToken = TestContext.Current.CancellationToken;
var connectionString = await _postgres.CreateDatabaseAsync("quickconnect_replica_reexchange", cancellationToken);
await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build();
var user = await CreateSchemaWithUserAsync(dataSource, cancellationToken);
var replicaA = await CreateReplicaAsync(dataSource, user);
var replicaB = await CreateReplicaAsync(dataSource, user);
var replicaC = await CreateReplicaAsync(dataSource, user);
for (var attempt = 0; attempt < 10; attempt++)
{
var authorizationInfo = AuthorizationInfoFor(attempt);
var initiated = await replicaA.Manager.TryConnect(authorizationInfo);
await replicaB.Manager.AuthorizeRequest(user.Id, initiated.Code);
var exchanged = await Task.WhenAll(
Task.Run(() => ExchangeAsync(replicaA.Manager, initiated.Secret), cancellationToken),
Task.Run(() => ExchangeAsync(replicaC.Manager, initiated.Secret), cancellationToken));
Assert.All(exchanged, outcome => Assert.NotNull(outcome));
Assert.Equal(exchanged[0]!.AccessToken, exchanged[1]!.AccessToken);
// Still there afterwards, on a replica that has not exchanged it yet.
var later = await replicaB.Manager.GetAuthorizedRequest(initiated.Secret);
Assert.Equal(exchanged[0]!.AccessToken, later.AccessToken);
var devices = await replicaA.Devices.GetDevices(new DeviceQuery { DeviceId = authorizationInfo.DeviceId });
Assert.Equal(later.AccessToken, Assert.Single(devices.Items).AccessToken);
}
}
/// <summary>
/// Two replicas authorizing one code at the same time mint one access token between them. A second
/// one would be live, attached to the same device and reachable by nobody.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Authorize_RacedOnTwoReplicas_MintsOneAccessToken()
{
var cancellationToken = TestContext.Current.CancellationToken;
var connectionString = await _postgres.CreateDatabaseAsync("quickconnect_replica_authorize_race", cancellationToken);
await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build();
var user = await CreateSchemaWithUserAsync(dataSource, cancellationToken);
var replicaA = await CreateReplicaAsync(dataSource, user);
var replicaB = await CreateReplicaAsync(dataSource, user);
var replicaC = await CreateReplicaAsync(dataSource, user);
for (var attempt = 0; attempt < 20; attempt++)
{
var authorizationInfo = AuthorizationInfoFor(attempt);
var initiated = await replicaA.Manager.TryConnect(authorizationInfo);
var outcomes = await Task.WhenAll(
Task.Run(() => AuthorizeAsync(replicaB.Manager, user.Id, initiated.Code), cancellationToken),
Task.Run(() => AuthorizeAsync(replicaC.Manager, user.Id, initiated.Code), cancellationToken));
Assert.Single(outcomes, authorized => authorized);
var devices = await replicaA.Devices.GetDevices(new DeviceQuery { DeviceId = authorizationInfo.DeviceId });
var device = Assert.Single(devices.Items);
var exchanged = await replicaA.Manager.GetAuthorizedRequest(initiated.Secret);
Assert.Equal(device.AccessToken, exchanged.AccessToken);
}
}
/// <summary>
/// An expired request is rejected on a replica that never saw it created, rather than resolving to a
/// stale authorization.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task ExpiredRequest_IsRejectedOnEveryReplica()
{
var cancellationToken = TestContext.Current.CancellationToken;
var connectionString = await _postgres.CreateDatabaseAsync("quickconnect_replica_expiry", cancellationToken);
await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build();
var user = await CreateSchemaWithUserAsync(dataSource, cancellationToken);
var replicaA = await CreateReplicaAsync(dataSource, user);
var replicaB = await CreateReplicaAsync(dataSource, user);
var initiated = await replicaA.Manager.TryConnect(_authorizationInfo);
Assert.NotNull(await replicaB.Manager.CheckRequestStatus(initiated.Secret));
// Shorten the stored expiry instead of waiting out the ten minute timeout.
await replicaA.Store.SetRequestAsync(initiated, DateTime.UtcNow.AddSeconds(1), cancellationToken);
await Task.Delay(TimeSpan.FromSeconds(2), cancellationToken);
await Assert.ThrowsAsync<ResourceNotFoundException>(() => replicaB.Manager.CheckRequestStatus(initiated.Secret));
await Assert.ThrowsAsync<ResourceNotFoundException>(() => replicaB.Manager.AuthorizeRequest(user.Id, initiated.Code));
}
/// <summary>
/// An authorization that was never exchanged expires too, so a code authorized and then abandoned
/// cannot be redeemed later from another replica.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task ExpiredAuthorization_IsRejectedOnEveryReplica()
{
var cancellationToken = TestContext.Current.CancellationToken;
var connectionString = await _postgres.CreateDatabaseAsync("quickconnect_replica_auth_expiry", cancellationToken);
await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build();
var user = await CreateSchemaWithUserAsync(dataSource, cancellationToken);
var replicaA = await CreateReplicaAsync(dataSource, user);
var replicaB = await CreateReplicaAsync(dataSource, user);
var initiated = await replicaA.Manager.TryConnect(_authorizationInfo);
await replicaA.Manager.AuthorizeRequest(user.Id, initiated.Code);
var stored = await replicaA.Store.GetRequestBySecretAsync(initiated.Secret, cancellationToken);
Assert.True(stored?.Authenticated);
await replicaA.Store.SetAuthorizationAsync(
initiated.Secret,
new AuthenticationResult { AccessToken = "stale" },
DateTime.UtcNow.AddSeconds(1),
cancellationToken);
await Task.Delay(TimeSpan.FromSeconds(2), cancellationToken);
await Assert.ThrowsAsync<ResourceNotFoundException>(() => replicaB.Manager.GetAuthorizedRequest(initiated.Secret));
}
private static AuthorizationInfo AuthorizationInfoFor(int attempt) => new AuthorizationInfo
{
Device = _authorizationInfo.Device,
DeviceId = string.Create(CultureInfo.InvariantCulture, $"device-{attempt}"),
Client = _authorizationInfo.Client,
Version = _authorizationInfo.Version
};
private static async Task<bool> AuthorizeAsync(IQuickConnect manager, Guid userId, string code)
{
try
{
return await manager.AuthorizeRequest(userId, code).ConfigureAwait(false);
}
catch (ConflictException)
{
return false;
}
}
private static async Task<AuthenticationResult?> ExchangeAsync(IQuickConnect manager, string secret)
{
try
{
return await manager.GetAuthorizedRequest(secret).ConfigureAwait(false);
}
catch (ResourceNotFoundException)
{
return null;
}
}
private static async Task<User> 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("quickconnect-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<JellyfinDbContext>();
var provider = new PostgreSqlDatabaseProvider(dataSource);
provider.Initialise(optionsBuilder, new DatabaseConfigurationOptions { DatabaseType = "PostgreSQL" });
return new JellyfinDbContext(
optionsBuilder.Options,
NullLogger<JellyfinDbContext>.Instance,
provider,
new NoLockBehavior(NullLogger<NoLockBehavior>.Instance));
}
private async Task<Replica> CreateReplicaAsync(NpgsqlDataSource dataSource, User user)
{
var connection = await _redis.ConnectAsync().ConfigureAwait(false);
_connections.Add(connection);
var userManager = new Mock<IUserManager>();
userManager.Setup(manager => manager.GetUserById(user.Id)).Returns(user);
var deviceManager = new DeviceManager(new DataSourceContextFactory(dataSource), userManager.Object);
var configManager = new Mock<IServerConfigurationManager>();
configManager.Setup(manager => manager.Configuration).Returns(new ServerConfiguration { QuickConnectAvailable = true });
// Stands in for SessionManager.AuthenticateDirect: the token has to be minted into the shared
// database, because the replica that exchanges the secret is not the one that authorized it.
var sessionManager = new Mock<ISessionManager>();
sessionManager
.Setup(manager => manager.AuthenticateDirect(It.IsAny<AuthenticationRequest>()))
.Returns<AuthenticationRequest>(async request =>
{
var device = await deviceManager.CreateDevice(
new Device(request.UserId, request.App, request.AppVersion, request.DeviceName, request.DeviceId)).ConfigureAwait(false);
return new AuthenticationResult
{
AccessToken = device.AccessToken,
ServerId = "server-1",
User = new UserDto { Id = user.Id, Name = user.Username, ServerId = "server-1" },
SessionInfo = new SessionInfoDto
{
Id = device.Id.ToString(CultureInfo.InvariantCulture),
UserId = user.Id,
UserName = user.Username,
Client = request.App,
DeviceId = request.DeviceId,
DeviceName = request.DeviceName,
ApplicationVersion = request.AppVersion
}
};
});
var store = new RedisQuickConnectStore(connection, NullLogger<RedisQuickConnectStore>.Instance);
var manager = new QuickConnectManager(
configManager.Object,
NullLogger<QuickConnectManager>.Instance,
sessionManager.Object,
store);
return new Replica(manager, store, deviceManager);
}
private sealed record Replica(IQuickConnect Manager, IQuickConnectStore Store, IDeviceManager Devices);
private sealed class DataSourceContextFactory : IDbContextFactory<JellyfinDbContext>
{
private readonly NpgsqlDataSource _dataSource;
public DataSourceContextFactory(NpgsqlDataSource dataSource)
{
_dataSource = dataSource;
}
public JellyfinDbContext CreateDbContext() => CreateContext(_dataSource);
}
}
@@ -0,0 +1,467 @@
using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Diagnostics;
using System.IO;
using System.Linq;
using System.Net;
using System.Net.Sockets;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using System.Xml.Serialization;
using Emby.Server.Implementations;
using Emby.Server.Implementations.QuickConnect;
using Jellyfin.Server.Extensions;
using Jellyfin.Server.Helpers;
using Jellyfin.Server.Migrations.Stages;
using Jellyfin.Server.ServerSetupApp;
using Jellyfin.Server.Tests.HighAvailability;
using MediaBrowser.Common.Configuration;
using MediaBrowser.Common.Extensions;
using MediaBrowser.Common.Net;
using MediaBrowser.Controller.Net;
using MediaBrowser.Controller.QuickConnect;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.Mvc.Testing;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
using StackExchange.Redis;
using Xunit;
namespace Jellyfin.Server.Tests.QuickConnect;
/// <summary>
/// Brings the server up the way <c>Program</c> does - the real host over <see cref="Startup"/>, the
/// startup and core migrations, then <see cref="ApplicationHost.InitializeServices"/> - to pin down what
/// a pod does when the quick connect store it is configured against cannot be reached.
/// </summary>
[Trait("Category", "RequiresDocker")]
[Collection("JellyfinSectionConfiguration")]
public sealed class QuickConnectStartupTests : IAsyncLifetime
{
private const string RedisConnectionStringVariable = "Jellyfin__TranscodeStore__RedisConnectionString";
private const string FfmpegNoValidationVariable = "JELLYFIN_FFMPEG__NOVALIDATION";
private const string DeadStore = "127.0.0.1:1,abortConnect=false,connectTimeout=250,connectRetry=0,syncTimeout=250";
private const string EagerDeadStore = "127.0.0.1:1,connectTimeout=250,connectRetry=0,syncTimeout=250";
private RedisTestServer _redis = null!;
/// <inheritdoc/>
public async ValueTask InitializeAsync()
{
Environment.SetEnvironmentVariable(FfmpegNoValidationVariable, "true");
_redis = await RedisTestServer.StartAsync().ConfigureAwait(false);
}
/// <inheritdoc/>
public async ValueTask DisposeAsync()
{
Environment.SetEnvironmentVariable(RedisConnectionStringVariable, null);
Environment.SetEnvironmentVariable(FfmpegNoValidationVariable, null);
await _redis.DisposeAsync().ConfigureAwait(false);
}
/// <summary>
/// A store that stays away past the probe's deadline stops the server coming up, so the outage is
/// visible where the server is started instead of arriving later as a failure on every request that
/// needs the store. The connection string carries the <c>abortConnect=false</c> a deployment uses, so
/// the multiplexer connects lazily and only a real read settles whether the store can be served.
/// </summary>
[Fact]
public void StoreDownPastTheDeadline_StopsStartup()
{
Environment.SetEnvironmentVariable(RedisConnectionStringVariable, DeadStore);
using var server = new StartupHarness();
Assert.ThrowsAny<ServiceUnavailableException>(() => server.Services);
Assert.Contains(
server.CriticalEntries,
entry => entry.Contains("Quick connect", StringComparison.Ordinal)
&& entry.Contains("UNREACHABLE", StringComparison.Ordinal)
&& entry.Contains("valkey", StringComparison.Ordinal));
}
/// <summary>
/// Three of the four connection strings the chart documents leave <c>abortConnect</c> at its default,
/// which connects eagerly, so the multiplexer is what fails and it fails while the store is being
/// built rather than on a read. The probe has to retry the build as well as the read and end on the
/// same message, not let a bare <see cref="RedisConnectionException"/> out.
/// </summary>
/// <remarks>
/// The core initialisation migrations are skipped here because <c>JellyfinMigrationService</c> takes
/// an <c>IBackupService</c> eagerly, which reaches the multiplexer through the library manager, so on
/// this shape they fail before the probe is reached at all. That ordering is a separate problem from
/// what the probe does when it runs.
/// </remarks>
[Fact]
public void EagerlyConnectingStoreDownPastTheDeadline_StopsStartupWithTheSameMessage()
{
Environment.SetEnvironmentVariable(RedisConnectionStringVariable, EagerDeadStore);
using var server = new StartupHarness(runCoreInitialisationMigrations: false);
Assert.ThrowsAny<ServiceUnavailableException>(() => server.Services);
Assert.Contains(
server.CriticalEntries,
entry => entry.Contains("Quick connect", StringComparison.Ordinal)
&& entry.Contains("UNREACHABLE", StringComparison.Ordinal)
&& entry.Contains("valkey", StringComparison.Ordinal));
}
/// <summary>
/// A store that is away when the probe first reads it but back inside the deadline lets the server
/// come up. A running instance already rides out a valkey blip; a starting one has to as well, or a
/// rollout is hostage to valkey restarting at the same time.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task StoreBackInsideTheDeadline_StartsAnyway()
{
await using var proxy = RedisFaultProxy.Start(_redis.ConnectionString);
Environment.SetEnvironmentVariable(RedisConnectionStringVariable, proxy.ConnectionString);
using var server = new StartupHarness(() =>
{
proxy.Cut();
_ = Task.Run(async () =>
{
await Task.Delay(TimeSpan.FromSeconds(3)).ConfigureAwait(false);
proxy.Restore();
});
});
Assert.IsType<RedisQuickConnectStore>(server.Services.GetRequiredService<IQuickConnectStore>());
Assert.Contains(server.WarningEntries, entry => entry.Contains("not reachable yet", StringComparison.Ordinal));
Assert.Empty(server.CriticalEntries);
}
/// <summary>
/// A start that failed has to look like a failure to whatever supervises the process. The server runs
/// as its own process here because the exit code is not observable anywhere else, and in a container
/// it must not sit on the setup server for ten minutes before getting there.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task StoreDownPastTheDeadline_ExitsNonZero()
{
var root = Path.Combine(Path.GetTempPath(), "jellyfin-quickconnect-exit", Path.GetRandomFileName());
var configDirectory = Path.Combine(root, "config");
Directory.CreateDirectory(configDirectory);
Directory.CreateDirectory(Path.Combine(root, "cache"));
WriteNetworkConfiguration(configDirectory);
var startInfo = new ProcessStartInfo(Environment.GetEnvironmentVariable("DOTNET_HOST_PATH") ?? "dotnet")
{
RedirectStandardOutput = true,
RedirectStandardError = true,
UseShellExecute = false,
WorkingDirectory = AppContext.BaseDirectory
};
foreach (var argument in new[]
{
"exec",
Path.Combine(AppContext.BaseDirectory, "jellyfin.dll"),
"--datadir", root,
"--cachedir", Path.Combine(root, "cache"),
"--nowebclient"
})
{
startInfo.ArgumentList.Add(argument);
}
startInfo.Environment[RedisConnectionStringVariable] = DeadStore;
startInfo.Environment[FfmpegNoValidationVariable] = "true";
startInfo.Environment["DOTNET_RUNNING_IN_CONTAINER"] = "true";
try
{
var (exitCode, output) = await RunToCompletionAsync(startInfo, TimeSpan.FromMinutes(5));
Assert.Contains("UNREACHABLE", output, StringComparison.Ordinal);
Assert.Equal(1, exitCode);
}
finally
{
TryDelete(root);
}
}
/// <summary>
/// A reachable store lets the server come up, and the quick connect it comes up with holds its
/// requests in the shared store every instance reads.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task ReachableStore_StartsAndServesQuickConnect()
{
Environment.SetEnvironmentVariable(RedisConnectionStringVariable, _redis.ConnectionString);
using var server = new StartupHarness();
Assert.IsType<RedisQuickConnectStore>(server.Services.GetRequiredService<IQuickConnectStore>());
var quickConnect = server.Services.GetRequiredService<IQuickConnect>();
var request = await quickConnect.TryConnect(NewAuthorizationInfo());
Assert.Equal(request.Code, (await quickConnect.CheckRequestStatus(request.Secret)).Code);
await using var redis = await _redis.ConnectAsync();
Assert.True(await redis.GetDatabase().KeyExistsAsync("jellyfin:quickconnect:request:" + request.Secret));
}
/// <summary>
/// Without a connection string the deployment is single-instance, and it starts on the process-local
/// store upstream uses.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task NoConnectionString_StartsOnTheProcessLocalStore()
{
Environment.SetEnvironmentVariable(RedisConnectionStringVariable, null);
using var server = new StartupHarness();
Assert.IsType<InMemoryQuickConnectStore>(server.Services.GetRequiredService<IQuickConnectStore>());
var quickConnect = server.Services.GetRequiredService<IQuickConnect>();
var request = await quickConnect.TryConnect(NewAuthorizationInfo());
Assert.Equal(request.Code, (await quickConnect.CheckRequestStatus(request.Secret)).Code);
}
private static AuthorizationInfo NewAuthorizationInfo() => new AuthorizationInfo
{
DeviceId = Guid.NewGuid().ToString("N"),
Device = "Living Room TV",
Client = "Jellyfin Web",
Version = "1.0.0"
};
// The setup server binds before anything else runs, so the spawned server gets a port of its own.
private static void WriteNetworkConfiguration(string configDirectory)
{
var listener = new TcpListener(IPAddress.Loopback, 0);
listener.Start();
var port = ((IPEndPoint)listener.LocalEndpoint).Port;
listener.Stop();
var configuration = new NetworkConfiguration
{
InternalHttpPort = port,
PublicHttpPort = port,
EnableHttps = false,
AutoDiscovery = false
};
using var writer = new StreamWriter(Path.Combine(configDirectory, "network.xml"));
new XmlSerializer(typeof(NetworkConfiguration)).Serialize(writer, configuration);
}
private static async Task<(int ExitCode, string Output)> RunToCompletionAsync(ProcessStartInfo startInfo, TimeSpan timeout)
{
var output = new StringBuilder();
using var process = new Process { StartInfo = startInfo, EnableRaisingEvents = true };
process.OutputDataReceived += (_, e) => Append(output, e.Data);
process.ErrorDataReceived += (_, e) => Append(output, e.Data);
process.Start();
process.BeginOutputReadLine();
process.BeginErrorReadLine();
using var cancellation = new CancellationTokenSource(timeout);
try
{
await process.WaitForExitAsync(cancellation.Token).ConfigureAwait(false);
}
catch (OperationCanceledException)
{
process.Kill(true);
throw new TimeoutException($"The server did not exit within {timeout}. Output:\n{output}");
}
return (process.ExitCode, output.ToString());
}
private static void Append(StringBuilder output, string? line)
{
if (line is not null)
{
lock (output)
{
output.AppendLine(line);
}
}
}
private static void TryDelete(string path)
{
try
{
Directory.Delete(path, true);
}
catch (IOException)
{
// A temporary directory left behind is not worth failing a test over.
}
}
private sealed class StartupHarness : WebApplicationFactory<Startup>
{
private readonly ConcurrentBag<IDisposable> _disposables = new();
private readonly ConcurrentQueue<(LogLevel Level, string Message)> _entries = new();
private readonly Action? _beforeInitializeServices;
private readonly bool _runCoreInitialisationMigrations;
private readonly string _root = Path.Combine(
Path.GetTempPath(),
"jellyfin-quickconnect-startup",
Path.GetRandomFileName());
static StartupHarness()
{
StartupHelpers.PerformStaticInitialization();
}
public StartupHarness(Action? beforeInitializeServices = null, bool runCoreInitialisationMigrations = true)
{
_beforeInitializeServices = beforeInitializeServices;
_runCoreInitialisationMigrations = runCoreInitialisationMigrations;
}
public IReadOnlyCollection<string> CriticalEntries => Messages(LogLevel.Critical);
public IReadOnlyCollection<string> WarningEntries => Messages(LogLevel.Warning);
protected override IHostBuilder CreateHostBuilder() => new HostBuilder();
protected override void ConfigureWebHost(IWebHostBuilder builder)
{
var commandLineOpts = new StartupOptions();
Directory.CreateDirectory(Path.Combine(_root, "logs"));
Directory.CreateDirectory(Path.Combine(_root, "config"));
Directory.CreateDirectory(Path.Combine(_root, "cache"));
Directory.CreateDirectory(Path.Combine(_root, "jellyfin-web"));
var appPaths = new ServerApplicationPaths(
_root,
Path.Combine(_root, "logs"),
Path.Combine(_root, "config"),
Path.Combine(_root, "cache"),
Path.Combine(_root, "jellyfin-web"));
StartupHelpers.InitLoggingConfigFile(appPaths).GetAwaiter().GetResult();
var startupConfig = Program.CreateAppConfiguration(commandLineOpts, appPaths);
ILoggerFactory loggerFactory = LoggerFactory.Create(logging => logging
.SetMinimumLevel(LogLevel.Warning)
.AddProvider(new RecordingProvider(_entries)));
_disposables.Add(loggerFactory);
var appHost = new CoreAppHost(appPaths, loggerFactory, commandLineOpts, startupConfig);
_disposables.Add(appHost);
builder.ConfigureServices(services => appHost.Init(services))
.ConfigureWebHostBuilder(appHost, startupConfig, appPaths, NullLogger.Instance)
.ConfigureAppConfiguration((context, configuration) => configuration
.SetBasePath(appPaths.ConfigurationDirectoryPath)
.AddInMemoryCollection(Emby.Server.Implementations.ConfigurationOptions.DefaultConfiguration)
.AddEnvironmentVariables("JELLYFIN_")
.AddInMemoryCollection(commandLineOpts.ConvertToConfig()))
.ConfigureServices(services => services.RegisterStartupLogger());
}
protected override IHost CreateHost(IHostBuilder builder)
{
var host = builder.Build();
try
{
var appHost = (CoreAppHost)host.Services.GetRequiredService<MediaBrowser.Common.IApplicationHost>();
appHost.ServiceProvider = host.Services;
var appPaths = (ServerApplicationPaths)host.Services.GetRequiredService<IApplicationPaths>();
var configuration = host.Services.GetRequiredService<IConfiguration>();
Program.ApplyStartupMigrationAsync(appPaths, configuration, new StartupOptions()).GetAwaiter().GetResult();
if (_runCoreInitialisationMigrations)
{
Program.ApplyCoreMigrationsAsync(host.Services, JellyfinMigrationStageTypes.CoreInitialisation).GetAwaiter().GetResult();
}
_beforeInitializeServices?.Invoke();
appHost.InitializeServices(configuration).GetAwaiter().GetResult();
Program.ApplyCoreMigrationsAsync(host.Services, JellyfinMigrationStageTypes.AppInitialisation).GetAwaiter().GetResult();
host.Start();
return host;
}
catch
{
// The factory only owns the host once this returns, so a failed start disposes it here
// or the multiplexer it holds outlives the test.
host.Dispose();
throw;
}
}
protected override void Dispose(bool disposing)
{
base.Dispose(disposing);
foreach (var disposable in _disposables)
{
disposable.Dispose();
}
_disposables.Clear();
TryDelete(_root);
}
private IReadOnlyCollection<string> Messages(LogLevel level)
=> _entries.Where(entry => entry.Level == level).Select(entry => entry.Message).ToArray();
private sealed class RecordingProvider : ILoggerProvider
{
private readonly ConcurrentQueue<(LogLevel Level, string Message)> _entries;
public RecordingProvider(ConcurrentQueue<(LogLevel Level, string Message)> entries)
{
_entries = entries;
}
public ILogger CreateLogger(string categoryName) => new RecordingLogger(_entries);
public void Dispose()
{
}
private sealed class RecordingLogger : ILogger
{
private readonly ConcurrentQueue<(LogLevel Level, string Message)> _entries;
public RecordingLogger(ConcurrentQueue<(LogLevel Level, string Message)> entries)
{
_entries = entries;
}
public IDisposable? BeginScope<TState>(TState state)
where TState : notnull
=> null;
public bool IsEnabled(LogLevel logLevel) => logLevel >= LogLevel.Warning;
public void Log<TState>(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func<TState, Exception?, string> formatter)
{
if (IsEnabled(logLevel))
{
_entries.Enqueue((logLevel, formatter!(state, exception)));
}
}
}
}
}
}
@@ -0,0 +1,261 @@
using System;
using System.IO;
using System.Threading.Tasks;
using Emby.Server.Implementations.QuickConnect;
using Jellyfin.Api.Controllers;
using Jellyfin.Api.Middleware;
using Jellyfin.Server.Tests.HighAvailability;
using MediaBrowser.Common.Extensions;
using MediaBrowser.Controller;
using MediaBrowser.Controller.Authentication;
using MediaBrowser.Controller.Configuration;
using MediaBrowser.Controller.Net;
using MediaBrowser.Controller.QuickConnect;
using MediaBrowser.Controller.Session;
using MediaBrowser.Model.Configuration;
using MediaBrowser.Model.Dto;
using MediaBrowser.Model.QuickConnect;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.Http;
using Microsoft.AspNetCore.Mvc;
using Microsoft.AspNetCore.Mvc.Infrastructure;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging.Abstractions;
using Moq;
using StackExchange.Redis;
using Xunit;
namespace Jellyfin.Server.Tests.QuickConnect;
/// <summary>
/// The status a client actually sees. A missing key and an unreachable valkey are different answers and
/// must not collapse into one: telling a polling client its secret is unknown ends its flow, while 503
/// tells it to keep trying. The exception the store really throws is run through the real exception
/// middleware, so the mapping is exercised rather than assumed.
/// </summary>
[Trait("Category", "RequiresDocker")]
public sealed class QuickConnectStatusCodeTests : IAsyncLifetime
{
private readonly Mock<ISessionManager> _sessionManager = new();
private RedisTestServer _redis = null!;
private RedisFaultProxy _proxy = null!;
private IConnectionMultiplexer _connection = null!;
private QuickConnectManager _manager = null!;
private QuickConnectController _controller = null!;
/// <inheritdoc/>
public async ValueTask InitializeAsync()
{
_redis = await RedisTestServer.StartAsync().ConfigureAwait(false);
_proxy = RedisFaultProxy.Start(_redis.ConnectionString);
_connection = await ConnectionMultiplexer.ConnectAsync(_proxy.ConnectionString).ConfigureAwait(false);
var configManager = new Mock<IServerConfigurationManager>();
configManager.Setup(manager => manager.Configuration).Returns(new ServerConfiguration { QuickConnectAvailable = true });
_manager = new QuickConnectManager(
configManager.Object,
NullLogger<QuickConnectManager>.Instance,
_sessionManager.Object,
new RedisQuickConnectStore(_connection, NullLogger<RedisQuickConnectStore>.Instance));
_controller = new QuickConnectController(_manager, Mock.Of<IAuthorizationContext>());
}
/// <inheritdoc/>
public async ValueTask DisposeAsync()
{
await _connection.DisposeAsync().ConfigureAwait(false);
await _proxy.DisposeAsync().ConfigureAwait(false);
await _redis.DisposeAsync().ConfigureAwait(false);
}
/// <summary>
/// A secret valkey has never heard of is a 404, which is what ends a flow the user abandoned.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Poll_UnknownSecret_IsNotFound()
{
Assert.Equal(
StatusCodes.Status404NotFound,
await StatusCodeAsync(async () => StatusOf(await _controller.GetQuickConnectState(NewSecret()))));
}
/// <summary>
/// The same poll while valkey is unreachable is a 503. This is the bug the shared store is here to
/// avoid: a blip must not tell every polling client that its secret is invalid.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Poll_WhileRedisIsUnreachable_IsServiceUnavailable()
{
var secret = await InitiateAsync();
_proxy.Cut();
Assert.Equal(
StatusCodes.Status503ServiceUnavailable,
await StatusCodeAsync(async () => StatusOf(await _controller.GetQuickConnectState(secret))));
}
/// <summary>
/// The exchange leg tells the two apart the same way.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Exchange_UnknownSecret_IsNotFound()
{
Assert.Equal(
StatusCodes.Status404NotFound,
await StatusCodeAsync(async () =>
{
await _manager.GetAuthorizedRequest(NewSecret()).ConfigureAwait(false);
return StatusCodes.Status200OK;
}));
}
/// <summary>
/// The exchange leg while valkey is unreachable.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Exchange_WhileRedisIsUnreachable_IsServiceUnavailable()
{
var secret = await InitiateAsync();
_proxy.Cut();
Assert.Equal(
StatusCodes.Status503ServiceUnavailable,
await StatusCodeAsync(async () =>
{
await _manager.GetAuthorizedRequest(secret).ConfigureAwait(false);
return StatusCodes.Status200OK;
}));
}
/// <summary>
/// The authorize leg while valkey is unreachable.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Authorize_WhileRedisIsUnreachable_IsServiceUnavailable()
{
var initiated = await _manager.TryConnect(AuthorizationInfo());
_proxy.Cut();
Assert.Equal(
StatusCodes.Status503ServiceUnavailable,
await StatusCodeAsync(async () =>
{
await _manager.AuthorizeRequest(Guid.NewGuid(), initiated.Code).ConfigureAwait(false);
return StatusCodes.Status200OK;
}));
}
/// <summary>
/// A mint that threw leaves its claim taken on purpose, because the write it failed on may have
/// landed. Retrying then has to say so and be a 409 the client can act on, not a 500 and not the
/// untrue claim that the request is already authorized.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Authorize_AfterAMintThatFailed_IsConflictAndSaysToStartAgain()
{
var initiated = await _manager.TryConnect(AuthorizationInfo());
_sessionManager
.Setup(manager => manager.AuthenticateDirect(It.IsAny<AuthenticationRequest>()))
.ThrowsAsync(new InvalidOperationException("mint failed"));
await Assert.ThrowsAsync<InvalidOperationException>(() => _manager.AuthorizeRequest(Guid.NewGuid(), initiated.Code));
var retry = await Record.ExceptionAsync(() => _manager.AuthorizeRequest(Guid.NewGuid(), initiated.Code));
var conflict = Assert.IsType<ConflictException>(retry);
Assert.DoesNotContain("already authorized", conflict.Message, StringComparison.Ordinal);
Assert.Contains("Start quick connect again", conflict.Message, StringComparison.Ordinal);
Assert.Equal(
StatusCodes.Status409Conflict,
await StatusCodeAsync(async () =>
{
await _manager.AuthorizeRequest(Guid.NewGuid(), initiated.Code).ConfigureAwait(false);
return StatusCodes.Status200OK;
}));
}
/// <summary>
/// A request that really was authorized still says so, so the accurate message above is not just a
/// blanket replacement for the old one.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Authorize_OfAnAuthorizedRequest_IsConflictAndSaysAlreadyAuthorized()
{
var initiated = await _manager.TryConnect(AuthorizationInfo());
var userId = Guid.NewGuid();
_sessionManager
.Setup(manager => manager.AuthenticateDirect(It.IsAny<AuthenticationRequest>()))
.ReturnsAsync(new AuthenticationResult
{
AccessToken = "token-1",
ServerId = "server-1",
User = new UserDto { Id = userId, Name = "user", ServerId = "server-1" }
});
Assert.True(await _manager.AuthorizeRequest(userId, initiated.Code));
var retry = await Record.ExceptionAsync(() => _manager.AuthorizeRequest(userId, initiated.Code));
Assert.Equal("Request is already authorized", Assert.IsType<ConflictException>(retry).Message);
}
private static AuthorizationInfo AuthorizationInfo() => new AuthorizationInfo
{
Device = "Living Room TV",
DeviceId = Guid.NewGuid().ToString("N"),
Client = "Jellyfin Web",
Version = "1.0.0"
};
private static string NewSecret() => Guid.NewGuid().ToString("N");
private static int StatusOf(ActionResult<QuickConnectResult> result)
=> result.Result is IStatusCodeActionResult status
? status.StatusCode ?? StatusCodes.Status200OK
: StatusCodes.Status200OK;
private static async Task<int> StatusCodeAsync(Func<Task<int>> action)
{
var appPaths = new Mock<IServerApplicationPaths>();
appPaths.Setup(paths => paths.ProgramSystemPath).Returns("/program");
appPaths.Setup(paths => paths.ProgramDataPath).Returns("/data");
var configManager = new Mock<IServerConfigurationManager>();
configManager.Setup(manager => manager.ApplicationPaths).Returns(appPaths.Object);
var hostEnvironment = new Mock<IWebHostEnvironment>();
hostEnvironment.SetupGet(environment => environment.EnvironmentName).Returns(Environments.Production);
var context = new DefaultHttpContext();
context.Response.Body = new MemoryStream();
var middleware = new ExceptionMiddleware(
async _ =>
{
context.Response.StatusCode = await action().ConfigureAwait(false);
},
NullLogger<ExceptionMiddleware>.Instance,
configManager.Object,
hostEnvironment.Object);
await middleware.Invoke(context).ConfigureAwait(false);
return context.Response.StatusCode;
}
private async Task<string> InitiateAsync()
=> (await _manager.TryConnect(AuthorizationInfo()).ConfigureAwait(false)).Secret;
}
@@ -0,0 +1,150 @@
using System;
using System.IO;
using System.Threading.Tasks;
using Emby.Server.Implementations.QuickConnect;
using Jellyfin.Server.Extensions;
using Jellyfin.Server.Tests.HighAvailability;
using MediaBrowser.Common.Configuration;
using MediaBrowser.Controller.Authentication;
using MediaBrowser.Controller.QuickConnect;
using MediaBrowser.Model.QuickConnect;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging.Abstractions;
using Moq;
using StackExchange.Redis;
using Xunit;
namespace Jellyfin.Server.Tests.QuickConnect;
/// <summary>
/// Drives the whole configuration path a deployment uses: a bare
/// <c>Jellyfin__TranscodeStore__RedisConnectionString</c> environment variable, the server's own
/// configuration builder, the store registration, and a quick connect flow against a real valkey.
/// </summary>
[Trait("Category", "RequiresDocker")]
[Collection("JellyfinSectionConfiguration")]
public sealed class QuickConnectStoreWiringTests : IAsyncLifetime
{
private const string RedisConnectionStringVariable = "Jellyfin__TranscodeStore__RedisConnectionString";
private RedisTestServer _redis = null!;
private string _configDirectory = string.Empty;
/// <inheritdoc/>
public async ValueTask InitializeAsync()
{
_redis = await RedisTestServer.StartAsync().ConfigureAwait(false);
_configDirectory = Directory.CreateTempSubdirectory("jellyfin-quickconnect-wiring").FullName;
await File.WriteAllTextAsync(Path.Combine(_configDirectory, "logging.default.json"), "{}").ConfigureAwait(false);
}
/// <inheritdoc/>
public async ValueTask DisposeAsync()
{
Environment.SetEnvironmentVariable(RedisConnectionStringVariable, null);
if (_configDirectory.Length > 0)
{
Directory.Delete(_configDirectory, true);
}
await _redis.DisposeAsync().ConfigureAwait(false);
}
/// <summary>
/// The variable form deployments set selects the shared store, and that store really talks to valkey.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task ManifestStyleEnvironmentVariable_SelectsTheSharedStore()
{
Environment.SetEnvironmentVariable(RedisConnectionStringVariable, _redis.ConnectionString);
await using var provider = BuildProvider();
var store = provider.GetRequiredService<IQuickConnectStore>();
Assert.IsType<RedisQuickConnectStore>(store);
var request = NewRequest();
await store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), TestContext.Current.CancellationToken);
Assert.True(await store.TryClaimAuthorizationAsync(request.Secret, DateTime.UtcNow.AddMinutes(10), TestContext.Current.CancellationToken));
var redis = provider.GetRequiredService<IConnectionMultiplexer>();
Assert.True(await redis.GetDatabase().KeyExistsAsync("jellyfin:quickconnect:request:" + request.Secret));
}
/// <summary>
/// Without the variable the deployment is single-instance and gets the process-local store, which
/// runs a whole flow on its own.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task NoEnvironmentVariable_SelectsTheProcessLocalStore()
{
Environment.SetEnvironmentVariable(RedisConnectionStringVariable, null);
await using var provider = BuildProvider();
var store = provider.GetRequiredService<IQuickConnectStore>();
Assert.IsType<InMemoryQuickConnectStore>(store);
var cancellationToken = TestContext.Current.CancellationToken;
var request = NewRequest();
await store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), cancellationToken);
Assert.Equal(request.Secret, (await store.GetRequestByCodeAsync(request.Code, cancellationToken))?.Secret);
Assert.True(await store.TryClaimAuthorizationAsync(request.Secret, DateTime.UtcNow.AddMinutes(10), cancellationToken));
await store.SetAuthorizationAsync(
request.Secret,
new AuthenticationResult { AccessToken = "token-1" },
DateTime.UtcNow.AddMinutes(10),
cancellationToken);
Assert.Equal("token-1", (await store.GetAuthorizationAsync(request.Secret, cancellationToken))?.AccessToken);
Assert.Equal("token-1", (await store.GetAuthorizationAsync(request.Secret, cancellationToken))?.AccessToken);
}
/// <summary>
/// An unreachable connection string that connects eagerly, the default, throws while the store is
/// being built rather than on a read. A failed singleton factory is not cached, so every resolve
/// throws afresh, which is what lets the startup probe in <see cref="QuickConnectStartupTests"/>
/// retry the build and report an eager store's outage as the same operator-facing failure it reports
/// for the lazily connecting form.
/// </summary>
[Fact]
public void UnreachableRedisAtStartup_FailsClosed()
{
Environment.SetEnvironmentVariable(RedisConnectionStringVariable, "127.0.0.1:1,connectTimeout=250,connectRetry=0");
using var provider = BuildProvider();
Assert.ThrowsAny<RedisConnectionException>(() => provider.GetRequiredService<IQuickConnectStore>());
Assert.ThrowsAny<RedisConnectionException>(() => provider.GetRequiredService<IQuickConnectStore>());
}
private static QuickConnectResult NewRequest() => new QuickConnectResult(
Guid.NewGuid().ToString("N"),
Guid.NewGuid().ToString("N").Substring(0, 6),
DateTime.UtcNow,
"device-1",
"Living Room TV",
"Jellyfin Web",
"1.0.0");
private ServiceProvider BuildProvider()
{
var appPaths = new Mock<IApplicationPaths>();
appPaths.Setup(paths => paths.ConfigurationDirectoryPath).Returns(_configDirectory);
IConfiguration configuration = Jellyfin.Server.Program.CreateAppConfiguration(new StartupOptions(), appPaths.Object);
var services = new ServiceCollection();
services.AddLogging();
services.AddTranscodeSessionStore(configuration, NullLogger.Instance);
services.AddQuickConnectStore(configuration, NullLogger.Instance);
return services.BuildServiceProvider();
}
}
@@ -0,0 +1,302 @@
using System;
using System.Collections.Generic;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using Emby.Server.Implementations.QuickConnect;
using Jellyfin.Server.Tests.HighAvailability;
using MediaBrowser.Common.Extensions;
using MediaBrowser.Controller.Authentication;
using MediaBrowser.Controller.QuickConnect;
using MediaBrowser.Model.QuickConnect;
using Microsoft.Extensions.Logging.Abstractions;
using StackExchange.Redis;
using Xunit;
namespace Jellyfin.Server.Tests.QuickConnect;
/// <summary>
/// What a <see cref="RedisQuickConnectStore"/> does while its Redis is unreachable. Each instance talks
/// to the one real server through a proxy of its own, so an outage can be given to one instance and not
/// the others, and then taken back.
/// </summary>
[Trait("Category", "RequiresDocker")]
public sealed class RedisQuickConnectStoreDegradedTests : IAsyncLifetime
{
private readonly List<RedisFaultProxy> _proxies = new();
private readonly List<IConnectionMultiplexer> _connections = new();
private RedisTestServer _redis = null!;
private static CancellationToken CancellationToken => TestContext.Current.CancellationToken;
/// <inheritdoc/>
public async ValueTask InitializeAsync()
{
_redis = await RedisTestServer.StartAsync().ConfigureAwait(false);
}
/// <inheritdoc/>
public async ValueTask DisposeAsync()
{
foreach (var connection in _connections)
{
await connection.DisposeAsync().ConfigureAwait(false);
}
foreach (var proxy in _proxies)
{
await proxy.DisposeAsync().ConfigureAwait(false);
}
await _redis.DisposeAsync().ConfigureAwait(false);
}
/// <summary>
/// Every leg of the flow fails closed while Redis is unreachable. None of them may answer as though
/// Redis had said the request is unknown, because that tells a polling client its secret is invalid.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task EveryLeg_WhileRedisIsUnreachable_ReportsUnavailable()
{
var instance = await CreateInstanceAsync();
var request = NewRequest();
var expiresUtc = DateTime.UtcNow.AddMinutes(10);
instance.Proxy.Cut();
await AssertUnavailableAsync(() => instance.Store.SetRequestAsync(request, expiresUtc, CancellationToken));
await AssertUnavailableAsync(() => instance.Store.GetRequestBySecretAsync(request.Secret, CancellationToken));
await AssertUnavailableAsync(() => instance.Store.GetRequestByCodeAsync(request.Code, CancellationToken));
await AssertUnavailableAsync(() => instance.Store.TryClaimAuthorizationAsync(request.Secret, expiresUtc, CancellationToken));
await AssertUnavailableAsync(() => instance.Store.SetAuthorizationAsync(
request.Secret,
new AuthenticationResult { AccessToken = "token-1" },
expiresUtc,
CancellationToken));
await AssertUnavailableAsync(() => instance.Store.GetAuthorizationAsync(request.Secret, CancellationToken));
}
/// <summary>
/// The same three reads against a Redis that is answering report a genuine miss as a miss, which is
/// what makes an outage and an unknown secret tellable apart by the callers above.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task EveryRead_AgainstAHealthyRedis_ReportsAMissAsAMiss()
{
var instance = await CreateInstanceAsync();
var request = NewRequest();
Assert.Null(await instance.Store.GetRequestBySecretAsync(request.Secret, CancellationToken));
Assert.Null(await instance.Store.GetRequestByCodeAsync(request.Code, CancellationToken));
Assert.Null(await instance.Store.GetAuthorizationAsync(request.Secret, CancellationToken));
}
/// <summary>
/// A request is resolvable by its secret and by its code together or not at all: a failure part way
/// through storing it must not leave a code on the user's screen that resolves to nothing for the
/// whole ten minutes the poll keeps succeeding.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task PendingRequest_WhenRedisGoesAwayMidWrite_IsResolvableByBothKeysOrNeither()
{
var instance = await CreateInstanceAsync();
var request = NewRequest();
// Redis is taken away the instant after it has applied the write of the secret key, which is
// where a two round trip write loses the code key.
instance.Proxy.CutAfterForwarding("request:" + request.Secret);
await Record.ExceptionAsync(() => instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken));
// Read straight from the server: the instance's own connection is the one that was cut.
var direct = await _redis.ConnectAsync();
_connections.Add(direct);
var bySecret = await direct.GetDatabase().KeyExistsAsync("jellyfin:quickconnect:request:" + request.Secret);
var byCode = await direct.GetDatabase().KeyExistsAsync("jellyfin:quickconnect:code:" + request.Code);
Assert.Equal(bySecret, byCode);
}
/// <summary>
/// A malformed stored value is a fault of its own, not Redis being unavailable, so it is not reported
/// as either a miss or an outage.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task PendingRequest_ThatIsMalformedInRedis_SurfacesAsItsOwnFault()
{
var instance = await CreateInstanceAsync();
var request = NewRequest();
await instance.Connection.GetDatabase().StringSetAsync(
"jellyfin:quickconnect:request:" + request.Secret,
"{ not json",
TimeSpan.FromMinutes(10));
await Assert.ThrowsAsync<JsonException>(() => instance.Store.GetRequestBySecretAsync(request.Secret, CancellationToken));
}
/// <summary>
/// A store whose Redis comes back answers from Redis again, with nothing carried over from the outage.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Store_AfterAnOutage_WorksAgain()
{
var instance = await CreateInstanceAsync();
var request = NewRequest();
instance.Proxy.Cut();
await AssertUnavailableAsync(() => instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken));
await RestoreAsync(instance);
Assert.Null(await instance.Store.GetRequestBySecretAsync(request.Secret, CancellationToken));
await instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken);
Assert.Equal(request.Secret, (await instance.Store.GetRequestByCodeAsync(request.Code, CancellationToken))?.Secret);
}
/// <summary>
/// Two instances racing to authorize one request: exactly one of them may go on to mint a token.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Claim_RacedOnTwoInstances_SucceedsOnce()
{
var first = await CreateInstanceAsync();
var second = await CreateInstanceAsync();
for (var attempt = 0; attempt < 25; attempt++)
{
var request = NewRequest();
var expiresUtc = DateTime.UtcNow.AddMinutes(10);
await first.Store.SetRequestAsync(request, expiresUtc, CancellationToken);
var claims = await Task.WhenAll(
Task.Run(() => first.Store.TryClaimAuthorizationAsync(request.Secret, expiresUtc, CancellationToken), CancellationToken),
Task.Run(() => second.Store.TryClaimAuthorizationAsync(request.Secret, expiresUtc, CancellationToken), CancellationToken));
Assert.Single(claims, claimed => claimed);
}
}
/// <summary>
/// A request that is unknown, already claimed or already authorized cannot be claimed.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Claim_IsRefused_ForUnknownClaimedAndAuthorizedRequests()
{
var instance = await CreateInstanceAsync();
var expiresUtc = DateTime.UtcNow.AddMinutes(10);
Assert.False(await instance.Store.TryClaimAuthorizationAsync("unknown-secret", expiresUtc, CancellationToken));
var request = NewRequest();
await instance.Store.SetRequestAsync(request, expiresUtc, CancellationToken);
Assert.True(await instance.Store.TryClaimAuthorizationAsync(request.Secret, expiresUtc, CancellationToken));
Assert.False(await instance.Store.TryClaimAuthorizationAsync(request.Secret, expiresUtc, CancellationToken));
var authorized = NewRequest();
authorized.Authenticated = true;
await instance.Store.SetRequestAsync(authorized, expiresUtc, CancellationToken);
Assert.False(await instance.Store.TryClaimAuthorizationAsync(authorized.Secret, expiresUtc, CancellationToken));
}
/// <summary>
/// An authorization is read, not spent: the same secret exchanged again on the same instance returns
/// the same access token for as long as the authorization lives.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task Authorization_IsReadRepeatedly_WithoutBeingSpent()
{
var instance = await CreateInstanceAsync();
var request = NewRequest();
await instance.Store.SetAuthorizationAsync(
request.Secret,
new AuthenticationResult { AccessToken = "token-1" },
DateTime.UtcNow.AddMinutes(10),
CancellationToken);
Assert.Equal("token-1", (await instance.Store.GetAuthorizationAsync(request.Secret, CancellationToken))?.AccessToken);
Assert.Equal("token-1", (await instance.Store.GetAuthorizationAsync(request.Secret, CancellationToken))?.AccessToken);
}
/// <summary>
/// A write whose expiry has already passed is ignored, by the shared store and the process-local one
/// alike: a deployment must not get a different answer out of the two.
/// </summary>
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
[Fact]
public async Task ElapsedExpiry_IsIgnoredByBothStores()
{
var shared = (await CreateInstanceAsync()).Store;
var local = new InMemoryQuickConnectStore();
foreach (var store in new IQuickConnectStore[] { shared, local })
{
var request = NewRequest();
var elapsed = DateTime.UtcNow.AddSeconds(-1);
await store.SetRequestAsync(request, elapsed, CancellationToken);
await store.SetAuthorizationAsync(request.Secret, new AuthenticationResult { AccessToken = "token-1" }, elapsed, CancellationToken);
Assert.Null(await store.GetRequestBySecretAsync(request.Secret, CancellationToken));
Assert.Null(await store.GetRequestByCodeAsync(request.Code, CancellationToken));
Assert.Null(await store.GetAuthorizationAsync(request.Secret, CancellationToken));
}
}
private static QuickConnectResult NewRequest() => new QuickConnectResult(
Guid.NewGuid().ToString("N"),
Guid.NewGuid().ToString("N").Substring(0, 6),
DateTime.UtcNow,
"device-1",
"Living Room TV",
"Jellyfin Web",
"1.0.0");
private static async Task AssertUnavailableAsync(Func<Task> operation)
{
var exception = await Record.ExceptionAsync(operation);
Assert.NotNull(exception);
Assert.IsType<ServiceUnavailableException>(exception);
}
private static async Task RestoreAsync(Instance instance)
{
instance.Proxy.Restore();
for (var attempt = 1; ; attempt++)
{
try
{
await instance.Connection.GetDatabase().PingAsync();
return;
}
catch (Exception exception) when (exception is RedisException or TimeoutException && attempt < 60)
{
await Task.Delay(TimeSpan.FromMilliseconds(500), CancellationToken);
}
}
}
private async Task<Instance> CreateInstanceAsync()
{
var proxy = RedisFaultProxy.Start(_redis.ConnectionString);
_proxies.Add(proxy);
var connection = await ConnectionMultiplexer.ConnectAsync(proxy.ConnectionString).ConfigureAwait(false);
_connections.Add(connection);
return new Instance(proxy, connection, new RedisQuickConnectStore(connection, NullLogger<RedisQuickConnectStore>.Instance));
}
private sealed record Instance(RedisFaultProxy Proxy, IConnectionMultiplexer Connection, RedisQuickConnectStore Store);
}