From e1ca272d6ca39f7bf63cf98bcae1a35d33bdfbab Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Sat, 26 Sep 2026 14:28:45 +1000 Subject: [PATCH] fail quick connect closed on an unreachable valkey and restore the idempotent exchange --- .../QuickConnect/QuickConnectManager.cs | 20 +- .../QuickConnect/RedisQuickConnectStore.cs | 93 +++---- .../Controllers/QuickConnectController.cs | 8 + Jellyfin.Api/Controllers/UserController.cs | 4 + .../Middleware/ExceptionMiddleware.cs | 2 + ...ConnectStoreServiceCollectionExtensions.cs | 27 +- .../Extensions/ConflictException.cs | 36 +++ .../Extensions/ServiceUnavailableException.cs | 37 +++ .../QuickConnect/IQuickConnect.cs | 8 +- .../QuickConnect/IQuickConnectStore.cs | 21 +- .../QuickConnect/InMemoryQuickConnectStore.cs | 27 +- .../QuickConnect/QuickConnectManagerTests.cs | 20 +- .../HighAvailability/RedisFaultProxy.cs | 57 +++- .../QuickConnect/QuickConnectReplicaTests.cs | 29 +- .../QuickConnectStatusCodeTests.cs | 261 ++++++++++++++++++ .../QuickConnectStoreWiringTests.cs | 53 ++-- .../RedisQuickConnectStoreDegradedTests.cs | 206 +++++++------- 17 files changed, 664 insertions(+), 245 deletions(-) create mode 100644 MediaBrowser.Common/Extensions/ConflictException.cs create mode 100644 MediaBrowser.Common/Extensions/ServiceUnavailableException.cs create mode 100644 tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectStatusCodeTests.cs diff --git a/Emby.Server.Implementations/QuickConnect/QuickConnectManager.cs b/Emby.Server.Implementations/QuickConnect/QuickConnectManager.cs index 5ca0c7de73..2c619f4178 100644 --- a/Emby.Server.Implementations/QuickConnect/QuickConnectManager.cs +++ b/Emby.Server.Implementations/QuickConnect/QuickConnectManager.cs @@ -141,7 +141,7 @@ namespace Emby.Server.Implementations.QuickConnect if (result.Authenticated) { - throw new InvalidOperationException("Request is already authorized"); + throw new ConflictException("Request is already authorized"); } // Change the time on the request so it expires one minute into the future. It can't expire immediately as otherwise some clients wouldn't ever see that they have been authenticated. @@ -150,7 +150,7 @@ namespace Emby.Server.Implementations.QuickConnect // The guard above is a read on shared state, so it cannot settle a race between instances; the claim can. if (!await _store.TryClaimAuthorizationAsync(result.Secret, ExpiryOf(result)).ConfigureAwait(false)) { - throw new InvalidOperationException("Request is already authorized"); + throw await RefusedClaimAsync(result.Secret).ConfigureAwait(false); } var authenticationResult = await _sessionManager.AuthenticateDirect(new AuthenticationRequest @@ -177,7 +177,7 @@ namespace Emby.Server.Implementations.QuickConnect { AssertActive(); - var result = await _store.TryConsumeAuthorizationAsync(secret).ConfigureAwait(false); + var result = await _store.GetAuthorizationAsync(secret).ConfigureAwait(false); if (result is null) { throw new ResourceNotFoundException("Unable to find request"); @@ -188,6 +188,20 @@ namespace Emby.Server.Implementations.QuickConnect private static DateTime ExpiryOf(QuickConnectResult request) => request.DateAdded.AddMinutes(Timeout); + /// + /// Explains a refused claim. The claim outlives a failed mint on purpose, so it can mean either + /// that the request is authorized or that authorizing it did not finish; the two are told apart + /// by re-reading the request rather than reported as the same thing. + /// + private async Task RefusedClaimAsync(string secret) + { + var current = await _store.GetRequestBySecretAsync(secret).ConfigureAwait(false); + + return current?.Authenticated == true + ? new ConflictException("Request is already authorized") + : new ConflictException("Request is being authorized elsewhere, or an earlier attempt to authorize it did not complete. Start quick connect again for a new code."); + } + private string GenerateSecureRandom(int length = 32) { Span bytes = stackalloc byte[length]; diff --git a/Emby.Server.Implementations/QuickConnect/RedisQuickConnectStore.cs b/Emby.Server.Implementations/QuickConnect/RedisQuickConnectStore.cs index e3396be7ba..f536fdb55a 100644 --- a/Emby.Server.Implementations/QuickConnect/RedisQuickConnectStore.cs +++ b/Emby.Server.Implementations/QuickConnect/RedisQuickConnectStore.cs @@ -3,6 +3,7 @@ using System.Text.Json; using System.Threading; using System.Threading.Tasks; using Jellyfin.Extensions.Json; +using MediaBrowser.Common.Extensions; using MediaBrowser.Controller.Authentication; using MediaBrowser.Controller.QuickConnect; using MediaBrowser.Model.QuickConnect; @@ -13,20 +14,27 @@ namespace Emby.Server.Implementations.QuickConnect; /// /// A Redis-backed that lets the initiate, authorize and exchange legs -/// of a quick connect flow land on different instances. Expiry is the key TTL, an authorization is -/// claimed with a Lua check-and-set and consumed with GETDEL, so only one instance can ever mint -/// a given secret's access token and only one can ever hand it out. +/// of a quick connect flow land on different instances. Expiry is the key TTL and an authorization is +/// claimed with a Lua check-and-set, so only one instance can ever mint a given secret's access token. /// /// -/// A pending request survives an unreachable Redis through a process-local fallback, because a second -/// copy of it is harmless. An authorization has none: a second copy of it is a second access token, and -/// a write whose response timed out may well have been applied, so a transport failure on that path is -/// surfaced rather than degraded. +/// There is no local fallback: a call Redis did not answer is inconclusive, and reporting it as a miss +/// would tell a polling client its secret is invalid. Quick connect is unavailable for as long as Redis +/// is, which password login is not. /// public sealed class RedisQuickConnectStore : IQuickConnectStore { private const string KeyPrefix = "jellyfin:quickconnect:"; + /// + /// Lua script writing the two keys a request is resolvable by in one step, so it can never be + /// reachable by its secret while the code the user is reading off the screen resolves to nothing. + /// + private const string SetRequestScript = @" +redis.call('SET', KEYS[1], ARGV[1], 'PX', ARGV[3]) +redis.call('SET', KEYS[2], ARGV[2], 'PX', ARGV[3]) +return 1"; + /// /// Lua script for the atomic claim of the sole right to authorize a request: the request has to /// exist and not already be authorized, and the claim marker is taken with SET NX, so of two @@ -40,7 +48,6 @@ if redis.call('SET', KEYS[2], '1', 'NX', 'PX', ARGV[1]) then return 1 end return 0"; private readonly IDatabase _db; - private readonly InMemoryQuickConnectStore _fallback; private readonly ILogger _logger; /// @@ -53,41 +60,23 @@ return 0"; ArgumentNullException.ThrowIfNull(redis); _db = redis.GetDatabase(); - _fallback = new InMemoryQuickConnectStore(); _logger = logger; } /// public async Task GetRequestBySecretAsync(string secret, CancellationToken cancellationToken = default) { - RedisValue raw; - try - { - raw = await _db.StringGetAsync(RequestKey(secret)).ConfigureAwait(false); - } - catch (Exception ex) when (IsTransportFailure(ex)) - { - LogDegraded(ex); - return await _fallback.GetRequestBySecretAsync(secret, cancellationToken).ConfigureAwait(false); - } + var raw = await CallAsync(() => _db.StringGetAsync(RequestKey(secret))).ConfigureAwait(false); - // A miss is an answer rather than a transport failure, so the fallback is not consulted for it. + // Deserialization is outside the guard: a malformed stored value is a fault of its own, not Redis + // being unavailable. return raw.HasValue ? JsonSerializer.Deserialize(raw.ToString(), JsonDefaults.Options) : null; } /// public async Task GetRequestByCodeAsync(string code, CancellationToken cancellationToken = default) { - RedisValue secret; - try - { - secret = await _db.StringGetAsync(CodeKey(code)).ConfigureAwait(false); - } - catch (Exception ex) when (IsTransportFailure(ex)) - { - LogDegraded(ex); - return await _fallback.GetRequestByCodeAsync(code, cancellationToken).ConfigureAwait(false); - } + var secret = await CallAsync(() => _db.StringGetAsync(CodeKey(code))).ConfigureAwait(false); return secret.HasValue ? await GetRequestBySecretAsync(secret.ToString(), cancellationToken).ConfigureAwait(false) @@ -105,17 +94,11 @@ return 0"; return; } - try - { - var json = JsonSerializer.Serialize(request, JsonDefaults.Options); - await _db.StringSetAsync(RequestKey(request.Secret), json, ttl).ConfigureAwait(false); - await _db.StringSetAsync(CodeKey(request.Code), request.Secret, ttl).ConfigureAwait(false); - } - catch (Exception ex) when (IsTransportFailure(ex)) - { - LogDegraded(ex); - await _fallback.SetRequestAsync(request, expiresUtc, cancellationToken).ConfigureAwait(false); - } + var json = JsonSerializer.Serialize(request, JsonDefaults.Options); + await CallAsync(() => _db.ScriptEvaluateAsync( + SetRequestScript, + keys: new RedisKey[] { RequestKey(request.Secret), CodeKey(request.Code) }, + values: new RedisValue[] { json, request.Secret, (long)ttl.TotalMilliseconds })).ConfigureAwait(false); } /// @@ -127,10 +110,10 @@ return 0"; return false; } - var claimed = (long?)await _db.ScriptEvaluateAsync( + var claimed = (long?)await CallAsync(() => _db.ScriptEvaluateAsync( ClaimAuthorizationScript, keys: new RedisKey[] { RequestKey(secret), ClaimKey(secret) }, - values: new RedisValue[] { (long)ttl.TotalMilliseconds }).ConfigureAwait(false); + values: new RedisValue[] { (long)ttl.TotalMilliseconds })).ConfigureAwait(false); return claimed == 1; } @@ -145,23 +128,19 @@ return 0"; } var json = JsonSerializer.Serialize(authenticationResult, JsonDefaults.Options); - await _db.StringSetAsync(AuthorizationKey(secret), json, ttl).ConfigureAwait(false); + await CallAsync(() => _db.StringSetAsync(AuthorizationKey(secret), json, ttl)).ConfigureAwait(false); } /// - public async Task TryConsumeAuthorizationAsync(string secret, CancellationToken cancellationToken = default) + public async Task GetAuthorizationAsync(string secret, CancellationToken cancellationToken = default) { - var raw = await _db.StringGetDeleteAsync(AuthorizationKey(secret)).ConfigureAwait(false); + var raw = await CallAsync(() => _db.StringGetAsync(AuthorizationKey(secret))).ConfigureAwait(false); return raw.HasValue ? JsonSerializer.Deserialize(raw.ToString(), JsonDefaults.Options) : null; } - // Deliberately excludes a malformed stored value, which is a fault of its own rather than a reason - // to answer from this instance. - private static bool IsTransportFailure(Exception exception) => exception is RedisException or TimeoutException; - private static string RequestKey(string secret) => KeyPrefix + "request:" + secret; private static string CodeKey(string code) => KeyPrefix + "code:" + code; @@ -170,6 +149,16 @@ return 0"; private static string AuthorizationKey(string secret) => KeyPrefix + "auth:" + secret; - private void LogDegraded(Exception exception) - => _logger.LogWarning(exception, "Quick connect request state could not be shared through Redis; falling back to this instance only."); + private async Task CallAsync(Func> call) + { + try + { + return await call().ConfigureAwait(false); + } + catch (Exception exception) when (exception is RedisException or RedisCommandException or TimeoutException) + { + _logger.LogError(exception, "Quick connect state could not be reached in Redis."); + throw new ServiceUnavailableException("Quick connect is temporarily unavailable.", exception); + } + } } diff --git a/Jellyfin.Api/Controllers/QuickConnectController.cs b/Jellyfin.Api/Controllers/QuickConnectController.cs index eebda12ff3..54424deb06 100644 --- a/Jellyfin.Api/Controllers/QuickConnectController.cs +++ b/Jellyfin.Api/Controllers/QuickConnectController.cs @@ -50,10 +50,12 @@ public class QuickConnectController : BaseJellyfinApiController /// /// Quick connect request successfully created. /// Quick connect is not active on this server. + /// Quick connect state is unavailable. /// A with a secret and code for future use or an error message. [HttpPost("Initiate")] [ProducesResponseType(StatusCodes.Status200OK)] [ProducesResponseType(StatusCodes.Status401Unauthorized)] + [ProducesResponseType(StatusCodes.Status503ServiceUnavailable)] public async Task> InitiateQuickConnect() { try @@ -73,10 +75,12 @@ public class QuickConnectController : BaseJellyfinApiController /// Secret previously returned from the Initiate endpoint. /// Quick connect result returned. /// Unknown quick connect secret. + /// Quick connect state is unavailable. /// An updated . [HttpGet("Connect")] [ProducesResponseType(StatusCodes.Status200OK)] [ProducesResponseType(StatusCodes.Status404NotFound)] + [ProducesResponseType(StatusCodes.Status503ServiceUnavailable)] public async Task> GetQuickConnectState([FromQuery, Required] string secret) { try @@ -100,11 +104,15 @@ public class QuickConnectController : BaseJellyfinApiController /// The user the authorize. Access to the requested user is required. /// Quick connect result authorized successfully. /// Unknown user id. + /// Request is already authorized, or authorizing it did not complete. + /// Quick connect state is unavailable. /// Boolean indicating if the authorization was successful. [HttpPost("Authorize")] [Authorize] [ProducesResponseType(StatusCodes.Status200OK)] [ProducesResponseType(StatusCodes.Status403Forbidden)] + [ProducesResponseType(StatusCodes.Status409Conflict)] + [ProducesResponseType(StatusCodes.Status503ServiceUnavailable)] public async Task> AuthorizeQuickConnect([FromQuery, Required] string code, [FromQuery] Guid? userId = null) { userId = RequestHelpers.GetUserId(User, userId); diff --git a/Jellyfin.Api/Controllers/UserController.cs b/Jellyfin.Api/Controllers/UserController.cs index c8aeb616a4..48e473557b 100644 --- a/Jellyfin.Api/Controllers/UserController.cs +++ b/Jellyfin.Api/Controllers/UserController.cs @@ -241,9 +241,13 @@ public class UserController : BaseJellyfinApiController /// The request. /// User authenticated. /// Missing token. + /// Unknown or unauthorized quick connect secret. + /// Quick connect state is unavailable. /// A containing an with information about the new session. [HttpPost("AuthenticateWithQuickConnect")] [ProducesResponseType(StatusCodes.Status200OK)] + [ProducesResponseType(StatusCodes.Status404NotFound)] + [ProducesResponseType(StatusCodes.Status503ServiceUnavailable)] [Tags("Authentication")] public async Task> AuthenticateWithQuickConnect([FromBody, Required] QuickConnectDto request) { diff --git a/Jellyfin.Api/Middleware/ExceptionMiddleware.cs b/Jellyfin.Api/Middleware/ExceptionMiddleware.cs index acbb4877d4..8f5da04575 100644 --- a/Jellyfin.Api/Middleware/ExceptionMiddleware.cs +++ b/Jellyfin.Api/Middleware/ExceptionMiddleware.cs @@ -131,6 +131,8 @@ public class ExceptionMiddleware FileNotFoundException => StatusCodes.Status404NotFound, ResourceNotFoundException => StatusCodes.Status404NotFound, MethodNotAllowedException => StatusCodes.Status405MethodNotAllowed, + ConflictException => StatusCodes.Status409Conflict, + ServiceUnavailableException => StatusCodes.Status503ServiceUnavailable, _ => StatusCodes.Status500InternalServerError }; } diff --git a/Jellyfin.Server/Extensions/QuickConnectStoreServiceCollectionExtensions.cs b/Jellyfin.Server/Extensions/QuickConnectStoreServiceCollectionExtensions.cs index 36f270e66e..aa13634a9e 100644 --- a/Jellyfin.Server/Extensions/QuickConnectStoreServiceCollectionExtensions.cs +++ b/Jellyfin.Server/Extensions/QuickConnectStoreServiceCollectionExtensions.cs @@ -20,7 +20,9 @@ public static class QuickConnectStoreServiceCollectionExtensions /// /// /// The connection string is only set for a multi-instance deployment, which is the only shape where - /// the initiate, authorize and exchange legs of one flow can land on different instances. + /// the initiate, authorize and exchange legs of one flow can land on different instances. Set but + /// unreachable is a misconfigured deployment rather than a single-instance one, so it fails rather + /// than quietly handing out a store the other instances cannot see. /// /// The service collection. /// The configuration to read the Redis connection string from. @@ -48,25 +50,8 @@ public static class QuickConnectStoreServiceCollectionExtensions "Quick connect store: {Store}. Quick connect flows complete across any instance.", nameof(RedisQuickConnectStore)); - return serviceCollection.AddSingleton(sp => - { - try - { - return new RedisQuickConnectStore( - sp.GetRequiredService(), - sp.GetRequiredService>()); - } - catch (Exception ex) - { - // Fail open: an unreachable Redis degrades to the single-instance behaviour of a flow - // having to complete against one instance, rather than taking quick connect down. - sp.GetRequiredService>().LogError( - ex, - "Redis is configured but unavailable, so quick connect flows will not complete across instances. Check {Key}.", - TranscodeStoreOptions.RedisConnectionStringKey); - - return new InMemoryQuickConnectStore(); - } - }); + return serviceCollection.AddSingleton(sp => new RedisQuickConnectStore( + sp.GetRequiredService(), + sp.GetRequiredService>())); } } diff --git a/MediaBrowser.Common/Extensions/ConflictException.cs b/MediaBrowser.Common/Extensions/ConflictException.cs new file mode 100644 index 0000000000..433fc13643 --- /dev/null +++ b/MediaBrowser.Common/Extensions/ConflictException.cs @@ -0,0 +1,36 @@ +using System; + +namespace MediaBrowser.Common.Extensions +{ + /// + /// Thrown when the current state of a resource does not allow the requested operation. + /// + public class ConflictException : Exception + { + /// + /// Initializes a new instance of the class. + /// + public ConflictException() + { + } + + /// + /// Initializes a new instance of the class. + /// + /// The message. + public ConflictException(string message) + : base(message) + { + } + + /// + /// Initializes a new instance of the class. + /// + /// The message. + /// The exception that caused this one. + public ConflictException(string message, Exception innerException) + : base(message, innerException) + { + } + } +} diff --git a/MediaBrowser.Common/Extensions/ServiceUnavailableException.cs b/MediaBrowser.Common/Extensions/ServiceUnavailableException.cs new file mode 100644 index 0000000000..34fc8049d6 --- /dev/null +++ b/MediaBrowser.Common/Extensions/ServiceUnavailableException.cs @@ -0,0 +1,37 @@ +using System; + +namespace MediaBrowser.Common.Extensions +{ + /// + /// Thrown when an operation cannot be answered because a backing service is unreachable, rather than + /// because the thing it was asked about does not exist. + /// + public class ServiceUnavailableException : Exception + { + /// + /// Initializes a new instance of the class. + /// + public ServiceUnavailableException() + { + } + + /// + /// Initializes a new instance of the class. + /// + /// The message. + public ServiceUnavailableException(string message) + : base(message) + { + } + + /// + /// Initializes a new instance of the class. + /// + /// The message. + /// The exception that caused this one. + public ServiceUnavailableException(string message, Exception innerException) + : base(message, innerException) + { + } + } +} diff --git a/MediaBrowser.Controller/QuickConnect/IQuickConnect.cs b/MediaBrowser.Controller/QuickConnect/IQuickConnect.cs index b585c30dbd..eba710cd81 100644 --- a/MediaBrowser.Controller/QuickConnect/IQuickConnect.cs +++ b/MediaBrowser.Controller/QuickConnect/IQuickConnect.cs @@ -1,5 +1,6 @@ using System; using System.Threading.Tasks; +using MediaBrowser.Common.Extensions; using MediaBrowser.Controller.Authentication; using MediaBrowser.Controller.Net; using MediaBrowser.Model.QuickConnect; @@ -31,7 +32,9 @@ namespace MediaBrowser.Controller.QuickConnect Task CheckRequestStatus(string secret); /// - /// Authorizes a quick connect request to connect as the calling user. + /// Authorizes a quick connect request to connect as the calling user. A request can be authorized + /// once: a second attempt, including one following an attempt that failed part way, throws + /// and the user has to start quick connect again for a new code. /// /// User id. /// Identifying code for the request. @@ -39,7 +42,8 @@ namespace MediaBrowser.Controller.QuickConnect Task AuthorizeRequest(Guid userId, string code); /// - /// Gets the authorized request for the secret. + /// Gets the authorized request for the secret. The read does not consume the authorization, so a + /// client that retries the exchange gets the same access token until the authorization expires. /// /// The secret. /// The authentication result. diff --git a/MediaBrowser.Controller/QuickConnect/IQuickConnectStore.cs b/MediaBrowser.Controller/QuickConnect/IQuickConnectStore.cs index 4bc30d9e3d..fd2254eb17 100644 --- a/MediaBrowser.Controller/QuickConnect/IQuickConnectStore.cs +++ b/MediaBrowser.Controller/QuickConnect/IQuickConnectStore.cs @@ -1,6 +1,7 @@ using System; using System.Threading; using System.Threading.Tasks; +using MediaBrowser.Common.Extensions; using MediaBrowser.Controller.Authentication; using MediaBrowser.Model.QuickConnect; @@ -11,6 +12,10 @@ namespace MediaBrowser.Controller.QuickConnect; /// initiate, authorize and exchange - can each land on a different instance, so the state has to be /// reachable from all of them. /// +/// +/// A shared implementation that cannot reach its backend throws +/// rather than reporting a miss, because a miss tells a polling client its secret is invalid. +/// public interface IQuickConnectStore { /// @@ -30,7 +35,8 @@ public interface IQuickConnectStore Task GetRequestByCodeAsync(string code, CancellationToken cancellationToken = default); /// - /// Stores a new or updated request until . + /// Stores a new or updated request until , resolvable by both its secret + /// and its code or by neither. A request already past is not stored. /// /// The request to store. /// The instant the request stops being resolvable. @@ -41,12 +47,13 @@ public interface IQuickConnectStore /// /// Atomically claims the sole right to authorize the request behind , so /// that two instances racing on one code cannot both mint an access token. The claim is never - /// released: a mint that failed after writing its token would otherwise be retried into a second one. + /// released: a mint that failed after writing its token would otherwise be retried into a second one, + /// so a request whose authorization failed has to be started again. /// /// The request secret. /// The instant the claim lapses, after which the request can be authorized again. /// A cancellation token. - /// true when this caller may go on to authorize the request; false when it is unknown, expired, already authorized or being authorized elsewhere. + /// true when this caller may go on to authorize the request; false when it is unknown, expired, already authorized or claimed elsewhere. Task TryClaimAuthorizationAsync(string secret, DateTime expiresUtc, CancellationToken cancellationToken = default); /// @@ -60,11 +67,11 @@ public interface IQuickConnectStore Task SetAuthorizationAsync(string secret, AuthenticationResult authenticationResult, DateTime expiresUtc, CancellationToken cancellationToken = default); /// - /// Atomically takes the authentication for and removes it, so that two - /// instances racing on the same secret cannot both hand out an access token. + /// Reads the authentication for . The read does not consume it, so a client + /// that retries an exchange gets the same access token for as long as the authentication lives. /// /// The request secret. /// A cancellation token. - /// The authentication, or null when the secret is unknown, expired or already exchanged. - Task TryConsumeAuthorizationAsync(string secret, CancellationToken cancellationToken = default); + /// The authentication, or null when the secret is unknown or has expired. + Task GetAuthorizationAsync(string secret, CancellationToken cancellationToken = default); } diff --git a/MediaBrowser.Controller/QuickConnect/InMemoryQuickConnectStore.cs b/MediaBrowser.Controller/QuickConnect/InMemoryQuickConnectStore.cs index 2296449406..0eef358730 100644 --- a/MediaBrowser.Controller/QuickConnect/InMemoryQuickConnectStore.cs +++ b/MediaBrowser.Controller/QuickConnect/InMemoryQuickConnectStore.cs @@ -9,9 +9,9 @@ using MediaBrowser.Model.QuickConnect; namespace MediaBrowser.Controller.QuickConnect; /// -/// A process-local . It is the single-instance default, and the -/// fallback a shared store degrades to while its backend is unreachable, so quick connect keeps -/// working for clients whose three legs happen to land on one instance. +/// A process-local , and the single-instance default. Nothing it holds +/// is visible to another instance, so a deployment running more than one has to configure a shared +/// store instead. /// public sealed class InMemoryQuickConnectStore : IQuickConnectStore { @@ -41,7 +41,11 @@ public sealed class InMemoryQuickConnectStore : IQuickConnectStore ArgumentNullException.ThrowIfNull(request); Expire(); - _requests[request.Secret] = new Entry(expiresUtc, request); + if (expiresUtc > DateTime.UtcNow) + { + _requests[request.Secret] = new Entry(expiresUtc, request); + } + return Task.CompletedTask; } @@ -61,20 +65,19 @@ public sealed class InMemoryQuickConnectStore : IQuickConnectStore public Task SetAuthorizationAsync(string secret, AuthenticationResult authenticationResult, DateTime expiresUtc, CancellationToken cancellationToken = default) { Expire(); - _authorizations[secret] = new Entry(expiresUtc, authenticationResult); + if (expiresUtc > DateTime.UtcNow) + { + _authorizations[secret] = new Entry(expiresUtc, authenticationResult); + } + return Task.CompletedTask; } /// - public Task TryConsumeAuthorizationAsync(string secret, CancellationToken cancellationToken = default) + public Task GetAuthorizationAsync(string secret, CancellationToken cancellationToken = default) { Expire(); - if (!_authorizations.TryRemove(secret, out var entry) || entry.ExpiresUtc <= DateTime.UtcNow) - { - return Task.FromResult(null); - } - - return Task.FromResult(entry.Value); + return Task.FromResult(_authorizations.TryGetValue(secret, out var entry) ? entry.Value : null); } private void Expire() diff --git a/tests/Jellyfin.Server.Implementations.Tests/QuickConnect/QuickConnectManagerTests.cs b/tests/Jellyfin.Server.Implementations.Tests/QuickConnect/QuickConnectManagerTests.cs index 31eac9bba0..512a42ed77 100644 --- a/tests/Jellyfin.Server.Implementations.Tests/QuickConnect/QuickConnectManagerTests.cs +++ b/tests/Jellyfin.Server.Implementations.Tests/QuickConnect/QuickConnectManagerTests.cs @@ -154,14 +154,26 @@ namespace Jellyfin.Server.Implementations.Tests.QuickConnect } [Fact] - public async Task GetAuthorizedRequest_SecondExchange_ThrowsResourceNotFoundException() + public async Task GetAuthorizedRequest_SecondExchange_ReturnsTheSameResult() { _config.QuickConnectAvailable = true; var res = await _quickConnectManager.TryConnect(_quickConnectAuthInfo); await _quickConnectManager.AuthorizeRequest(Guid.Empty, res.Code); - Assert.NotNull(await _quickConnectManager.GetAuthorizedRequest(res.Secret)); - await Assert.ThrowsAsync(() => _quickConnectManager.GetAuthorizedRequest(res.Secret)); + var first = await _quickConnectManager.GetAuthorizedRequest(res.Secret); + var second = await _quickConnectManager.GetAuthorizedRequest(res.Secret); + + Assert.Same(first, second); + } + + [Fact] + public async Task AuthorizeRequest_OfAnAuthorizedRequest_ThrowsConflictException() + { + _config.QuickConnectAvailable = true; + var res = await _quickConnectManager.TryConnect(_quickConnectAuthInfo); + await _quickConnectManager.AuthorizeRequest(Guid.Empty, res.Code); + + await Assert.ThrowsAsync(() => _quickConnectManager.AuthorizeRequest(Guid.Empty, res.Code)); } private async Task AuthorizeAsync(string code) @@ -170,7 +182,7 @@ namespace Jellyfin.Server.Implementations.Tests.QuickConnect { return await _quickConnectManager.AuthorizeRequest(Guid.Empty, code).ConfigureAwait(false); } - catch (InvalidOperationException) + catch (ConflictException) { return false; } diff --git a/tests/Jellyfin.Server.Tests/HighAvailability/RedisFaultProxy.cs b/tests/Jellyfin.Server.Tests/HighAvailability/RedisFaultProxy.cs index 14d6c1a22c..6c4c6cb4a3 100644 --- a/tests/Jellyfin.Server.Tests/HighAvailability/RedisFaultProxy.cs +++ b/tests/Jellyfin.Server.Tests/HighAvailability/RedisFaultProxy.cs @@ -4,6 +4,7 @@ using System.Globalization; using System.IO; using System.Net; using System.Net.Sockets; +using System.Text; using System.Threading; using System.Threading.Tasks; using StackExchange.Redis; @@ -25,6 +26,7 @@ public sealed class RedisFaultProxy : IAsyncDisposable private readonly int _port; private volatile bool _cut; + private volatile byte[]? _cutAfterMarker; private RedisFaultProxy(TcpListener listener, int port, string targetHost, int targetPort) { @@ -74,11 +76,23 @@ public sealed class RedisFaultProxy : IAsyncDisposable DropLiveConnections(); } + /// + /// Arms a cut for the moment after a command containing has been forwarded + /// and answered, so a test can take Redis away between two round trips of one operation rather than + /// only before or after all of them. + /// + /// Text that identifies the command to cut after. + public void CutAfterForwarding(string marker) => _cutAfterMarker = Encoding.UTF8.GetBytes(marker); + /// /// Lets connections through again. Clients reconnect on their own schedule, so callers have to wait /// for the connection to come back rather than assume it already has. /// - public void Restore() => _cut = false; + public void Restore() + { + _cutAfterMarker = null; + _cut = false; + } /// public async ValueTask DisposeAsync() @@ -136,10 +150,17 @@ public sealed class RedisFaultProxy : IAsyncDisposable _live[client] = 0; _live[upstream] = 0; + // Registered first, then rechecked: a cut concurrent with this connect would otherwise drop + // the live connections before this pair joined them and leave it running through the outage. + if (_cut) + { + return; + } + var clientStream = client.GetStream(); var upstreamStream = upstream.GetStream(); await Task.WhenAny( - CopyAsync(clientStream, upstreamStream), + CopyFromClientAsync(clientStream, upstreamStream), CopyAsync(upstreamStream, clientStream)).ConfigureAwait(false); } catch (Exception exception) when (exception is IOException or SocketException or OperationCanceledException or ObjectDisposedException) @@ -157,6 +178,38 @@ public sealed class RedisFaultProxy : IAsyncDisposable } } + private async Task CopyFromClientAsync(NetworkStream from, NetworkStream to) + { + var buffer = new byte[16 * 1024]; + try + { + while (true) + { + var read = await from.ReadAsync(buffer, _cts.Token).ConfigureAwait(false); + if (read == 0) + { + return; + } + + await to.WriteAsync(buffer.AsMemory(0, read), _cts.Token).ConfigureAwait(false); + + var marker = _cutAfterMarker; + if (marker is not null && buffer.AsSpan(0, read).IndexOf(marker) >= 0) + { + _cutAfterMarker = null; + + // Long enough for the server to have applied the command that was just forwarded. + await Task.Delay(TimeSpan.FromMilliseconds(250), _cts.Token).ConfigureAwait(false); + Cut(); + return; + } + } + } + catch (Exception exception) when (exception is IOException or SocketException or OperationCanceledException or ObjectDisposedException) + { + } + } + private async Task CopyAsync(NetworkStream from, NetworkStream to) { try diff --git a/tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectReplicaTests.cs b/tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectReplicaTests.cs index 9d8c3d2d49..31341139cc 100644 --- a/tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectReplicaTests.cs +++ b/tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectReplicaTests.cs @@ -112,16 +112,15 @@ public sealed class QuickConnectReplicaTests : IAsyncLifetime } /// - /// A secret is single use across the whole deployment: two replicas racing to exchange it must not - /// both hand out an access token. One scheduling of one race settles nothing either way, so the race - /// is run repeatedly. + /// Exchanging a secret does not spend it: a client that retries, or whose retry lands on another + /// replica, gets the same access token back rather than a 404, and the device is minted once. /// /// A representing the asynchronous operation. [Fact] - public async Task Exchange_RacedOnTwoReplicas_SucceedsOnce() + public async Task Exchange_RepeatedOnTwoReplicas_ReturnsTheSameToken() { var cancellationToken = TestContext.Current.CancellationToken; - var connectionString = await _postgres.CreateDatabaseAsync("quickconnect_replica_race", cancellationToken); + var connectionString = await _postgres.CreateDatabaseAsync("quickconnect_replica_reexchange", cancellationToken); await using var dataSource = new NpgsqlDataSourceBuilder(connectionString).Build(); var user = await CreateSchemaWithUserAsync(dataSource, cancellationToken); @@ -130,19 +129,25 @@ public sealed class QuickConnectReplicaTests : IAsyncLifetime var replicaB = await CreateReplicaAsync(dataSource, user); var replicaC = await CreateReplicaAsync(dataSource, user); - for (var attempt = 0; attempt < 25; attempt++) + for (var attempt = 0; attempt < 10; attempt++) { - var initiated = await replicaA.Manager.TryConnect(AuthorizationInfoFor(attempt)); + var authorizationInfo = AuthorizationInfoFor(attempt); + var initiated = await replicaA.Manager.TryConnect(authorizationInfo); await replicaB.Manager.AuthorizeRequest(user.Id, initiated.Code); - var outcomes = await Task.WhenAll( + var exchanged = await Task.WhenAll( Task.Run(() => ExchangeAsync(replicaA.Manager, initiated.Secret), cancellationToken), Task.Run(() => ExchangeAsync(replicaC.Manager, initiated.Secret), cancellationToken)); - Assert.Single(outcomes, outcome => outcome is not null); + Assert.All(exchanged, outcome => Assert.NotNull(outcome)); + Assert.Equal(exchanged[0]!.AccessToken, exchanged[1]!.AccessToken); - // And it stays consumed for every later attempt, on any replica. - await Assert.ThrowsAsync(() => replicaB.Manager.GetAuthorizedRequest(initiated.Secret)); + // Still there afterwards, on a replica that has not exchanged it yet. + var later = await replicaB.Manager.GetAuthorizedRequest(initiated.Secret); + Assert.Equal(exchanged[0]!.AccessToken, later.AccessToken); + + var devices = await replicaA.Devices.GetDevices(new DeviceQuery { DeviceId = authorizationInfo.DeviceId }); + Assert.Equal(later.AccessToken, Assert.Single(devices.Items).AccessToken); } } @@ -258,7 +263,7 @@ public sealed class QuickConnectReplicaTests : IAsyncLifetime { return await manager.AuthorizeRequest(userId, code).ConfigureAwait(false); } - catch (InvalidOperationException) + catch (ConflictException) { return false; } diff --git a/tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectStatusCodeTests.cs b/tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectStatusCodeTests.cs new file mode 100644 index 0000000000..55d3284de7 --- /dev/null +++ b/tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectStatusCodeTests.cs @@ -0,0 +1,261 @@ +using System; +using System.IO; +using System.Threading.Tasks; +using Emby.Server.Implementations.QuickConnect; +using Jellyfin.Api.Controllers; +using Jellyfin.Api.Middleware; +using Jellyfin.Server.Tests.HighAvailability; +using MediaBrowser.Common.Extensions; +using MediaBrowser.Controller; +using MediaBrowser.Controller.Authentication; +using MediaBrowser.Controller.Configuration; +using MediaBrowser.Controller.Net; +using MediaBrowser.Controller.QuickConnect; +using MediaBrowser.Controller.Session; +using MediaBrowser.Model.Configuration; +using MediaBrowser.Model.Dto; +using MediaBrowser.Model.QuickConnect; +using Microsoft.AspNetCore.Hosting; +using Microsoft.AspNetCore.Http; +using Microsoft.AspNetCore.Mvc; +using Microsoft.AspNetCore.Mvc.Infrastructure; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging.Abstractions; +using Moq; +using StackExchange.Redis; +using Xunit; + +namespace Jellyfin.Server.Tests.QuickConnect; + +/// +/// The status a client actually sees. A missing key and an unreachable valkey are different answers and +/// must not collapse into one: telling a polling client its secret is unknown ends its flow, while 503 +/// tells it to keep trying. The exception the store really throws is run through the real exception +/// middleware, so the mapping is exercised rather than assumed. +/// +[Trait("Category", "RequiresDocker")] +public sealed class QuickConnectStatusCodeTests : IAsyncLifetime +{ + private readonly Mock _sessionManager = new(); + + private RedisTestServer _redis = null!; + private RedisFaultProxy _proxy = null!; + private IConnectionMultiplexer _connection = null!; + private QuickConnectManager _manager = null!; + private QuickConnectController _controller = null!; + + /// + public async ValueTask InitializeAsync() + { + _redis = await RedisTestServer.StartAsync().ConfigureAwait(false); + _proxy = RedisFaultProxy.Start(_redis.ConnectionString); + _connection = await ConnectionMultiplexer.ConnectAsync(_proxy.ConnectionString).ConfigureAwait(false); + + var configManager = new Mock(); + configManager.Setup(manager => manager.Configuration).Returns(new ServerConfiguration { QuickConnectAvailable = true }); + + _manager = new QuickConnectManager( + configManager.Object, + NullLogger.Instance, + _sessionManager.Object, + new RedisQuickConnectStore(_connection, NullLogger.Instance)); + + _controller = new QuickConnectController(_manager, Mock.Of()); + } + + /// + public async ValueTask DisposeAsync() + { + await _connection.DisposeAsync().ConfigureAwait(false); + await _proxy.DisposeAsync().ConfigureAwait(false); + await _redis.DisposeAsync().ConfigureAwait(false); + } + + /// + /// A secret valkey has never heard of is a 404, which is what ends a flow the user abandoned. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task Poll_UnknownSecret_IsNotFound() + { + Assert.Equal( + StatusCodes.Status404NotFound, + await StatusCodeAsync(async () => StatusOf(await _controller.GetQuickConnectState(NewSecret())))); + } + + /// + /// The same poll while valkey is unreachable is a 503. This is the bug the shared store is here to + /// avoid: a blip must not tell every polling client that its secret is invalid. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task Poll_WhileRedisIsUnreachable_IsServiceUnavailable() + { + var secret = await InitiateAsync(); + _proxy.Cut(); + + Assert.Equal( + StatusCodes.Status503ServiceUnavailable, + await StatusCodeAsync(async () => StatusOf(await _controller.GetQuickConnectState(secret)))); + } + + /// + /// The exchange leg tells the two apart the same way. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task Exchange_UnknownSecret_IsNotFound() + { + Assert.Equal( + StatusCodes.Status404NotFound, + await StatusCodeAsync(async () => + { + await _manager.GetAuthorizedRequest(NewSecret()).ConfigureAwait(false); + return StatusCodes.Status200OK; + })); + } + + /// + /// The exchange leg while valkey is unreachable. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task Exchange_WhileRedisIsUnreachable_IsServiceUnavailable() + { + var secret = await InitiateAsync(); + _proxy.Cut(); + + Assert.Equal( + StatusCodes.Status503ServiceUnavailable, + await StatusCodeAsync(async () => + { + await _manager.GetAuthorizedRequest(secret).ConfigureAwait(false); + return StatusCodes.Status200OK; + })); + } + + /// + /// The authorize leg while valkey is unreachable. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task Authorize_WhileRedisIsUnreachable_IsServiceUnavailable() + { + var initiated = await _manager.TryConnect(AuthorizationInfo()); + _proxy.Cut(); + + Assert.Equal( + StatusCodes.Status503ServiceUnavailable, + await StatusCodeAsync(async () => + { + await _manager.AuthorizeRequest(Guid.NewGuid(), initiated.Code).ConfigureAwait(false); + return StatusCodes.Status200OK; + })); + } + + /// + /// A mint that threw leaves its claim taken on purpose, because the write it failed on may have + /// landed. Retrying then has to say so and be a 409 the client can act on, not a 500 and not the + /// untrue claim that the request is already authorized. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task Authorize_AfterAMintThatFailed_IsConflictAndSaysToStartAgain() + { + var initiated = await _manager.TryConnect(AuthorizationInfo()); + + _sessionManager + .Setup(manager => manager.AuthenticateDirect(It.IsAny())) + .ThrowsAsync(new InvalidOperationException("mint failed")); + + await Assert.ThrowsAsync(() => _manager.AuthorizeRequest(Guid.NewGuid(), initiated.Code)); + + var retry = await Record.ExceptionAsync(() => _manager.AuthorizeRequest(Guid.NewGuid(), initiated.Code)); + + var conflict = Assert.IsType(retry); + Assert.DoesNotContain("already authorized", conflict.Message, StringComparison.Ordinal); + Assert.Contains("Start quick connect again", conflict.Message, StringComparison.Ordinal); + + Assert.Equal( + StatusCodes.Status409Conflict, + await StatusCodeAsync(async () => + { + await _manager.AuthorizeRequest(Guid.NewGuid(), initiated.Code).ConfigureAwait(false); + return StatusCodes.Status200OK; + })); + } + + /// + /// A request that really was authorized still says so, so the accurate message above is not just a + /// blanket replacement for the old one. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task Authorize_OfAnAuthorizedRequest_IsConflictAndSaysAlreadyAuthorized() + { + var initiated = await _manager.TryConnect(AuthorizationInfo()); + var userId = Guid.NewGuid(); + + _sessionManager + .Setup(manager => manager.AuthenticateDirect(It.IsAny())) + .ReturnsAsync(new AuthenticationResult + { + AccessToken = "token-1", + ServerId = "server-1", + User = new UserDto { Id = userId, Name = "user", ServerId = "server-1" } + }); + + Assert.True(await _manager.AuthorizeRequest(userId, initiated.Code)); + + var retry = await Record.ExceptionAsync(() => _manager.AuthorizeRequest(userId, initiated.Code)); + + Assert.Equal("Request is already authorized", Assert.IsType(retry).Message); + } + + private static AuthorizationInfo AuthorizationInfo() => new AuthorizationInfo + { + Device = "Living Room TV", + DeviceId = Guid.NewGuid().ToString("N"), + Client = "Jellyfin Web", + Version = "1.0.0" + }; + + private static string NewSecret() => Guid.NewGuid().ToString("N"); + + private static int StatusOf(ActionResult result) + => result.Result is IStatusCodeActionResult status + ? status.StatusCode ?? StatusCodes.Status200OK + : StatusCodes.Status200OK; + + private static async Task StatusCodeAsync(Func> action) + { + var appPaths = new Mock(); + appPaths.Setup(paths => paths.ProgramSystemPath).Returns("/program"); + appPaths.Setup(paths => paths.ProgramDataPath).Returns("/data"); + + var configManager = new Mock(); + configManager.Setup(manager => manager.ApplicationPaths).Returns(appPaths.Object); + + var hostEnvironment = new Mock(); + hostEnvironment.SetupGet(environment => environment.EnvironmentName).Returns(Environments.Production); + + var context = new DefaultHttpContext(); + context.Response.Body = new MemoryStream(); + + var middleware = new ExceptionMiddleware( + async _ => + { + context.Response.StatusCode = await action().ConfigureAwait(false); + }, + NullLogger.Instance, + configManager.Object, + hostEnvironment.Object); + + await middleware.Invoke(context).ConfigureAwait(false); + + return context.Response.StatusCode; + } + + private async Task InitiateAsync() + => (await _manager.TryConnect(AuthorizationInfo()).ConfigureAwait(false)).Secret; +} diff --git a/tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectStoreWiringTests.cs b/tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectStoreWiringTests.cs index fb58e45374..b492f44d98 100644 --- a/tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectStoreWiringTests.cs +++ b/tests/Jellyfin.Server.Tests/QuickConnect/QuickConnectStoreWiringTests.cs @@ -59,7 +59,7 @@ public sealed class QuickConnectStoreWiringTests : IAsyncLifetime [Fact] public async Task ManifestStyleEnvironmentVariable_SelectsTheSharedStore() { - Environment.SetEnvironmentVariable(RedisConnectionStringVariable, _redis.ConnectionString + ",abortConnect=false"); + Environment.SetEnvironmentVariable(RedisConnectionStringVariable, _redis.ConnectionString); await using var provider = BuildProvider(); @@ -76,45 +76,50 @@ public sealed class QuickConnectStoreWiringTests : IAsyncLifetime } /// - /// Without the variable the deployment is single-instance and gets the process-local store. - /// - [Fact] - public void NoEnvironmentVariable_SelectsTheProcessLocalStore() - { - Environment.SetEnvironmentVariable(RedisConnectionStringVariable, null); - - using var provider = BuildProvider(); - - Assert.IsType(provider.GetRequiredService()); - } - - /// - /// A configured but unreachable Redis degrades to the single-instance behaviour of a flow having to - /// complete against one instance, rather than taking quick connect down at startup. + /// Without the variable the deployment is single-instance and gets the process-local store, which + /// runs a whole flow on its own. /// /// A representing the asynchronous operation. [Fact] - public async Task UnreachableRedisAtStartup_DegradesToTheProcessLocalStore() + public async Task NoEnvironmentVariable_SelectsTheProcessLocalStore() { - Environment.SetEnvironmentVariable(RedisConnectionStringVariable, "127.0.0.1:1,connectTimeout=250,connectRetry=0"); + Environment.SetEnvironmentVariable(RedisConnectionStringVariable, null); await using var provider = BuildProvider(); var store = provider.GetRequiredService(); Assert.IsType(store); - // Quick connect still works, it just cannot span instances. + var cancellationToken = TestContext.Current.CancellationToken; var request = NewRequest(); - await store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), TestContext.Current.CancellationToken); - Assert.True(await store.TryClaimAuthorizationAsync(request.Secret, DateTime.UtcNow.AddMinutes(10), TestContext.Current.CancellationToken)); + await store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), cancellationToken); + + Assert.Equal(request.Secret, (await store.GetRequestByCodeAsync(request.Code, cancellationToken))?.Secret); + Assert.True(await store.TryClaimAuthorizationAsync(request.Secret, DateTime.UtcNow.AddMinutes(10), cancellationToken)); + await store.SetAuthorizationAsync( request.Secret, new AuthenticationResult { AccessToken = "token-1" }, DateTime.UtcNow.AddMinutes(10), - TestContext.Current.CancellationToken); + cancellationToken); - Assert.Equal("token-1", (await store.TryConsumeAuthorizationAsync(request.Secret, TestContext.Current.CancellationToken))?.AccessToken); - Assert.Null(await store.TryConsumeAuthorizationAsync(request.Secret, TestContext.Current.CancellationToken)); + Assert.Equal("token-1", (await store.GetAuthorizationAsync(request.Secret, cancellationToken))?.AccessToken); + Assert.Equal("token-1", (await store.GetAuthorizationAsync(request.Secret, cancellationToken))?.AccessToken); + } + + /// + /// A connection string that is set but unreachable is a misconfigured multi-instance deployment. It + /// fails rather than handing out a store the other instances cannot see, which would put quick + /// connect back on the cross-instance behaviour this configuration exists to fix. + /// + [Fact] + public void UnreachableRedisAtStartup_FailsClosed() + { + Environment.SetEnvironmentVariable(RedisConnectionStringVariable, "127.0.0.1:1,connectTimeout=250,connectRetry=0"); + + using var provider = BuildProvider(); + + Assert.ThrowsAny(() => provider.GetRequiredService()); } private static QuickConnectResult NewRequest() => new QuickConnectResult( diff --git a/tests/Jellyfin.Server.Tests/QuickConnect/RedisQuickConnectStoreDegradedTests.cs b/tests/Jellyfin.Server.Tests/QuickConnect/RedisQuickConnectStoreDegradedTests.cs index 3c07faeae7..bfd4999705 100644 --- a/tests/Jellyfin.Server.Tests/QuickConnect/RedisQuickConnectStoreDegradedTests.cs +++ b/tests/Jellyfin.Server.Tests/QuickConnect/RedisQuickConnectStoreDegradedTests.cs @@ -5,7 +5,9 @@ using System.Threading; using System.Threading.Tasks; using Emby.Server.Implementations.QuickConnect; using Jellyfin.Server.Tests.HighAvailability; +using MediaBrowser.Common.Extensions; using MediaBrowser.Controller.Authentication; +using MediaBrowser.Controller.QuickConnect; using MediaBrowser.Model.QuickConnect; using Microsoft.Extensions.Logging.Abstractions; using StackExchange.Redis; @@ -51,58 +53,83 @@ public sealed class RedisQuickConnectStoreDegradedTests : IAsyncLifetime } /// - /// A request stored while Redis is unreachable is still resolvable on the instance that stored it, - /// so a flow whose three legs happen to land on one instance keeps working through the outage. + /// Every leg of the flow fails closed while Redis is unreachable. None of them may answer as though + /// Redis had said the request is unknown, because that tells a polling client its secret is invalid. /// /// A representing the asynchronous operation. [Fact] - public async Task PendingRequest_SurvivesAnOutage_OnTheInstanceThatStoredIt() + public async Task EveryLeg_WhileRedisIsUnreachable_ReportsUnavailable() { var instance = await CreateInstanceAsync(); var request = NewRequest(); + var expiresUtc = DateTime.UtcNow.AddMinutes(10); instance.Proxy.Cut(); - await instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken); - Assert.Equal(request.Secret, (await instance.Store.GetRequestBySecretAsync(request.Secret, CancellationToken))?.Secret); - Assert.Equal(request.Secret, (await instance.Store.GetRequestByCodeAsync(request.Code, CancellationToken))?.Secret); + await AssertUnavailableAsync(() => instance.Store.SetRequestAsync(request, expiresUtc, CancellationToken)); + await AssertUnavailableAsync(() => instance.Store.GetRequestBySecretAsync(request.Secret, CancellationToken)); + await AssertUnavailableAsync(() => instance.Store.GetRequestByCodeAsync(request.Code, CancellationToken)); + await AssertUnavailableAsync(() => instance.Store.TryClaimAuthorizationAsync(request.Secret, expiresUtc, CancellationToken)); + await AssertUnavailableAsync(() => instance.Store.SetAuthorizationAsync( + request.Secret, + new AuthenticationResult { AccessToken = "token-1" }, + expiresUtc, + CancellationToken)); + await AssertUnavailableAsync(() => instance.Store.GetAuthorizationAsync(request.Secret, CancellationToken)); } /// - /// Once Redis answers again it is the only authority: a miss is a miss, not a reason to serve the - /// copy this instance kept while it was unreachable. + /// The same three reads against a Redis that is answering report a genuine miss as a miss, which is + /// what makes an outage and an unknown secret tellable apart by the callers above. /// /// A representing the asynchronous operation. [Fact] - public async Task PendingRequest_StoredDuringAnOutage_IsNotServedOnceRedisAnswersAgain() + public async Task EveryRead_AgainstAHealthyRedis_ReportsAMissAsAMiss() { var instance = await CreateInstanceAsync(); var request = NewRequest(); - instance.Proxy.Cut(); - await instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken); - Assert.NotNull(await instance.Store.GetRequestBySecretAsync(request.Secret, CancellationToken)); - - await RestoreAsync(instance); - Assert.Null(await instance.Store.GetRequestBySecretAsync(request.Secret, CancellationToken)); Assert.Null(await instance.Store.GetRequestByCodeAsync(request.Code, CancellationToken)); + Assert.Null(await instance.Store.GetAuthorizationAsync(request.Secret, CancellationToken)); } /// - /// A malformed stored value is a fault of its own, not a transport failure, so it is surfaced rather - /// than answered from the copy this instance happens to hold. + /// A request is resolvable by its secret and by its code together or not at all: a failure part way + /// through storing it must not leave a code on the user's screen that resolves to nothing for the + /// whole ten minutes the poll keeps succeeding. /// /// A representing the asynchronous operation. [Fact] - public async Task PendingRequest_ThatIsMalformedInRedis_SurfacesInsteadOfDegrading() + public async Task PendingRequest_WhenRedisGoesAwayMidWrite_IsResolvableByBothKeysOrNeither() { var instance = await CreateInstanceAsync(); var request = NewRequest(); - instance.Proxy.Cut(); - await instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken); - await RestoreAsync(instance); + // Redis is taken away the instant after it has applied the write of the secret key, which is + // where a two round trip write loses the code key. + instance.Proxy.CutAfterForwarding("request:" + request.Secret); + await Record.ExceptionAsync(() => instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken)); + + // Read straight from the server: the instance's own connection is the one that was cut. + var direct = await _redis.ConnectAsync(); + _connections.Add(direct); + var bySecret = await direct.GetDatabase().KeyExistsAsync("jellyfin:quickconnect:request:" + request.Secret); + var byCode = await direct.GetDatabase().KeyExistsAsync("jellyfin:quickconnect:code:" + request.Code); + + Assert.Equal(bySecret, byCode); + } + + /// + /// A malformed stored value is a fault of its own, not Redis being unavailable, so it is not reported + /// as either a miss or an outage. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task PendingRequest_ThatIsMalformedInRedis_SurfacesAsItsOwnFault() + { + var instance = await CreateInstanceAsync(); + var request = NewRequest(); await instance.Connection.GetDatabase().StringSetAsync( "jellyfin:quickconnect:request:" + request.Secret, @@ -113,102 +140,24 @@ public sealed class RedisQuickConnectStoreDegradedTests : IAsyncLifetime } /// - /// An authorization write that failed leaves nothing behind on the instance, because the response - /// that never arrived may still have been applied and a second copy of an authorization is a second - /// access token. + /// A store whose Redis comes back answers from Redis again, with nothing carried over from the outage. /// /// A representing the asynchronous operation. [Fact] - public async Task Authorization_ThatFailedToStore_LeavesNothingOnTheInstance() + public async Task Store_AfterAnOutage_WorksAgain() { var instance = await CreateInstanceAsync(); var request = NewRequest(); - await instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken); instance.Proxy.Cut(); - await AssertTransportFailureAsync(() => instance.Store.SetAuthorizationAsync( - request.Secret, - new AuthenticationResult { AccessToken = "token-1" }, - DateTime.UtcNow.AddMinutes(10), - CancellationToken)); + await AssertUnavailableAsync(() => instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken)); await RestoreAsync(instance); - Assert.Null(await instance.Store.TryConsumeAuthorizationAsync(request.Secret, CancellationToken)); - } + Assert.Null(await instance.Store.GetRequestBySecretAsync(request.Secret, CancellationToken)); - /// - /// The instance whose authorization write failed while the write landed anyway still hands the token - /// out exactly once, rather than once from Redis and again from a copy of its own. - /// - /// A representing the asynchronous operation. - [Fact] - public async Task Authorization_IsHandedOutOnce_EvenAfterAFailedWriteOnTheSameInstance() - { - var instance = await CreateInstanceAsync(); - var request = NewRequest(); - var authentication = new AuthenticationResult { AccessToken = "token-1" }; await instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken); - - instance.Proxy.Cut(); - await AssertTransportFailureAsync(() => instance.Store.SetAuthorizationAsync( - request.Secret, - authentication, - DateTime.UtcNow.AddMinutes(10), - CancellationToken)); - - await RestoreAsync(instance); - - // Stands in for that write having been applied before the response was lost. - await instance.Store.SetAuthorizationAsync(request.Secret, authentication, DateTime.UtcNow.AddMinutes(10), CancellationToken); - - Assert.Equal("token-1", (await instance.Store.TryConsumeAuthorizationAsync(request.Secret, CancellationToken))?.AccessToken); - Assert.Null(await instance.Store.TryConsumeAuthorizationAsync(request.Secret, CancellationToken)); - } - - /// - /// Exchanging during an outage fails loudly and spends nothing, so the token is still there to be - /// handed out once when Redis comes back. - /// - /// A representing the asynchronous operation. - [Fact] - public async Task Exchange_DuringAnOutage_SurfacesTheFailureAndLeavesTheTokenUnspent() - { - var instance = await CreateInstanceAsync(); - var request = NewRequest(); - await instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken); - await instance.Store.SetAuthorizationAsync( - request.Secret, - new AuthenticationResult { AccessToken = "token-1" }, - DateTime.UtcNow.AddMinutes(10), - CancellationToken); - - instance.Proxy.Cut(); - await AssertTransportFailureAsync(() => instance.Store.TryConsumeAuthorizationAsync(request.Secret, CancellationToken)); - - await RestoreAsync(instance); - - Assert.Equal("token-1", (await instance.Store.TryConsumeAuthorizationAsync(request.Secret, CancellationToken))?.AccessToken); - Assert.Null(await instance.Store.TryConsumeAuthorizationAsync(request.Secret, CancellationToken)); - } - - /// - /// Authorizing during an outage fails loudly rather than claiming locally, because a claim only this - /// instance knows about does not stop another one minting a second access token. - /// - /// A representing the asynchronous operation. - [Fact] - public async Task Claim_DuringAnOutage_SurfacesTheFailure() - { - var instance = await CreateInstanceAsync(); - var request = NewRequest(); - await instance.Store.SetRequestAsync(request, DateTime.UtcNow.AddMinutes(10), CancellationToken); - - instance.Proxy.Cut(); - await AssertTransportFailureAsync(() => instance.Store.TryClaimAuthorizationAsync( - request.Secret, - DateTime.UtcNow.AddMinutes(10), - CancellationToken)); + Assert.Equal(request.Secret, (await instance.Store.GetRequestByCodeAsync(request.Code, CancellationToken))?.Secret); } /// @@ -258,6 +207,51 @@ public sealed class RedisQuickConnectStoreDegradedTests : IAsyncLifetime Assert.False(await instance.Store.TryClaimAuthorizationAsync(authorized.Secret, expiresUtc, CancellationToken)); } + /// + /// An authorization is read, not spent: the same secret exchanged again on the same instance returns + /// the same access token for as long as the authorization lives. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task Authorization_IsReadRepeatedly_WithoutBeingSpent() + { + var instance = await CreateInstanceAsync(); + var request = NewRequest(); + await instance.Store.SetAuthorizationAsync( + request.Secret, + new AuthenticationResult { AccessToken = "token-1" }, + DateTime.UtcNow.AddMinutes(10), + CancellationToken); + + Assert.Equal("token-1", (await instance.Store.GetAuthorizationAsync(request.Secret, CancellationToken))?.AccessToken); + Assert.Equal("token-1", (await instance.Store.GetAuthorizationAsync(request.Secret, CancellationToken))?.AccessToken); + } + + /// + /// A write whose expiry has already passed is ignored, by the shared store and the process-local one + /// alike: a deployment must not get a different answer out of the two. + /// + /// A representing the asynchronous operation. + [Fact] + public async Task ElapsedExpiry_IsIgnoredByBothStores() + { + var shared = (await CreateInstanceAsync()).Store; + var local = new InMemoryQuickConnectStore(); + + foreach (var store in new IQuickConnectStore[] { shared, local }) + { + var request = NewRequest(); + var elapsed = DateTime.UtcNow.AddSeconds(-1); + + await store.SetRequestAsync(request, elapsed, CancellationToken); + await store.SetAuthorizationAsync(request.Secret, new AuthenticationResult { AccessToken = "token-1" }, elapsed, CancellationToken); + + Assert.Null(await store.GetRequestBySecretAsync(request.Secret, CancellationToken)); + Assert.Null(await store.GetRequestByCodeAsync(request.Code, CancellationToken)); + Assert.Null(await store.GetAuthorizationAsync(request.Secret, CancellationToken)); + } + } + private static QuickConnectResult NewRequest() => new QuickConnectResult( Guid.NewGuid().ToString("N"), Guid.NewGuid().ToString("N").Substring(0, 6), @@ -267,12 +261,12 @@ public sealed class RedisQuickConnectStoreDegradedTests : IAsyncLifetime "Jellyfin Web", "1.0.0"); - private static async Task AssertTransportFailureAsync(Func operation) + private static async Task AssertUnavailableAsync(Func operation) { var exception = await Record.ExceptionAsync(operation); Assert.NotNull(exception); - Assert.True(exception is RedisException or TimeoutException, exception.ToString()); + Assert.IsType(exception); } private static async Task RestoreAsync(Instance instance)