Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions docs/Configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
40 changes: 39 additions & 1 deletion src/StackExchange.Redis/ConnectionMultiplexer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -2419,6 +2421,42 @@ internal void UpdateClusterRange(ClusterConfiguration configuration)
}
}

/// <summary>
/// Brings the primary/replica role of each known node into line with the <c>CLUSTER NODES</c> view.
/// </summary>
/// <remarks>
/// 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 <c>INFO replication</c> 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.
/// </remarks>
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;
}
}
}

/// <summary>
/// Applies the slot map from the <c>CLUSTER SLOTS</c> view, which supersedes
/// <see cref="UpdateClusterRange(ClusterConfiguration)"/> when the answering server supplied one.
Expand Down
2 changes: 1 addition & 1 deletion src/StackExchange.Redis/RedisServer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down
27 changes: 27 additions & 0 deletions src/StackExchange.Redis/ServerEndPoint.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down Expand Up @@ -1044,6 +1045,28 @@ internal void OnRepeatedConnectFailure(int consecutiveFailures)
/// <summary>Zero means "never", so a tick count that lands on it moves by one.</summary>
private static int NudgeFromZeroTicks(int ticks) => ticks == 0 ? 1 : ticks;

/// <summary>
/// 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
/// <c>PUBLISH</c> that clients have long used to announce a topology change would reach nobody.
/// </summary>
/// <remarks>
/// 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.
/// </remarks>
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
Expand All @@ -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)
{
Expand Down
10 changes: 10 additions & 0 deletions src/StackExchange.Redis/ServerSelectionStrategy.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
{
Expand Down
12 changes: 11 additions & 1 deletion tests/StackExchange.Redis.Tests/ClientKillTests.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using System.Collections.Generic;
using System.Linq;
using System.Net;
using System.Threading;
using System.Threading.Tasks;
Expand All @@ -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);
}

Expand Down
66 changes: 66 additions & 0 deletions tests/StackExchange.Redis.Tests/ClusterFailoverRolesUnitTests.cs
Original file line number Diff line number Diff line change
@@ -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;

/// <summary>
/// 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 <c>-MOVED</c> comes back
/// and nothing prompts a refresh, leaving writes failing until the next scheduled role check. See #3254.
/// </summary>
[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");
}
}
}
73 changes: 73 additions & 0 deletions tests/StackExchange.Redis.Tests/ConfigurationChannelUnitTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
using System;
using System.Threading.Tasks;
using Xunit;

namespace StackExchange.Redis.Tests;

/// <summary>
/// The configuration-change channel is how a client is told, by hand, that the topology moved: a <c>PUBLISH</c> 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.
/// </summary>
[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<EndPointEventArgs>(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<EndPointEventArgs>(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");
}
}
41 changes: 40 additions & 1 deletion toys/StackExchange.Redis.Server/RedisServer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,32 @@ public bool Migrate(int hashSlot, EndPoint to)
}
throw new KeyNotFoundException($"Source node not found for slot {hashSlot}");
}
/// <summary>
/// Promotes <paramref name="replica"/> in place of its primary, as a completed <c>CLUSTER FAILOVER</c> 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 <c>-MOVED</c>, or by asking.
/// </summary>
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<byte> key, EndPoint to) => Migrate(ServerSelectionStrategy.GetClusterSlot(key), to);
public bool Migrate(in RedisKey key, EndPoint to) => Migrate(GetHashSlot(key), to);

Expand Down Expand Up @@ -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)
{
Expand Down Expand Up @@ -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);
Expand Down
Loading