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);