Files
jellyfin-ha-src/Emby.Server.Implementations/Configuration/RedisConfigurationInvalidationBus.cs
unkin-agent b662ffa48f
ci/woodpecker/push/ci Pipeline was successful
ci/woodpecker/pr/ci Pipeline was successful
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.
2026-09-20 23:47:21 +10:00

97 lines
3.7 KiB
C#

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.");
}
}
}
}