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)