using System; using System.Collections.Concurrent; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using Jellyfin.Extensions.Json; using MediaBrowser.Controller.Session; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; 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. /// A request is answered on the sender's own channel, so the sender learns what the receiver did with /// it rather than only that something was subscribed. /// public sealed class RedisPodMessageBus : IPodMessageBus { private const string ChannelPrefix = "jellyfin:pod:"; private static readonly JsonSerializerOptions _jsonOptions = JsonDefaults.Options; private readonly ConcurrentDictionary> _pending = new(StringComparer.Ordinal); private readonly ISubscriber _subscriber; private readonly ILogger _logger; private readonly TimeSpan _timeout; private Func>? _handler; /// /// Initializes a new instance of the class. /// /// The Redis connection multiplexer. /// The session directory configuration options. /// The identity of this instance. /// The logger. public RedisPodMessageBus( IConnectionMultiplexer redis, IOptions options, string podId, ILogger logger) { ArgumentNullException.ThrowIfNull(redis); ArgumentNullException.ThrowIfNull(options); ArgumentException.ThrowIfNullOrEmpty(podId); _subscriber = redis.GetSubscriber(); _logger = logger; _timeout = TimeSpan.FromSeconds(Math.Max(1, options.Value.OperationTimeoutSeconds)); PodId = podId; try { _subscriber.Subscribe(RedisChannel.Literal(ChannelPrefix + PodId), (_, value) => Dispatch(value)); } catch (Exception ex) { _logger.LogWarning(ex, "Failed to subscribe to {PodId}; messages routed here are dropped.", PodId); } } /// public string PodId { get; } /// public async Task RequestAsync(string targetPod, PodMessage message, CancellationToken cancellationToken = default) { ArgumentException.ThrowIfNullOrEmpty(targetPod); ArgumentNullException.ThrowIfNull(message); message.OriginPod = PodId; message.CorrelationId = Guid.NewGuid().ToString("N"); var acknowledged = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); _pending[message.CorrelationId] = acknowledged; try { var subscribers = await _subscriber.PublishAsync( RedisChannel.Literal(ChannelPrefix + targetPod), JsonSerializer.Serialize(message, _jsonOptions)).WaitAsync(_timeout, cancellationToken).ConfigureAwait(false); if (subscribers == 0) { return false; } return await acknowledged.Task.WaitAsync(_timeout, cancellationToken).ConfigureAwait(false); } catch (TimeoutException) { _logger.LogWarning("Instance {TargetPod} did not acknowledge a {Kind} message within {Timeout}.", targetPod, message.Kind, _timeout); return false; } catch (Exception ex) when (ex is not OperationCanceledException) { _logger.LogWarning(ex, "Failed to send a {Kind} message to {TargetPod}.", message.Kind, targetPod); return false; } finally { _pending.TryRemove(message.CorrelationId, out _); } } /// public void Subscribe(Func> handler) { ArgumentNullException.ThrowIfNull(handler); _handler = handler; } private async void Dispatch(RedisValue value) { PodMessage? message = null; try { message = JsonSerializer.Deserialize(value.ToString(), _jsonOptions); if (message is null) { return; } if (string.Equals(message.Kind, PodMessage.AckKind, StringComparison.Ordinal)) { if (_pending.TryRemove(message.CorrelationId, out var acknowledged)) { acknowledged.TrySetResult(message.Handled); } return; } var handler = _handler; var handled = handler is not null && await handler(message).ConfigureAwait(false); await AcknowledgeAsync(message, handled).ConfigureAwait(false); } catch (Exception ex) { _logger.LogWarning(ex, "Failed to handle a message routed to this instance."); if (message is not null && !string.Equals(message.Kind, PodMessage.AckKind, StringComparison.Ordinal)) { await AcknowledgeAsync(message, false).ConfigureAwait(false); } } } private async Task AcknowledgeAsync(PodMessage message, bool handled) { if (string.IsNullOrEmpty(message.CorrelationId) || string.IsNullOrEmpty(message.OriginPod)) { return; } var ack = new PodMessage { Kind = PodMessage.AckKind, OriginPod = PodId, CorrelationId = message.CorrelationId, Handled = handled }; try { await _subscriber.PublishAsync( RedisChannel.Literal(ChannelPrefix + message.OriginPod), JsonSerializer.Serialize(ack, _jsonOptions)).WaitAsync(_timeout).ConfigureAwait(false); } catch (Exception ex) { _logger.LogWarning(ex, "Failed to acknowledge a {Kind} message to {OriginPod}.", message.Kind, message.OriginPod); } } }