9e66708d87
Ownership is claimed with a Lua check-and-set keyed on the instance holding the websocket, routing prefers a live controller over a local copy, the session list deduplicates by owner, removal is ownership-checked, undelivered routed messages surface, single-session lookups stop scanning the keyspace and directory writes leave the request path bounded by a timeout.
96 lines
3.1 KiB
C#
96 lines
3.1 KiB
C#
using System;
|
|
using System.Text.Json;
|
|
using System.Threading;
|
|
using System.Threading.Tasks;
|
|
using Jellyfin.Extensions.Json;
|
|
using MediaBrowser.Controller.Session;
|
|
using Microsoft.Extensions.Logging;
|
|
using StackExchange.Redis;
|
|
|
|
namespace Emby.Server.Implementations.Session;
|
|
|
|
/// <summary>
|
|
/// A Redis pub/sub <see cref="IPodMessageBus"/>. Every instance subscribes to a channel named after
|
|
/// itself, which keeps addressed delivery working without the instances being routable to each other.
|
|
/// </summary>
|
|
public sealed class RedisPodMessageBus : IPodMessageBus
|
|
{
|
|
private const string ChannelPrefix = "jellyfin:pod:";
|
|
|
|
private static readonly TimeSpan _publishTimeout = TimeSpan.FromSeconds(5);
|
|
|
|
private static readonly JsonSerializerOptions _jsonOptions = JsonDefaults.Options;
|
|
|
|
private readonly ISubscriber _subscriber;
|
|
private readonly ILogger<RedisPodMessageBus> _logger;
|
|
|
|
/// <summary>
|
|
/// Initializes a new instance of the <see cref="RedisPodMessageBus"/> class.
|
|
/// </summary>
|
|
/// <param name="redis">The Redis connection multiplexer.</param>
|
|
/// <param name="logger">The logger.</param>
|
|
public RedisPodMessageBus(IConnectionMultiplexer redis, ILogger<RedisPodMessageBus> logger)
|
|
{
|
|
ArgumentNullException.ThrowIfNull(redis);
|
|
|
|
_subscriber = redis.GetSubscriber();
|
|
_logger = logger;
|
|
PodId = PodIdentity.Current;
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public string PodId { get; }
|
|
|
|
/// <inheritdoc />
|
|
public async Task<long> PublishAsync(string targetPod, PodMessage message, CancellationToken cancellationToken = default)
|
|
{
|
|
ArgumentException.ThrowIfNullOrEmpty(targetPod);
|
|
ArgumentNullException.ThrowIfNull(message);
|
|
|
|
message.OriginPod = PodId;
|
|
|
|
try
|
|
{
|
|
return await _subscriber.PublishAsync(
|
|
RedisChannel.Literal(ChannelPrefix + targetPod),
|
|
JsonSerializer.Serialize(message, _jsonOptions)).WaitAsync(_publishTimeout, cancellationToken).ConfigureAwait(false);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogWarning(ex, "Failed to send a {Kind} message to {TargetPod}.", message.Kind, targetPod);
|
|
return 0;
|
|
}
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public void Subscribe(Func<PodMessage, Task> 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<PodMessage, Task> handler, RedisValue value)
|
|
{
|
|
try
|
|
{
|
|
var message = JsonSerializer.Deserialize<PodMessage>(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.");
|
|
}
|
|
}
|
|
}
|