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; /// /// Round-trips through a real Redis, the transport two /// replicas actually use to tell each other that the shared configuration directory has changed. /// [Trait("Category", "RequiresDocker")] public sealed class RedisConfigurationInvalidationBusTests : IAsyncLifetime { private readonly RedisContainer _container; private IConnectionMultiplexer? _redis; /// /// Initializes a new instance of the class. /// public RedisConfigurationInvalidationBusTests() { _container = new RedisBuilder("redis:7-alpine").Build(); } /// public async ValueTask InitializeAsync() { await _container.StartAsync(); _redis = await ConnectionMultiplexer.ConnectAsync(_container.GetConnectionString()); } /// public async ValueTask DisposeAsync() { if (_redis is not null) { await _redis.DisposeAsync(); } await _container.DisposeAsync(); } /// /// A notice published by one replica reaches the other, carrying enough to invalidate one entry. /// /// A representing the asynchronous operation. [Fact] public async Task Publish_ReachesTheOtherInstance() { var received = new TaskCompletionSource(); 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); } /// /// 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. /// /// A representing the asynchronous operation. [Fact] public async Task Publish_IsNotDeliveredToThePublisher() { var ownNotice = new TaskCompletionSource(); var otherNotice = new TaskCompletionSource(); 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.Instance); } finally { Environment.SetEnvironmentVariable("JELLYFIN_INSTANCE_ID", null); } } }