propagate shared-config and library-option changes between instances
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.
This commit is contained in:
+73
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
+161
@@ -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;
|
||||
}
|
||||
}
|
||||
+104
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
+147
@@ -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;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user