Skip to content
This repository was archived by the owner on Dec 24, 2022. It is now read-only.
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
33 changes: 25 additions & 8 deletions src/ServiceStack.Redis/RedisSentinel.cs
Original file line number Diff line number Diff line change
Expand Up @@ -52,27 +52,37 @@ public IRedisClientsManager Setup()

private void GetValidSentinel()
{
while (this.clientManager == null && failures < RedisSentinel.MaxFailures)
while (this.clientManager == null && ShouldRetry())
{
try
{
worker = GetNextSentinel();
clientManager = worker.GetClientManager();
worker.BeginListeningForConfigurationChanges();
this.worker = GetNextSentinel();
this.clientManager = worker.GetClientManager();
this.worker.BeginListeningForConfigurationChanges();
}
catch (RedisException)
{
if (worker != null)
if (this.worker != null)
{
worker.SentinelError -= Worker_SentinelError;
worker.Dispose();
this.worker.SentinelError -= Worker_SentinelError;
this.worker.Dispose();
}

failures++;
this.failures++;
}
}
}

/// <summary>
/// Check if GetValidSentinel should try the next sentinel server
/// </summary>
/// <returns></returns>
/// <remarks>This will be true if the failures is less than either RedisSentinel.MaxFailures or the # of sentinels, whatever is greater</remarks>
private bool ShouldRetry()
{
return this.failures < Math.Max(RedisSentinel.MaxFailures, this.sentinels.Count);
}

private RedisSentinelWorker GetNextSentinel()
{
sentinelIndex++;
Expand All @@ -88,15 +98,22 @@ private RedisSentinelWorker GetNextSentinel()
return sentinelWorker;
}

/// <summary>
/// Raised if there is an error from a sentinel worker
/// </summary>
/// <param name="sender"></param>
/// <param name="e"></param>
private void Worker_SentinelError(object sender, EventArgs e)
{
var worker = sender as RedisSentinelWorker;

if (worker != null)
{
// dispose the worker
worker.SentinelError -= Worker_SentinelError;
worker.Dispose();

// get a new worker and start looking for more changes
this.worker = GetNextSentinel();
this.worker.BeginListeningForConfigurationChanges();
}
Expand Down
142 changes: 50 additions & 92 deletions src/ServiceStack.Redis/RedisSentinelWorker.cs
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,38 @@ private void ConfigureRedisFromSentinel()
}
}

private Dictionary<string, string> ParseDataArray(object[] items)
{
Dictionary<string, string> data = new Dictionary<string, string>();
bool isKey = false;
string key = null;
string value = null;

foreach (var item in items)
{
if (item is byte[])
{
isKey = !isKey;

if (isKey)
{
key = Encoding.UTF8.GetString((byte[])item);
}
else
{
value = Encoding.UTF8.GetString((byte[])item);

if (!data.ContainsKey(key))
{
data.Add(key, value);
}
}
}
}

return data;
}

/// <summary>
/// Takes output from sentinel slaves command and converts into a list of servers
/// </summary>
Expand All @@ -113,69 +145,26 @@ private void ConfigureRedisFromSentinel()
private IEnumerable<string> ConvertSlaveArrayToList(object[] slaves)
{
var servers = new List<string>();
bool fetchIP = false;
bool fetchPort = false;
bool fetchFlags = false;
string ip = null;
string port = null;
string value = null;
string flags = null;

foreach (var slave in slaves.OfType<object[]>())
{
fetchIP = false;
fetchPort = false;
ip = null;
port = null;
var data = ParseDataArray(slave);

foreach (var item in slave)
{
if (item is byte[])
{
value = Encoding.UTF8.GetString((byte[])item);
if (value == "ip")
{
fetchIP = true;
continue;
}
else if (value == "port")
{
fetchPort = true;
continue;
}
else if (value == "flags")
{
fetchFlags = true;
continue;
}
else if (fetchIP)
{
ip = value;

if (ip == "127.0.0.1")
{
ip = this.sentinelClient.Host;
}
fetchIP = false;
}
else if (fetchPort)
{
port = value;
fetchPort = false;
}
else if (fetchFlags)
{
flags = value;
fetchFlags = false;

if (ip != null && port != null && !flags.Contains("s_down"))
{
servers.Add("{0}:{1}".Fmt(ip, port));
}
}
data.TryGetValue("flags", out flags);
data.TryGetValue("ip", out ip);
data.TryGetValue("port", out port);

if (ip == "127.0.0.1")
{
ip = this.sentinelClient.Host;
}

}
if (ip != null && port != null && !flags.Contains("s_down") && !flags.Contains("o_down"))
{
servers.Add("{0}:{1}".Fmt(ip, port));
}
}

Expand All @@ -190,48 +179,17 @@ private IEnumerable<string> ConvertSlaveArrayToList(object[] slaves)
private IEnumerable<string> ConvertMasterArrayToList(object[] items)
{
var servers = new List<string>();
bool fetchIP = false;
bool fetchPort = false;
string ip = null;
string port = null;
string value = null;

foreach (var item in items)
{
if (item is byte[])
{
value = Encoding.UTF8.GetString((byte[])item);
if (value == "ip")
{
fetchIP = true;
continue;
}
else if (value == "port")
{
fetchPort = true;
continue;
}
else if (fetchIP)
{
ip = value;
if (ip == "127.0.0.1")
{
ip = this.sentinelClient.Host;
}
fetchIP = false;
}
else if (fetchPort)
{
port = value;
fetchPort = false;
}
var data = ParseDataArray(items);

if (ip != null && port != null)
{
servers.Add("{0}:{1}".Fmt(ip, port));
break;
}
}
data.TryGetValue("ip", out ip);
data.TryGetValue("port", out port);

if (ip != null && port != null)
{
servers.Add("{0}:{1}".Fmt(ip, port));
}

return servers;
Expand Down
15 changes: 15 additions & 0 deletions tests/ServiceStack.Redis.Tests/RedisSentinelTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -62,5 +62,20 @@ public void Can_Get_Master_Addr()
Assert.AreEqual(host, "127.0.0.1"); // IP of localhost
Assert.AreEqual(port, TestConfig.RedisPort.ToString());
}

[Test]
public void Can_Get_Redis_ClientsManager()
{
var sentinel = new RedisSentinel(new[] { "{0}:{1}".Fmt(TestConfig.SingleHost, TestConfig.RedisSentinelPort) }, TestConfig.MasterName);

var clientsManager = sentinel.Setup();
var client = clientsManager.GetClient();

Assert.AreEqual(client.Host, "127.0.0.1");
Assert.AreEqual(client.Port, TestConfig.RedisPort);

client.Dispose();
sentinel.Dispose();
}
}
}