using System;
using System.Text.Json;
using System.Threading.Tasks;
using Jellyfin.Extensions.Json;
using MediaBrowser.Controller.Session;
using Microsoft.Extensions.Logging;
using StackExchange.Redis;
namespace Emby.Server.Implementations.Session;
///
/// A Redis pub/sub . Every instance subscribes to a channel named after
/// itself, which keeps addressed delivery working without the instances being routable to each other.
///
public sealed class RedisPodMessageBus : IPodMessageBus
{
private const string ChannelPrefix = "jellyfin:pod:";
private static readonly JsonSerializerOptions _jsonOptions = JsonDefaults.Options;
private readonly ISubscriber _subscriber;
private readonly ILogger _logger;
///
/// Initializes a new instance of the class.
///
/// The Redis connection multiplexer.
/// The logger.
public RedisPodMessageBus(IConnectionMultiplexer redis, ILogger logger)
{
ArgumentNullException.ThrowIfNull(redis);
_subscriber = redis.GetSubscriber();
_logger = logger;
PodId = PodIdentity.Current;
}
///
public string PodId { get; }
///
public void Publish(string targetPod, PodMessage message)
{
ArgumentException.ThrowIfNullOrEmpty(targetPod);
ArgumentNullException.ThrowIfNull(message);
message.OriginPod = PodId;
try
{
_subscriber.Publish(
RedisChannel.Literal(ChannelPrefix + targetPod),
JsonSerializer.Serialize(message, _jsonOptions),
CommandFlags.FireAndForget);
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Failed to send a {Kind} message to {TargetPod}.", message.Kind, targetPod);
}
}
///
public void Subscribe(Func handler)
{
ArgumentNullException.ThrowIfNull(handler);
try
{
_subscriber.Subscribe(RedisChannel.Literal(ChannelPrefix + PodId), (_, value) => Dispatch(handler, value));
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Failed to subscribe to {PodId}; messages routed here are dropped.", PodId);
}
}
private async void Dispatch(Func handler, RedisValue value)
{
try
{
var message = JsonSerializer.Deserialize(value.ToString(), _jsonOptions);
if (message is not null)
{
await handler(message).ConfigureAwait(false);
}
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Failed to handle a message routed to this instance.");
}
}
}