diff --git a/docs/Configuration.md b/docs/Configuration.md index 9297386a8..fa71d9ab5 100644 --- a/docs/Configuration.md +++ b/docs/Configuration.md @@ -266,6 +266,13 @@ Both options can be customized or disabled (set to `""`), via the `.Configuratio These settings are also used by the `IServer.MakeMaster()` method, which can set the tie-breaker in the database and broadcast the configuration change message. The configuration message can also be used separately to primary/replica changes simply to request all nodes to refresh their configurations, via the `ConnectionMultiplexer.PublishReconfigure` method. +The configuration channel is subscribed to on every server connection, under both RESP2 and RESP3. Like any other pub/sub channel, it has the `ChannelPrefix` (if one is configured) applied - both when subscribing to it and when the library publishes to it - so a client with `channelPrefix=app1-` listens on `app1-__Booksleeve_MasterChanged`. Two consequences follow: + +- If you announce a change yourself rather than through the library (for example `PUBLISH app1-__Booksleeve_MasterChanged "*"` from `redis-cli`), publish to the *prefixed* name, once for each distinct prefix in use; clients with a different prefix will not hear it. A publish to the unprefixed name is only heard by clients that have no prefix. +- If you use ACLs that restrict channels, the user needs access to the *prefixed* channel (for example `&app1-__Booksleeve_MasterChanged`), in addition to any channels the application itself uses. + +The configuration channel is a convenience for prompt notification. Where the server can report the change itself (for example the `-MOVED` responses from a redis cluster), the library will refresh without it. + ## Refreshing the topology after repeated connect failures An endpoint that refuses every connection is evidence that what the client believes about the deployment may diff --git a/src/StackExchange.Redis/ConnectionMultiplexer.cs b/src/StackExchange.Redis/ConnectionMultiplexer.cs index 4645c5fc5..e24c5d360 100644 --- a/src/StackExchange.Redis/ConnectionMultiplexer.cs +++ b/src/StackExchange.Redis/ConnectionMultiplexer.cs @@ -329,7 +329,9 @@ async Task BroadcastAsync(ServerSnapshot serverNodes) && ConfigurationChangedChannel != null && CommandMap.IsAvailable(RedisCommand.PUBLISH)) { - RedisValue channel = ConfigurationChangedChannel; + // as a channel, not a value: the prefix applies to this one too, or it would be published to a + // name that nobody (who subscribed to the prefixed one) is listening on + var channel = RedisChannel.Literal(ConfigurationChangedChannel); foreach (var node in serverNodes) { if (!node.IsConnected) continue; @@ -2419,6 +2421,42 @@ internal void UpdateClusterRange(ClusterConfiguration configuration) } } + /// + /// Brings the primary/replica role of each known node into line with the CLUSTER NODES view. + /// + /// + /// The slot map alone is not enough: after a failover it can be entirely correct while the role flags are + /// stale, and a write routed to a node still flagged as a replica is refused client-side. Roles otherwise + /// only refresh from the periodic INFO replication check, so this is what lets a topology refresh + /// actually repair a failover. Never creates a server - a node we do not hold has no flag to be stale. + /// + internal void ApplyClusterRoles(ClusterConfiguration configuration, ClusterTopology? topology) + { + foreach (var node in configuration.Nodes) + { + if (node.IgnoreFromClient || node.EndPoint is null) continue; + + var listed = topology?[node.NodeId]; + var server = TryResolveServerEndPoint(node.EndPoint); + if (server is null && listed is not null) + { + foreach (var identity in listed.Identities) + { + if ((server = TryResolveServerEndPoint(identity)) is not null) break; + } + } + + // SLOTS wins where it has a view, as it does for the slot map: taking the role from a different + // reply would let the two disagree about a node, which is the failure being repaired. NODES speaks + // only for the nodes SLOTS omits, i.e. those serving no slots + var isReplica = listed?.IsReplica ?? node.IsReplica; + if (server is not null && server.ServerType == ServerType.Cluster && server.IsReplica != isReplica) + { + server.IsReplica = isReplica; + } + } + } + /// /// Applies the slot map from the CLUSTER SLOTS view, which supersedes /// when the answering server supplied one. diff --git a/src/StackExchange.Redis/RedisServer.cs b/src/StackExchange.Redis/RedisServer.cs index 73049aa5f..94fb73065 100644 --- a/src/StackExchange.Redis/RedisServer.cs +++ b/src/StackExchange.Redis/RedisServer.cs @@ -689,7 +689,7 @@ internal static Message CreateReplicaOfMessage(ServerEndPoint sendMessageTo, End var channel = multiplexer.ConfigurationChangedChannel; if (channel != null && multiplexer.CommandMap.IsAvailable(RedisCommand.PUBLISH)) { - var msg = Message.Create(-1, CommandFlags.FireAndForget | CommandFlags.NoRedirect, RedisCommand.PUBLISH, (RedisValue)channel, RedisLiterals.Wildcard); + var msg = Message.Create(-1, CommandFlags.FireAndForget | CommandFlags.NoRedirect, RedisCommand.PUBLISH, RedisChannel.Literal(channel), RedisLiterals.Wildcard); msg.SetInternalCall(); return msg; } diff --git a/src/StackExchange.Redis/ServerEndPoint.cs b/src/StackExchange.Redis/ServerEndPoint.cs index e7c33c4c6..109fd5129 100644 --- a/src/StackExchange.Redis/ServerEndPoint.cs +++ b/src/StackExchange.Redis/ServerEndPoint.cs @@ -526,6 +526,7 @@ public void SetClusterConfiguration(ClusterConfiguration configuration) { Multiplexer.UpdateClusterRange(configuration); } + Multiplexer.ApplyClusterRoles(configuration, ClusterTopology); Multiplexer.Trace("Resolving genealogy..."); UpdateNodeRelations(configuration); Multiplexer.Trace("Cluster configured"); @@ -1044,6 +1045,28 @@ internal void OnRepeatedConnectFailure(int consecutiveFailures) /// Zero means "never", so a tick count that lands on it moves by one. private static int NudgeFromZeroTicks(int ticks) => ticks == 0 ? 1 : ticks; + /// + /// Subscribes to the configuration-change broadcast on a RESP3 connection. With RESP3 there is no + /// separate subscription connection, and so no subscription handshake - which is where RESP2 subscribes + /// to it. Left at that, the channel would silently have no subscriber, and the manual + /// PUBLISH that clients have long used to announce a topology change would reach nobody. + /// + /// + /// Done once the connection is known to be RESP3 rather than as part of the handshake: a connection + /// that fell back to RESP2 must not be put into subscriber mode. + /// + private void SubscribeToConfigurationChannel(PhysicalBridge bridge) + { + var channel = Multiplexer.ConfigurationChangedChannel; + if (channel is null || !SupportsSubscriptions || !Multiplexer.CommandMap.IsAvailable(RedisCommand.SUBSCRIBE)) return; + + var msg = Message.Create(-1, CommandFlags.FireAndForget, RedisCommand.SUBSCRIBE, RedisChannel.Literal(channel)); + msg.SetSource(ResultProcessor.TrackSubscriptions, null); +#pragma warning disable CS0618 // Type or member is obsolete + bridge.TryWriteSync(msg, isReplica); +#pragma warning restore CS0618 + } + internal void OnFullyEstablished(PhysicalConnection connection, string source) { try @@ -1070,6 +1093,10 @@ internal void OnFullyEstablished(PhysicalConnection connection, string source) // TracerProcessor which is executing this line inside a SetResultCore(). // Since we're issuing commands inside a SetResult path in a message, we'd create a deadlock by waiting. Multiplexer.EnsureSubscriptions(CommandFlags.FireAndForget); + if (isResp3 && bridge == interactive) + { + SubscribeToConfigurationChannel(bridge); + } } else if (SupportsSubscriptions && Multiplexer.RawConfig.Protocol > RedisProtocol.Resp2) { diff --git a/src/StackExchange.Redis/ServerSelectionStrategy.cs b/src/StackExchange.Redis/ServerSelectionStrategy.cs index 4ec7f9757..2c79aa804 100644 --- a/src/StackExchange.Redis/ServerSelectionStrategy.cs +++ b/src/StackExchange.Redis/ServerSelectionStrategy.cs @@ -210,6 +210,16 @@ public bool TryResend(int hashSlot, Message message, EndPoint endpoint, bool isM ServerEndPoint? server = multiplexer?.GetServerEndPoint(endpoint, provenance: ServerProvenance.Redirect); if (server != null) { + // a MOVED names the node that now owns the slot, and only a primary can own one. If we still + // believe it is a replica we have not caught up with a failover, and the resend below would be + // refused client-side as a write to a replica - with no MOVED to follow, so nothing else would + // correct it before the next scheduled role check. The next topology refresh remains the + // authority and can overrule this + if (isMoved && ServerType == ServerType.Cluster && server.IsReplica) + { + server.IsReplica = false; + } + bool retry = false; if ((message.Flags & CommandFlags.NoRedirect) == 0) { diff --git a/tests/StackExchange.Redis.Tests/ClientKillTests.cs b/tests/StackExchange.Redis.Tests/ClientKillTests.cs index f10f69ef6..8ac4f2c39 100644 --- a/tests/StackExchange.Redis.Tests/ClientKillTests.cs +++ b/tests/StackExchange.Redis.Tests/ClientKillTests.cs @@ -1,4 +1,5 @@ using System.Collections.Generic; +using System.Linq; using System.Net; using System.Threading; using System.Threading.Tasks; @@ -19,7 +20,16 @@ public async Task ClientKill() await using var conn = Create(allowAdmin: true, shared: false, backlogPolicy: BacklogPolicy.FailFast); var server = conn.GetServer(conn.GetEndPoints()[0]); - long result = server.ClientKill(id.AsInt64(), ClientType.Normal, null, true); + + // RESP3 has no separate subscription connection, so the interactive connection carries the + // configuration-channel subscription - and the server counts a subscribed client as pubsub, not normal. + var protocol = TestContext.Current.GetProtocol(); + var client = server.ClientList().Single(x => x.Id == id.AsInt64()); + Assert.Equal(protocol, client.Protocol); + var expectedType = protocol == RedisProtocol.Resp3 ? ClientType.PubSub : ClientType.Normal; + Assert.Equal(expectedType, client.ClientType); + + long result = server.ClientKill(id.AsInt64(), expectedType, null, true); Assert.Equal(1, result); } diff --git a/tests/StackExchange.Redis.Tests/ClusterFailoverRolesUnitTests.cs b/tests/StackExchange.Redis.Tests/ClusterFailoverRolesUnitTests.cs new file mode 100644 index 000000000..8497bf4ce --- /dev/null +++ b/tests/StackExchange.Redis.Tests/ClusterFailoverRolesUnitTests.cs @@ -0,0 +1,66 @@ +using System.Net; +using System.Threading.Tasks; +using StackExchange.Redis.Server; +using Xunit; +using static StackExchange.Redis.Server.RedisServer; + +namespace StackExchange.Redis.Tests; + +/// +/// After a failover the slot map can be right while the primary/replica role of the two nodes is stale, and a +/// write routed to a node still believed to be a replica is refused client-side - so no -MOVED comes back +/// and nothing prompts a refresh, leaving writes failing until the next scheduled role check. See #3254. +/// +[RunPerProtocol] +public class ClusterFailoverRolesUnitTests(ITestOutputHelper log) +{ + private static (InProcessTestServer Server, EndPoint Primary, EndPoint Replica) Create(ITestOutputHelper log) + { + var server = new InProcessTestServer(log) { ServerType = ServerType.Cluster }; + GetHost(server.DefaultEndPoint, out var port); + var replica = server.AddReplicaNode(new IPEndPoint(IPAddress.Loopback, port + 1), server.DefaultEndPoint); + return (server, server.DefaultEndPoint, replica); + } + + [Fact] + public async Task MovedToAPromotedReplicaIsFollowed() + { + var (server, primary, replica) = Create(log); + using (server) + { + await using var conn = await server.ConnectAsync(withPubSub: false); + var db = conn.GetDatabase(); + Assert.True(conn.GetServer(replica).IsReplica); + await db.StringSetAsync("before", "1"); + + server.Failover(replica); + + // the old primary answers -MOVED naming a node we still believe is a replica; that redirect is + // evidence it is not, and the resend must not be refused as a write to a replica + await db.StringSetAsync("after", "2"); + Assert.Equal("2", await db.StringGetAsync("after")); + Assert.False(conn.GetServer(replica).IsReplica); + } + } + + [Fact] + public async Task TopologyRefreshRepairsRolesAfterAFailover() + { + var (server, primary, replica) = Create(log); + using (server) + { + await using var conn = await server.ConnectAsync(withPubSub: false); + Assert.False(conn.GetServer(primary).IsReplica); + Assert.True(conn.GetServer(replica).IsReplica); + + server.Failover(replica); + await conn.ReconfigureAsync("test"); + + // previously the slot map followed but the flags did not, until the periodic INFO replication check + Assert.True(conn.GetServer(primary).IsReplica); + Assert.False(conn.GetServer(replica).IsReplica); + + await conn.GetDatabase().StringSetAsync("after", "2"); + } + } +} diff --git a/tests/StackExchange.Redis.Tests/ConfigurationChannelUnitTests.cs b/tests/StackExchange.Redis.Tests/ConfigurationChannelUnitTests.cs new file mode 100644 index 000000000..7045d6cc6 --- /dev/null +++ b/tests/StackExchange.Redis.Tests/ConfigurationChannelUnitTests.cs @@ -0,0 +1,73 @@ +using System; +using System.Threading.Tasks; +using Xunit; + +namespace StackExchange.Redis.Tests; + +/// +/// The configuration-change channel is how a client is told, by hand, that the topology moved: a PUBLISH to +/// it makes every subscribed client refresh. Under RESP2 the dedicated subscription connection subscribes to it as +/// part of its handshake; RESP3 has no such connection, and the subscription was never made. See #3254. +/// +[RunPerProtocol] +public class ConfigurationChannelUnitTests(ITestOutputHelper log) +{ + private const string Channel = "__Booksleeve_MasterChanged"; + + private static ConfigurationOptions Configure(InProcessTestServer server, bool prefix) + { + var config = server.GetClientConfig(); + config.ConfigurationChannel = Channel; // the test server turns this off by default + if (prefix) config.ChannelPrefix = RedisChannel.Literal("testuser-"); + return config; + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task PublishingToTheConfigurationChannelIsHeard(bool prefix) + { + using var server = new InProcessTestServer(log); + var config = Configure(server, prefix); + await using var conn = await ConnectionMultiplexer.ConnectAsync(config); + + var heard = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + conn.ConfigurationChangedBroadcast += (_, e) => heard.TrySetResult(e); + + // the subscription is made in the background after connecting, so nobody is listening for an instant: + // publish until someone is, rather than assume the order + var subscriber = conn.GetSubscriber(); + var deadline = DateTime.UtcNow.AddSeconds(10); + while (!heard.Task.IsCompleted && DateTime.UtcNow < deadline) + { + await subscriber.PublishAsync(RedisChannel.Literal(Channel), "*"); + await Task.WhenAny(heard.Task, Task.Delay(100)); + } + + Assert.True(heard.Task.IsCompleted, "the configuration channel has no subscriber"); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task TheLibrarysOwnBroadcastIsHeard(bool prefix) + { + // the channel is subscribed with the channel prefix applied, so it must be fired with it too - or a + // client that changes a server's role announces it to a channel that nobody is listening on + using var server = new InProcessTestServer(log); + await using var conn = await ConnectionMultiplexer.ConnectAsync(Configure(server, prefix)); + + var heard = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + conn.ConfigurationChangedBroadcast += (_, e) => heard.TrySetResult(e); + + var admin = conn.GetServer(server.DefaultEndPoint); + var deadline = DateTime.UtcNow.AddSeconds(10); + while (!heard.Task.IsCompleted && DateTime.UtcNow < deadline) + { + await admin.ReplicaOfAsync(null!); // broadcasts the change once it is made + await Task.WhenAny(heard.Task, Task.Delay(100)); + } + + Assert.True(heard.Task.IsCompleted, "the broadcast was not heard"); + } +} diff --git a/toys/StackExchange.Redis.Server/RedisServer.cs b/toys/StackExchange.Redis.Server/RedisServer.cs index 9e7cebd85..7fe8c46a9 100644 --- a/toys/StackExchange.Redis.Server/RedisServer.cs +++ b/toys/StackExchange.Redis.Server/RedisServer.cs @@ -123,6 +123,32 @@ public bool Migrate(int hashSlot, EndPoint to) } throw new KeyNotFoundException($"Source node not found for slot {hashSlot}"); } + /// + /// Promotes in place of its primary, as a completed CLUSTER FAILOVER does: + /// it takes the old primary's slots, the old primary becomes a replica of it, and any other replicas of + /// the old primary follow. Nothing is announced; a client finds out from a -MOVED, or by asking. + /// + public void Failover(EndPoint replica) + { + if (ServerType != ServerType.Cluster) throw new InvalidOperationException($"Server mode is {ServerType}"); + if (!TryGetNode(replica ?? throw new ArgumentNullException(nameof(replica)), out var promoted)) throw new KeyNotFoundException($"Node not found: {Format.ToString(replica)}"); + var demoted = GetPrimaryOf(promoted) ?? throw new InvalidOperationException($"Not a replica: {Format.ToString(replica)}"); + + var slots = demoted.Slots.ToArray(); + foreach (var pair in _nodes) + { + if (pair.Value.PrimaryId == demoted.Id) pair.Value.PrimaryId = promoted.Id; + } + promoted.PrimaryId = null; + promoted.Flags &= ~NodeFlags.Replica; + promoted.UpdateSlots(slots); + + demoted.PrimaryId = promoted.Id; + demoted.Flags |= NodeFlags.Replica; + demoted.UpdateSlots([]); + Log($"failover: {Format.ToString(replica)} promoted, {demoted.Host}:{demoted.Port} demoted"); + } + public bool Migrate(Span key, EndPoint to) => Migrate(ServerSelectionStrategy.GetClusterSlot(key), to); public bool Migrate(in RedisKey key, EndPoint to) => Migrate(GetHashSlot(key), to); @@ -713,6 +739,19 @@ protected virtual TypedRedisValue Watch(RedisClient client, in RedisRequest requ return TypedRedisValue.OK; } + // no read/write distinction is enforced (the topology decides who answers), but a client that has put a + // replica connection into read mode must be able to take it out again when that node is promoted + [RedisCommand(1)] + protected virtual TypedRedisValue Readonly(RedisClient client, in RedisRequest request) => TypedRedisValue.OK; + + [RedisCommand(1)] + protected virtual TypedRedisValue Readwrite(RedisClient client, in RedisRequest request) => TypedRedisValue.OK; + + // accepted and ignored: topology is whatever the test says it is, but a client that calls ReplicaOf expects + // an answer (and then broadcasts the change) + [RedisCommand(3)] + protected virtual TypedRedisValue Replicaof(RedisClient client, in RedisRequest request) => TypedRedisValue.OK; + [RedisCommand(1)] protected virtual TypedRedisValue Unwatch(RedisClient client, in RedisRequest request) { @@ -1226,7 +1265,7 @@ public override string ToString() private readonly RedisServer _server; public RedisServer Server => _server; - public NodeFlags Flags { get; } + public NodeFlags Flags { get; internal set; } public Node(RedisServer server, EndPoint endpoint, NodeFlags flags) { Host = GetHost(endpoint, out var port);