800 lines
36 KiB
C#
800 lines
36 KiB
C#
using System;
|
|
using System.Collections.Concurrent;
|
|
using System.Collections.Generic;
|
|
using System.Globalization;
|
|
using System.Linq;
|
|
using System.Text.Json;
|
|
using System.Threading;
|
|
using System.Threading.Tasks;
|
|
using Jellyfin.Database.Implementations;
|
|
using Jellyfin.Database.Implementations.DbConfiguration;
|
|
using Jellyfin.Database.Implementations.Entities;
|
|
using Jellyfin.Database.Implementations.Locking;
|
|
using Jellyfin.Database.Providers.PostgreSQL;
|
|
using Jellyfin.Server.Implementations.Devices;
|
|
using Jellyfin.Server.Tests.Migrations;
|
|
using MediaBrowser.Common.Extensions;
|
|
using MediaBrowser.Controller;
|
|
using MediaBrowser.Controller.Configuration;
|
|
using MediaBrowser.Controller.Drawing;
|
|
using MediaBrowser.Controller.Dto;
|
|
using MediaBrowser.Controller.Events;
|
|
using MediaBrowser.Controller.Library;
|
|
using MediaBrowser.Controller.Session;
|
|
using MediaBrowser.Model.Configuration;
|
|
using MediaBrowser.Model.Dto;
|
|
using MediaBrowser.Model.Session;
|
|
using MediaBrowser.Model.SyncPlay;
|
|
using Microsoft.EntityFrameworkCore;
|
|
using Microsoft.Extensions.Hosting;
|
|
using Microsoft.Extensions.Logging;
|
|
using Microsoft.Extensions.Logging.Abstractions;
|
|
using Microsoft.Extensions.Options;
|
|
using Moq;
|
|
using Npgsql;
|
|
using StackExchange.Redis;
|
|
using Xunit;
|
|
using RedisPodMessageBus = Emby.Server.Implementations.Session.RedisPodMessageBus;
|
|
using RedisSessionDirectory = Emby.Server.Implementations.Session.RedisSessionDirectory;
|
|
using SessionManager = Emby.Server.Implementations.Session.SessionManager;
|
|
|
|
namespace Jellyfin.Server.Tests.HighAvailability;
|
|
|
|
/// <summary>
|
|
/// Two independently constructed <see cref="SessionManager"/> instances over one PostgreSQL database and
|
|
/// one valkey are the in-process stand-in for two replicas without sticky sessions: a session either of
|
|
/// them holds has to be visible to, and controllable from, the other.
|
|
/// </summary>
|
|
[Trait("Category", "RequiresDocker")]
|
|
public sealed class SessionDirectoryReplicaTests : IAsyncLifetime
|
|
{
|
|
private const string AppName = "Jellyfin Web";
|
|
private const string AppVersion = "1.0.0";
|
|
private const string DeviceName = "Living Room TV";
|
|
private const string RemoteEndPoint = "127.0.0.1";
|
|
|
|
private PostgreSqlTestServer _postgres = null!;
|
|
private RedisTestServer _redis = null!;
|
|
private NpgsqlDataSource _dataSource = null!;
|
|
private IConnectionMultiplexer _connection = null!;
|
|
private ISessionDirectory _directory = null!;
|
|
private User _user = null!;
|
|
private User _guest = null!;
|
|
|
|
/// <inheritdoc/>
|
|
public async ValueTask InitializeAsync()
|
|
{
|
|
_postgres = await PostgreSqlTestServer.StartAsync();
|
|
_redis = await RedisTestServer.StartAsync();
|
|
_connection = await _redis.ConnectAsync();
|
|
_directory = new RedisSessionDirectory(
|
|
_connection,
|
|
Options.Create(new SessionDirectoryOptions()),
|
|
NullLogger<RedisSessionDirectory>.Instance);
|
|
|
|
var connectionString = await _postgres.CreateDatabaseAsync("session_directory", CancellationToken.None);
|
|
_dataSource = new NpgsqlDataSourceBuilder(connectionString).Build();
|
|
|
|
var context = CreateContext(_dataSource);
|
|
await using (context.ConfigureAwait(false))
|
|
{
|
|
await context.Database.EnsureCreatedAsync(CancellationToken.None);
|
|
|
|
_user = new User("replica-user", "provider", "provider");
|
|
_guest = new User("replica-guest", "provider", "provider");
|
|
context.Users.Add(_user);
|
|
context.Users.Add(_guest);
|
|
await context.SaveChangesAsync(CancellationToken.None);
|
|
}
|
|
}
|
|
|
|
/// <inheritdoc/>
|
|
public async ValueTask DisposeAsync()
|
|
{
|
|
await _connection.DisposeAsync();
|
|
await _dataSource.DisposeAsync();
|
|
await _redis.DisposeAsync();
|
|
await _postgres.DisposeAsync();
|
|
}
|
|
|
|
/// <summary>
|
|
/// Half the active playback is invisible when the session list only reports what the replica serving
|
|
/// the request happens to hold, so a session registered on one replica has to appear on the other.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task SessionRegisteredOnOneReplica_IsListedByAnother()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
await using var replicaA = CreateReplica("pod-a");
|
|
await using var replicaB = CreateReplica("pod-b");
|
|
|
|
var session = await Request(replicaA, "device-listed");
|
|
|
|
var listedByB = await replicaB.GetSessions(_user.Id, null, null, null, false, cancellationToken);
|
|
var listedByA = await replicaA.GetSessions(_user.Id, null, null, null, false, cancellationToken);
|
|
|
|
Assert.Contains(listedByB, i => string.Equals(i.Id, session.Id, StringComparison.Ordinal));
|
|
Assert.Contains(listedByA, i => string.Equals(i.Id, session.Id, StringComparison.Ordinal));
|
|
|
|
// The session is reported once, not once per replica that can see it.
|
|
Assert.Single(listedByB, i => string.Equals(i.Id, session.Id, StringComparison.Ordinal));
|
|
}
|
|
|
|
/// <summary>
|
|
/// The deployment has no sticky sessions, so one device's requests land on either replica while its
|
|
/// websocket stays on one of them. Ownership has to follow the connection rather than the last
|
|
/// request served, or the directory names the wrong replica, the session list doubles up and remote
|
|
/// control is delivered to a replica with nothing to deliver it to.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task RequestsAlternatingBetweenReplicas_KeepOwnershipWithTheConnection()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
var options = new SessionDirectoryOptions { EntryTtlSeconds = 60, RefreshIntervalSeconds = 2 };
|
|
await using var replicaA = CreateReplica("pod-a", options);
|
|
await using var replicaB = CreateReplica("pod-b", options);
|
|
|
|
// The device is first seen by the replica that will not hold its websocket.
|
|
await Request(replicaB, "device-roaming");
|
|
|
|
var session = await Request(replicaA, "device-roaming");
|
|
var controller = new RecordingSessionController();
|
|
session.AddController(controller);
|
|
await replicaA.OnSessionControllerConnected(session);
|
|
|
|
// The load balancer keeps handing the device's requests to whichever replica it likes, and the
|
|
// replica without the websocket must never take the session from the one that has it.
|
|
for (var i = 0; i < 8; i++)
|
|
{
|
|
await Request(replicaB, "device-roaming");
|
|
await Task.Delay(250, cancellationToken);
|
|
|
|
var entry = await _directory.GetAsync(session.Id, cancellationToken);
|
|
Assert.NotNull(entry);
|
|
Assert.Equal("pod-a", entry.OwnerPod);
|
|
Assert.True(entry.HoldsConnection);
|
|
|
|
await Request(replicaA, "device-roaming");
|
|
await Task.Delay(50, cancellationToken);
|
|
}
|
|
|
|
var listedByA = await replicaA.GetSessions(_user.Id, null, null, null, false, cancellationToken);
|
|
var listedByB = await replicaB.GetSessions(_user.Id, null, null, null, false, cancellationToken);
|
|
|
|
Assert.Single(listedByA, i => string.Equals(i.Id, session.Id, StringComparison.Ordinal));
|
|
Assert.Single(listedByB, i => string.Equals(i.Id, session.Id, StringComparison.Ordinal));
|
|
|
|
await replicaB.SendMessageCommand(
|
|
string.Empty,
|
|
session.Id,
|
|
new MessageCommand { Header = "Header", Text = "Dinner is ready" },
|
|
cancellationToken);
|
|
|
|
var (messageType, data) = await controller.WaitForMessageAsync(cancellationToken);
|
|
Assert.Equal(SessionMessageType.GeneralCommand, messageType);
|
|
Assert.Contains("Dinner is ready", data, StringComparison.Ordinal);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Remote control and "send message to session" used to succeed and do nothing when the device is
|
|
/// connected to another replica; the message has to reach the connection wherever it is held.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task MessageSentOnOneReplica_ReachesTheConnectionHeldByAnother()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
await using var replicaA = CreateReplica("pod-a");
|
|
await using var replicaB = CreateReplica("pod-b");
|
|
|
|
var session = await Request(replicaA, "device-controlled");
|
|
var controller = new RecordingSessionController();
|
|
session.AddController(controller);
|
|
await replicaA.OnSessionControllerConnected(session);
|
|
|
|
await replicaB.SendMessageCommand(
|
|
string.Empty,
|
|
session.Id,
|
|
new MessageCommand { Header = "Header", Text = "Dinner is ready" },
|
|
cancellationToken);
|
|
|
|
var (messageType, data) = await controller.WaitForMessageAsync(cancellationToken);
|
|
Assert.Equal(SessionMessageType.GeneralCommand, messageType);
|
|
Assert.Contains("Dinner is ready", data, StringComparison.Ordinal);
|
|
}
|
|
|
|
/// <summary>
|
|
/// An entry outlives the replica that wrote it by up to its expiry, and a command routed into that
|
|
/// gap reaches nobody. Reporting it as delivered is the failure this directory exists to remove.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task MessageRoutedToADeadOwner_IsReportedAsUndelivered()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
await using var replicaA = CreateReplica("pod-a");
|
|
await using var replicaB = CreateReplica("pod-b");
|
|
|
|
var session = await Request(replicaA, "device-dead-owner");
|
|
session.AddController(new RecordingSessionController());
|
|
await replicaA.OnSessionControllerConnected(session);
|
|
|
|
var entry = await _directory.GetAsync(session.Id, cancellationToken);
|
|
Assert.NotNull(entry);
|
|
|
|
// A replica that is no longer listening, holding the entry until it expires.
|
|
entry.OwnerPod = "pod-gone";
|
|
Assert.True(await _directory.PublishAsync(entry, DateTime.UtcNow.Ticks, cancellationToken));
|
|
|
|
await Assert.ThrowsAsync<ResourceNotFoundException>(
|
|
() => replicaB.SendMessageCommand(
|
|
string.Empty,
|
|
session.Id,
|
|
new MessageCommand { Header = "Header", Text = "Dinner is ready" },
|
|
cancellationToken));
|
|
}
|
|
|
|
/// <summary>
|
|
/// The owner's entry is only as fresh as its last refresh, so a websocket that closes in between
|
|
/// leaves an entry claiming a connection that is gone. The command has to be reported undelivered,
|
|
/// which only the replica that would have written it to the socket can say.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task MessageRoutedToAnOwnerWhoseSocketDied_IsReportedAsUndelivered()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
await using var replicaA = CreateReplica("pod-a");
|
|
await using var replicaB = CreateReplica("pod-b");
|
|
|
|
var session = await Request(replicaA, "device-dead-socket");
|
|
var controller = new RecordingSessionController();
|
|
session.AddController(controller);
|
|
await replicaA.OnSessionControllerConnected(session);
|
|
|
|
// The socket closes. Nothing rewrites the entry: it still names pod-a and still says the
|
|
// connection is held, exactly as it does for the rest of the refresh interval.
|
|
controller.IsSessionActive = false;
|
|
|
|
var entry = await _directory.GetAsync(session.Id, cancellationToken);
|
|
Assert.NotNull(entry);
|
|
Assert.Equal("pod-a", entry.OwnerPod);
|
|
Assert.True(entry.HoldsConnection);
|
|
|
|
await Assert.ThrowsAsync<ResourceNotFoundException>(
|
|
() => replicaB.SendMessageCommand(
|
|
string.Empty,
|
|
session.Id,
|
|
new MessageCommand { Header = "Header", Text = "Dinner is ready" },
|
|
cancellationToken));
|
|
}
|
|
|
|
/// <summary>
|
|
/// A device reconnecting lands on either replica, so both can hold a live connection for the same
|
|
/// deterministic session id at once. Exactly one of them owns the entry, and it stays that one.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task BothReplicasHoldingAConnection_AgreeOnOneOwner()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
var options = new SessionDirectoryOptions { EntryTtlSeconds = 60, RefreshIntervalSeconds = 1 };
|
|
await using var replicaA = CreateReplica("pod-a", options);
|
|
await using var replicaB = CreateReplica("pod-b", options);
|
|
|
|
var sessionA = await Request(replicaA, "device-two-sockets");
|
|
sessionA.AddController(new RecordingSessionController());
|
|
await replicaA.OnSessionControllerConnected(sessionA);
|
|
|
|
var sessionB = await Request(replicaB, "device-two-sockets");
|
|
sessionB.AddController(new RecordingSessionController());
|
|
await replicaB.OnSessionControllerConnected(sessionB);
|
|
|
|
Assert.Equal(sessionA.Id, sessionB.Id);
|
|
|
|
// The later connection owns the session; both replicas keep republishing theirs.
|
|
for (var i = 0; i < 8; i++)
|
|
{
|
|
await Task.Delay(500, cancellationToken);
|
|
|
|
var entry = await _directory.GetAsync(sessionA.Id, cancellationToken);
|
|
Assert.NotNull(entry);
|
|
Assert.Equal("pod-b", entry.OwnerPod);
|
|
Assert.True(entry.HoldsConnection);
|
|
}
|
|
|
|
var listedByA = await replicaA.GetSessions(_user.Id, null, null, null, false, cancellationToken);
|
|
Assert.Single(listedByA, i => string.Equals(i.Id, sessionA.Id, StringComparison.Ordinal));
|
|
}
|
|
|
|
/// <summary>
|
|
/// A session with no websocket is claimed with a zero epoch by every replica that serves a request
|
|
/// for it. The first claim has to stand, or the listed session flips between two partial copies.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task TwoReplicasWithoutAConnection_DoNotTakeTheSessionFromEachOther()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
var options = new SessionDirectoryOptions { EntryTtlSeconds = 60, RefreshIntervalSeconds = 1 };
|
|
await using var replicaA = CreateReplica("pod-a", options);
|
|
await using var replicaB = CreateReplica("pod-b", options);
|
|
|
|
var session = await Request(replicaA, "device-no-socket");
|
|
await Request(replicaB, "device-no-socket");
|
|
|
|
var claimed = await _directory.GetAsync(session.Id, cancellationToken);
|
|
Assert.NotNull(claimed);
|
|
Assert.False(claimed.HoldsConnection);
|
|
|
|
var owner = claimed.OwnerPod;
|
|
|
|
Assert.False(await _directory.PublishAsync(
|
|
new SessionDirectoryEntry
|
|
{
|
|
OwnerPod = owner == "pod-a" ? "pod-b" : "pod-a",
|
|
HoldsConnection = false,
|
|
Session = claimed.Session
|
|
},
|
|
0,
|
|
cancellationToken));
|
|
|
|
for (var i = 0; i < 6; i++)
|
|
{
|
|
await Request(replicaB, "device-no-socket");
|
|
await Task.Delay(400, cancellationToken);
|
|
await Request(replicaA, "device-no-socket");
|
|
|
|
var entry = await _directory.GetAsync(session.Id, cancellationToken);
|
|
Assert.NotNull(entry);
|
|
Assert.Equal(owner, entry.OwnerPod);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Without sticky sessions a playback report lands on either replica while the websocket stays on
|
|
/// one. The report belongs to the replica everyone else is shown, so it is applied there and the
|
|
/// session reads as playing from every replica rather than idle on all of them.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task PlaybackReportedToTheNonOwner_IsVisibleFromBothReplicas()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
await using var replicaA = CreateReplica("pod-a");
|
|
await using var replicaB = CreateReplica("pod-b");
|
|
|
|
var session = await Request(replicaA, "device-playing");
|
|
session.AddController(new RecordingSessionController());
|
|
await replicaA.OnSessionControllerConnected(session);
|
|
|
|
// The load balancer hands the playback report to the replica without the websocket.
|
|
await Request(replicaB, "device-playing");
|
|
await replicaB.OnPlaybackStart(new PlaybackStartInfo
|
|
{
|
|
SessionId = session.Id,
|
|
Item = new BaseItemDto { Id = Guid.NewGuid(), Name = "Routed Movie" },
|
|
PositionTicks = 0
|
|
});
|
|
|
|
Assert.Equal("Routed Movie", session.NowPlayingItem?.Name);
|
|
|
|
var listedByA = await replicaA.GetSessions(_user.Id, null, null, null, false, cancellationToken);
|
|
var listedByB = await replicaB.GetSessions(_user.Id, null, null, null, false, cancellationToken);
|
|
|
|
Assert.Equal("Routed Movie", Single(listedByA, session.Id).NowPlayingItem?.Name);
|
|
Assert.Equal("Routed Movie", Single(listedByB, session.Id).NowPlayingItem?.Name);
|
|
|
|
await replicaB.OnPlaybackStopped(new PlaybackStopInfo { SessionId = session.Id, PositionTicks = 1 });
|
|
|
|
Assert.Null(session.NowPlayingItem);
|
|
Assert.Null(Single(await replicaB.GetSessions(_user.Id, null, null, null, false, cancellationToken), session.Id).NowPlayingItem);
|
|
}
|
|
|
|
/// <summary>
|
|
/// A report that cannot be handed to the owner is still applied here, so an unreachable owner is
|
|
/// never worse than the single-instance behaviour of keeping the state on the replica that served
|
|
/// the request.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task PlaybackReportedWithTheOwnerUnreachable_IsAppliedLocally()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
await using var replicaA = CreateReplica("pod-a");
|
|
await using var replicaB = CreateReplica("pod-b");
|
|
|
|
var session = await Request(replicaA, "device-orphaned");
|
|
session.AddController(new RecordingSessionController());
|
|
await replicaA.OnSessionControllerConnected(session);
|
|
|
|
var local = await Request(replicaB, "device-orphaned");
|
|
|
|
var entry = await _directory.GetAsync(session.Id, cancellationToken);
|
|
Assert.NotNull(entry);
|
|
entry.OwnerPod = "pod-gone";
|
|
Assert.True(await _directory.PublishAsync(entry, long.MaxValue, cancellationToken));
|
|
|
|
await replicaB.OnPlaybackStart(new PlaybackStartInfo
|
|
{
|
|
SessionId = session.Id,
|
|
Item = new BaseItemDto { Id = Guid.NewGuid(), Name = "Orphaned Movie" },
|
|
PositionTicks = 0
|
|
});
|
|
|
|
Assert.Equal("Orphaned Movie", local.NowPlayingItem?.Name);
|
|
}
|
|
|
|
/// <summary>
|
|
/// A directory that cannot be read says nothing about where a session is. Treating the failure as
|
|
/// "no such entry" hands the command to a local copy with no connection, which reports success and
|
|
/// delivers nothing.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task DirectoryReadFailingDuringRemoteControl_DoesNotSilentlyDoNothing()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
var failing = new FailableSessionDirectory(new RedisSessionDirectory(
|
|
_connection,
|
|
Options.Create(new SessionDirectoryOptions()),
|
|
NullLogger<RedisSessionDirectory>.Instance));
|
|
|
|
await using var replicaA = CreateReplica("pod-a");
|
|
await using var replicaB = CreateReplica("pod-b", failing);
|
|
|
|
var session = await Request(replicaA, "device-unreadable");
|
|
session.AddController(new RecordingSessionController());
|
|
await replicaA.OnSessionControllerConnected(session);
|
|
|
|
// The replica serving the request holds a copy of the session, and only a copy.
|
|
await Request(replicaB, "device-unreadable");
|
|
|
|
failing.FailReads = true;
|
|
|
|
await Assert.ThrowsAsync<RedisTimeoutException>(
|
|
() => replicaB.SendMessageCommand(
|
|
string.Empty,
|
|
session.Id,
|
|
new MessageCommand { Header = "Header", Text = "Dinner is ready" },
|
|
cancellationToken));
|
|
}
|
|
|
|
/// <summary>
|
|
/// Ownership is decided by comparing the two replicas' connection epochs, so the epochs cannot come
|
|
/// from the replicas' own clocks: a lagging clock would keep a genuinely newer connection from ever
|
|
/// taking the session. They are handed out per session by the shared store instead.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task ConnectionEpochs_AreHandedOutByTheStore()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
await using var replicaA = CreateReplica("pod-a");
|
|
|
|
var session = await Request(replicaA, "device-epoch");
|
|
session.AddController(new RecordingSessionController());
|
|
await replicaA.OnSessionControllerConnected(session);
|
|
|
|
var recorded = await ReadOwnerEpoch(session.Id, cancellationToken);
|
|
var allocated = await _directory.AllocateConnectionEpochAsync(session.Id, cancellationToken);
|
|
|
|
// A counter the store owns, not a reading of any replica's clock.
|
|
Assert.Equal(1, recorded);
|
|
Assert.Equal(2, allocated);
|
|
Assert.Equal(1, await _directory.AllocateConnectionEpochAsync(session.Id + "-other", cancellationToken));
|
|
}
|
|
|
|
/// <summary>
|
|
/// SyncPlay groups are still instance-local, so the replica serving the request has to notice that
|
|
/// its copy of the session has no connection. Holding a copy is not holding the connection, and a
|
|
/// command handed to a copy would be dropped without a word.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task SyncPlayCommandOnTheReplicaWithoutTheConnection_IsSkippedAndLogged()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
var logger = new CapturingLogger<SessionManager>();
|
|
await using var replicaA = CreateReplica("pod-a");
|
|
await using var replicaB = CreateReplica("pod-b", logger: logger);
|
|
|
|
var session = await Request(replicaA, "device-syncplay");
|
|
var controller = new RecordingSessionController();
|
|
session.AddController(controller);
|
|
await replicaA.OnSessionControllerConnected(session);
|
|
|
|
await Request(replicaB, "device-syncplay");
|
|
|
|
await replicaB.SendSyncPlayCommand(
|
|
session.Id,
|
|
new SendCommand(Guid.NewGuid(), Guid.NewGuid(), DateTime.UtcNow, SendCommandType.Pause, 0, DateTime.UtcNow),
|
|
cancellationToken);
|
|
|
|
Assert.Contains(logger.Messages, i => i.Contains("SyncPlay command", StringComparison.Ordinal) && i.Contains(session.Id, StringComparison.Ordinal));
|
|
Assert.False(controller.HasMessage);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Both replicas keep a copy of a session whose requests they have served, so the replica ending its
|
|
/// own copy must not erase the entry of the one still holding the connection.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task ReplicaEndingItsOwnCopy_LeavesTheOwnersEntryAlone()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
await using var replicaA = CreateReplica("pod-a");
|
|
await using var replicaB = CreateReplica("pod-b");
|
|
|
|
var session = await Request(replicaA, "device-shared-end");
|
|
session.AddController(new RecordingSessionController());
|
|
await replicaA.OnSessionControllerConnected(session);
|
|
|
|
await Request(replicaB, "device-shared-end");
|
|
await replicaB.ReportSessionEnded(session.Id);
|
|
|
|
var entry = await _directory.GetAsync(session.Id, cancellationToken);
|
|
Assert.NotNull(entry);
|
|
Assert.Equal("pod-a", entry.OwnerPod);
|
|
|
|
await replicaA.ReportSessionEnded(session.Id);
|
|
|
|
Assert.Null(await _directory.GetAsync(session.Id, cancellationToken));
|
|
}
|
|
|
|
/// <summary>
|
|
/// The session list now shows sessions from every replica, so an action offered against one of them
|
|
/// has to reach it rather than fail as missing on the replica serving the request.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task AdditionalUserAddedOnOneReplica_ReachesTheOwner()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
await using var replicaA = CreateReplica("pod-a");
|
|
await using var replicaB = CreateReplica("pod-b");
|
|
|
|
var session = await Request(replicaA, "device-additional-user");
|
|
session.AddController(new RecordingSessionController());
|
|
await replicaA.OnSessionControllerConnected(session);
|
|
|
|
await replicaB.AddAdditionalUser(string.Empty, session.Id, _guest.Id);
|
|
|
|
await WaitUntil(() => session.AdditionalUsers.Any(i => i.UserId.Equals(_guest.Id)), cancellationToken);
|
|
|
|
await replicaB.RemoveAdditionalUser(string.Empty, session.Id, _guest.Id);
|
|
|
|
await WaitUntil(() => !session.AdditionalUsers.Any(i => i.UserId.Equals(_guest.Id)), cancellationToken);
|
|
}
|
|
|
|
/// <summary>
|
|
/// A replica that dies stops refreshing its entries, and the sessions it held have to leave the
|
|
/// directory rather than linger in every other replica's session list forever.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task SessionsOfAReplicaThatStopsRefreshing_LeaveTheDirectory()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
|
|
// A never refreshes within the test, so it stands in for a replica that crashed.
|
|
await using var replicaA = CreateReplica("pod-a", new SessionDirectoryOptions { EntryTtlSeconds = 1, RefreshIntervalSeconds = 3600 });
|
|
await using var replicaB = CreateReplica("pod-b");
|
|
|
|
var session = await Request(replicaA, "device-expiring");
|
|
|
|
var listedWhileAlive = await replicaB.GetSessions(_user.Id, null, null, null, false, cancellationToken);
|
|
Assert.Contains(listedWhileAlive, i => string.Equals(i.Id, session.Id, StringComparison.Ordinal));
|
|
|
|
await Task.Delay(TimeSpan.FromSeconds(2), cancellationToken);
|
|
|
|
var listedAfterExpiry = await replicaB.GetSessions(_user.Id, null, null, null, false, cancellationToken);
|
|
Assert.DoesNotContain(listedAfterExpiry, i => string.Equals(i.Id, session.Id, StringComparison.Ordinal));
|
|
}
|
|
|
|
/// <summary>
|
|
/// A deployment without a shared store keeps the single-instance behaviour: nothing is published and
|
|
/// the other instance sees nothing.
|
|
/// </summary>
|
|
/// <returns>A <see cref="Task"/> representing the asynchronous operation.</returns>
|
|
[Fact]
|
|
public async Task WithoutADirectory_ReplicasOnlyReportTheirOwnSessions()
|
|
{
|
|
var cancellationToken = TestContext.Current.CancellationToken;
|
|
await using var replicaA = CreateReplica("pod-a", directory: NullSessionDirectory.Instance, bus: NullPodMessageBus.Instance);
|
|
await using var replicaB = CreateReplica("pod-b", directory: NullSessionDirectory.Instance, bus: NullPodMessageBus.Instance);
|
|
|
|
var session = await Request(replicaA, "device-local");
|
|
|
|
var listedByB = await replicaB.GetSessions(_user.Id, null, null, null, false, cancellationToken);
|
|
Assert.DoesNotContain(listedByB, i => string.Equals(i.Id, session.Id, StringComparison.Ordinal));
|
|
}
|
|
|
|
private async Task<long> ReadOwnerEpoch(string sessionId, CancellationToken cancellationToken)
|
|
{
|
|
var raw = await _connection.GetDatabase().StringGetAsync("jellyfin:sessionowner:" + sessionId).WaitAsync(cancellationToken);
|
|
|
|
return long.Parse(raw.ToString().Split('|')[0], CultureInfo.InvariantCulture);
|
|
}
|
|
|
|
private static SessionInfoDto Single(IReadOnlyList<SessionInfoDto> sessions, string sessionId)
|
|
=> Assert.Single(sessions, i => string.Equals(i.Id, sessionId, StringComparison.Ordinal));
|
|
|
|
private static async Task WaitUntil(Func<bool> condition, CancellationToken cancellationToken)
|
|
{
|
|
var deadline = DateTime.UtcNow.AddSeconds(10);
|
|
while (!condition())
|
|
{
|
|
Assert.True(DateTime.UtcNow < deadline, "The expected change never arrived.");
|
|
await Task.Delay(50, cancellationToken);
|
|
}
|
|
}
|
|
|
|
private static JellyfinDbContext CreateContext(NpgsqlDataSource dataSource)
|
|
{
|
|
var optionsBuilder = new DbContextOptionsBuilder<JellyfinDbContext>();
|
|
var provider = new PostgreSqlDatabaseProvider(dataSource);
|
|
provider.Initialise(optionsBuilder, new DatabaseConfigurationOptions { DatabaseType = "PostgreSQL" });
|
|
return new JellyfinDbContext(
|
|
optionsBuilder.Options,
|
|
NullLogger<JellyfinDbContext>.Instance,
|
|
provider,
|
|
new NoLockBehavior(NullLogger<NoLockBehavior>.Instance));
|
|
}
|
|
|
|
private Task<SessionInfo> Request(SessionManager replica, string deviceId)
|
|
=> replica.LogSessionActivity(AppName, AppVersion, deviceId, DeviceName, RemoteEndPoint, _user);
|
|
|
|
private SessionManager CreateReplica(string podId, SessionDirectoryOptions? options = null, ILogger<SessionManager>? logger = null)
|
|
{
|
|
options ??= new SessionDirectoryOptions { EntryTtlSeconds = 60, RefreshIntervalSeconds = 3600 };
|
|
|
|
var directory = new RedisSessionDirectory(
|
|
_connection,
|
|
Options.Create(options),
|
|
NullLogger<RedisSessionDirectory>.Instance);
|
|
|
|
return CreateReplica(podId, options, directory, CreateBus(podId, options), logger);
|
|
}
|
|
|
|
private SessionManager CreateReplica(string podId, ISessionDirectory directory)
|
|
=> CreateReplica(podId, new SessionDirectoryOptions(), directory, CreateBus(podId, new SessionDirectoryOptions()), null);
|
|
|
|
private SessionManager CreateReplica(string podId, ISessionDirectory directory, IPodMessageBus bus)
|
|
=> CreateReplica(podId, new SessionDirectoryOptions(), directory, bus, null);
|
|
|
|
private SessionManager CreateReplica(string podId, SessionDirectoryOptions options, ISessionDirectory directory, IPodMessageBus bus, ILogger<SessionManager>? logger)
|
|
{
|
|
var userManager = new Mock<IUserManager>();
|
|
userManager.Setup(i => i.GetUserById(_user.Id)).Returns(_user);
|
|
userManager.Setup(i => i.GetUserById(_guest.Id)).Returns(_guest);
|
|
|
|
var appHost = new Mock<IServerApplicationHost>();
|
|
appHost.SetupGet(i => i.SystemId).Returns("server-" + podId);
|
|
|
|
var configurationManager = new Mock<IServerConfigurationManager>();
|
|
configurationManager.SetupGet(i => i.Configuration).Returns(new ServerConfiguration());
|
|
|
|
return new SessionManager(
|
|
logger ?? NullLogger<SessionManager>.Instance,
|
|
Mock.Of<IEventManager>(),
|
|
Mock.Of<IUserDataManager>(),
|
|
configurationManager.Object,
|
|
Mock.Of<ILibraryManager>(),
|
|
userManager.Object,
|
|
Mock.Of<IMusicManager>(),
|
|
Mock.Of<IDtoService>(),
|
|
Mock.Of<IImageProcessor>(),
|
|
appHost.Object,
|
|
new DeviceManager(new DataSourceContextFactory(_dataSource), userManager.Object),
|
|
Mock.Of<IMediaSourceManager>(),
|
|
Mock.Of<IHostApplicationLifetime>(),
|
|
directory,
|
|
bus,
|
|
Options.Create(options));
|
|
}
|
|
|
|
private IPodMessageBus CreateBus(string podId, SessionDirectoryOptions options)
|
|
=> new RedisPodMessageBus(
|
|
_connection,
|
|
Options.Create(options),
|
|
podId,
|
|
NullLogger<RedisPodMessageBus>.Instance);
|
|
|
|
/// <summary>
|
|
/// Hands every replica its own context over the one shared database, the way the pooled factory does
|
|
/// in the server.
|
|
/// </summary>
|
|
private sealed class DataSourceContextFactory : IDbContextFactory<JellyfinDbContext>
|
|
{
|
|
private readonly NpgsqlDataSource _dataSource;
|
|
|
|
public DataSourceContextFactory(NpgsqlDataSource dataSource)
|
|
{
|
|
_dataSource = dataSource;
|
|
}
|
|
|
|
public JellyfinDbContext CreateDbContext() => CreateContext(_dataSource);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Keeps what a replica logged so a skipped route can be told apart from a silent drop.
|
|
/// </summary>
|
|
/// <typeparam name="T">The category the logger belongs to.</typeparam>
|
|
private sealed class CapturingLogger<T> : ILogger<T>
|
|
{
|
|
private readonly ConcurrentQueue<string> _messages = new();
|
|
|
|
public IEnumerable<string> Messages => _messages;
|
|
|
|
public IDisposable BeginScope<TState>(TState state)
|
|
where TState : notnull
|
|
=> NullLogger.Instance.BeginScope(state);
|
|
|
|
public bool IsEnabled(LogLevel logLevel) => true;
|
|
|
|
public void Log<TState>(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func<TState, Exception?, string> formatter)
|
|
=> _messages.Enqueue(formatter(state, exception));
|
|
}
|
|
|
|
/// <summary>
|
|
/// A directory whose reads can be made to fail the way an unreachable valkey does.
|
|
/// </summary>
|
|
private sealed class FailableSessionDirectory : ISessionDirectory
|
|
{
|
|
private readonly ISessionDirectory _inner;
|
|
|
|
public FailableSessionDirectory(ISessionDirectory inner)
|
|
{
|
|
_inner = inner;
|
|
}
|
|
|
|
public bool FailReads { get; set; }
|
|
|
|
public Task<long> AllocateConnectionEpochAsync(string sessionId, CancellationToken cancellationToken = default)
|
|
=> _inner.AllocateConnectionEpochAsync(sessionId, cancellationToken);
|
|
|
|
public Task<bool> PublishAsync(SessionDirectoryEntry entry, long connectionEpoch, CancellationToken cancellationToken = default)
|
|
=> _inner.PublishAsync(entry, connectionEpoch, cancellationToken);
|
|
|
|
public Task RemoveAsync(string sessionId, string ownerPod, CancellationToken cancellationToken = default)
|
|
=> _inner.RemoveAsync(sessionId, ownerPod, cancellationToken);
|
|
|
|
public Task<SessionDirectoryEntry?> GetAsync(string sessionId, CancellationToken cancellationToken = default)
|
|
=> FailReads
|
|
? throw new RedisTimeoutException("The session directory is unreachable.", CommandStatus.Unknown)
|
|
: _inner.GetAsync(sessionId, cancellationToken);
|
|
|
|
public Task<IReadOnlyList<SessionDirectoryEntry>> GetAllAsync(CancellationToken cancellationToken = default)
|
|
=> FailReads
|
|
? throw new RedisTimeoutException("The session directory is unreachable.", CommandStatus.Unknown)
|
|
: _inner.GetAllAsync(cancellationToken);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Stands in for the websocket the owning replica holds.
|
|
/// </summary>
|
|
private sealed class RecordingSessionController : ISessionController
|
|
{
|
|
private readonly TaskCompletionSource<(SessionMessageType MessageType, string Data)> _received = new();
|
|
|
|
public bool IsSessionActive { get; set; } = true;
|
|
|
|
public bool SupportsMediaControl => true;
|
|
|
|
public bool HasMessage => _received.Task.IsCompleted;
|
|
|
|
public Task SendMessage<T>(SessionMessageType name, Guid messageId, T data, CancellationToken cancellationToken)
|
|
{
|
|
_received.TrySetResult((name, JsonSerializer.Serialize(data)));
|
|
return Task.CompletedTask;
|
|
}
|
|
|
|
public Task<(SessionMessageType MessageType, string Data)> WaitForMessageAsync(CancellationToken cancellationToken)
|
|
=> _received.Task.WaitAsync(TimeSpan.FromSeconds(10), cancellationToken);
|
|
}
|
|
}
|