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
11 changes: 5 additions & 6 deletions src/Cafe.Barista/Modules/BusModule.cs
Original file line number Diff line number Diff line change
Expand Up @@ -34,8 +34,7 @@ protected override void Load(ContainerBuilder builder)
// )

// Redis Transport
//.WithTransport(new RedisTransportConfiguration().WithConnectionString("bus.iymtwr.0001.apse2.cache.amazonaws.com"))
//.WithTransport(new RedisTransportConfiguration().WithConnectionString("localhost"))
.WithTransport(new RedisTransportConfiguration().WithConnectionString("localhost")).WithAutoDeleteOnIdle(TimeSpan.FromMinutes(10))

// ActiveMQ Transport
//.WithTransport(new AMQPTransportConfiguration()
Expand All @@ -53,10 +52,10 @@ protected override void Load(ContainerBuilder builder)
// .WithAutoCreateSchema())

// NATS Transport
.WithTransport(new NatsTransportConfiguration()
.WithUrl("nats://localhost:4222")
.WithCredentials("admin", "password")
.WithJetStream())
// .WithTransport(new NatsTransportConfiguration()
// .WithUrl("nats://localhost:4222")
// .WithCredentials("admin", "password")
// .WithJetStream())

.WithNames("Barista", Environment.MachineName)
.WithTypesFrom(handlerTypesProvider)
Expand Down
10 changes: 5 additions & 5 deletions src/Cafe.Cashier/Modules/BusModule.cs
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ protected override void Load(ContainerBuilder builder)
// )

// Redis Transport
//.WithTransport(new RedisTransportConfiguration().WithConnectionString("localhost"))
.WithTransport(new RedisTransportConfiguration().WithConnectionString("localhost")).WithAutoDeleteOnIdle(TimeSpan.FromMinutes(10))

// ActiveMQ Transport
//.WithTransport(new AMQPTransportConfiguration()
Expand All @@ -57,10 +57,10 @@ protected override void Load(ContainerBuilder builder)
// .WithAutoCreateSchema())

// NATS Transport
.WithTransport(new NatsTransportConfiguration()
.WithUrl("nats://localhost:4222")
.WithCredentials("admin", "password")
.WithJetStream())
// .WithTransport(new NatsTransportConfiguration()
// .WithUrl("nats://localhost:4222")
// .WithCredentials("admin", "password")
// .WithJetStream())

.WithNames("Cashier", Environment.MachineName)
.WithTypesFrom(handlerTypesProvider)
Expand Down
10 changes: 5 additions & 5 deletions src/Cafe.Waiter/Modules/BusModule.cs
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ protected override void Load(ContainerBuilder builder)
// )

// Redis Transport
//.WithTransport(new RedisTransportConfiguration().WithConnectionString("localhost"))
.WithTransport(new RedisTransportConfiguration().WithConnectionString("localhost")).WithAutoDeleteOnIdle(TimeSpan.FromMinutes(10))

// ActiveMQ Transport
//.WithTransport(new AMQPTransportConfiguration()
Expand All @@ -61,10 +61,10 @@ protected override void Load(ContainerBuilder builder)
// .WithAutoCreateSchema())

// NATS Transport
.WithTransport(new NatsTransportConfiguration()
.WithUrl("nats://localhost:4222")
.WithCredentials("admin", "password")
.WithJetStream())
// .WithTransport(new NatsTransportConfiguration()
// .WithUrl("nats://localhost:4222")
// .WithCredentials("admin", "password")
// .WithJetStream())

.WithNames("Waiter", Environment.MachineName)
.WithTypesFrom(handlerTypesProvider)
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,145 @@
using System;
using System.Threading.Tasks;
using Nimbus.Configuration.Settings;
using Nimbus.Infrastructure.Logging;
using Nimbus.Infrastructure.MessageSendersAndReceivers;
using Nimbus.Infrastructure.Retries;
using Nimbus.InfrastructureContracts;
using Nimbus.Serializers.Json;
using Nimbus.Tests.Integration.Configuration;
using Nimbus.Tests.Integration.TestUtilities;
using Nimbus.Transports.Redis.MessageSendersAndReceivers;
using Nimbus.Transports.Redis.QueueManagement;
using NUnit.Framework;
using Shouldly;
using StackExchange.Redis;

namespace Nimbus.Tests.Integration.Tests.RedisTransportTests
{
[TestFixture]
[RequiresRedis]
public class WhenAnIdleSubscriptionExpires
{
private ConnectionMultiplexer _multiplexer;
private IDatabase _database;
private Subscription _subscription;
private AutoDeleteOnIdleSetting _autoDeleteOnIdle;
private NullLogger _logger;
private IRetry _retry;
private JsonSerializer _serializer;

[SetUp]
public void SetUp()
{
var connectionString = AppSettingsLoader.Settings.Transports.Redis.ConnectionString;
_multiplexer = ConnectionMultiplexer.Connect(connectionString);
_database = _multiplexer.GetDatabase();

var topicPath = $"t.IdleSubscriptionTests.{Guid.NewGuid():N}";
var subscriptionName = $"sub.{Guid.NewGuid():N}";
_subscription = new Subscription(topicPath, subscriptionName);

_autoDeleteOnIdle = new AutoDeleteOnIdleSetting {Value = TimeSpan.FromSeconds(5)};
_logger = new NullLogger();
_retry = new DefaultRetry(_logger);
_serializer = new JsonSerializer();
}

[TearDown]
public void TearDown()
{
_database.KeyDelete(_subscription.TopicSubscribersRedisKey);
_database.KeyDelete(_subscription.SubscriptionMessagesRedisKey);
_database.KeyDelete(_subscription.SubscriberAliveRedisKey);
_multiplexer.Dispose();
}

private RedisSubscriptionReceiver CreateReceiver()
{
return new RedisSubscriptionReceiver(
_subscription,
() => _multiplexer,
() => _database,
_serializer,
new ConcurrentHandlerLimitSetting(),
new GlobalHandlerThrottle(new GlobalConcurrentHandlerLimitSetting()),
_logger,
_retry,
_autoDeleteOnIdle);
}

[Test]
public async Task WarmUpRegistersTheSubscriberAndSetsALivenessKeyWithTtl()
{
var receiver = CreateReceiver();

await receiver.Start(msg => Task.CompletedTask);
try
{
_database.SetContains(_subscription.TopicSubscribersRedisKey, _subscription.SubscriptionMessagesRedisKey).ShouldBeTrue();

var ttl = _database.KeyTimeToLive(_subscription.SubscriberAliveRedisKey);
ttl.HasValue.ShouldBeTrue();
ttl.Value.ShouldBeLessThanOrEqualTo(_autoDeleteOnIdle.Value);
}
finally
{
await receiver.Stop();
receiver.Dispose();
}
}

[Test]
public async Task SendPrunesADeadSubscriberInsteadOfPushingToIt()
{
// No liveness key is set up - simulates a subscriber whose TTL has already expired.
_database.SetAdd(_subscription.TopicSubscribersRedisKey, _subscription.SubscriptionMessagesRedisKey);

var sender = new RedisTopicSender(_subscription.TopicPath, _serializer, () => _database);
await sender.Send(new NimbusMessage(_subscription.TopicPath));

_database.SetContains(_subscription.TopicSubscribersRedisKey, _subscription.SubscriptionMessagesRedisKey).ShouldBeFalse();
_database.ListLength(_subscription.SubscriptionMessagesRedisKey).ShouldBe(0);
}

[Test]
public async Task SendStillDeliversToALiveSubscriber()
{
_database.SetAdd(_subscription.TopicSubscribersRedisKey, _subscription.SubscriptionMessagesRedisKey);
_database.StringSet(_subscription.SubscriberAliveRedisKey, true, _autoDeleteOnIdle.Value);

var sender = new RedisTopicSender(_subscription.TopicPath, _serializer, () => _database);
await sender.Send(new NimbusMessage(_subscription.TopicPath));

_database.SetContains(_subscription.TopicSubscribersRedisKey, _subscription.SubscriptionMessagesRedisKey).ShouldBeTrue();
_database.ListLength(_subscription.SubscriptionMessagesRedisKey).ShouldBe(1);
}

[Test]
public async Task ReaperRemovesAColdSubscriptionThatNobodyIsPublishingTo()
{
// Nothing publishes to this topic, so RedisTopicSender never gets a chance to prune it -
// only the reaper's own sweep can clean it up.
_database.SetAdd(_subscription.TopicSubscribersRedisKey, _subscription.SubscriptionMessagesRedisKey);
_database.ListRightPush(_subscription.SubscriptionMessagesRedisKey, "orphaned-message");

var reaper = new RedisIdleSubscriptionReaper(() => _multiplexer, _autoDeleteOnIdle, _logger);
await reaper.ReapOnce();

_database.SetContains(_subscription.TopicSubscribersRedisKey, _subscription.SubscriptionMessagesRedisKey).ShouldBeFalse();
_database.KeyExists(_subscription.SubscriptionMessagesRedisKey).ShouldBeFalse();
}

[Test]
public async Task ReaperLeavesALiveSubscriptionAlone()
{
_database.SetAdd(_subscription.TopicSubscribersRedisKey, _subscription.SubscriptionMessagesRedisKey);
_database.StringSet(_subscription.SubscriberAliveRedisKey, true, _autoDeleteOnIdle.Value);

var reaper = new RedisIdleSubscriptionReaper(() => _multiplexer, _autoDeleteOnIdle, _logger);
await reaper.ReapOnce();

_database.SetContains(_subscription.TopicSubscribersRedisKey, _subscription.SubscriptionMessagesRedisKey).ShouldBeTrue();
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -55,9 +55,15 @@ private void OnNotificationReceived(RedisChannel redisChannel, RedisValue redisV
_logger.Debug("Redis notification received in receiver for {RedisKey}", _redisKey);
_receiveSemaphore.Release();
}

protected virtual void OnPoll()
{
}

protected override async Task<NimbusMessage> Fetch(CancellationToken cancellationToken)
{
OnPoll();

if (_haveFetchedAllPreExistingMessages) await _receiveSemaphore.WaitAsync(_redisPollInterval, cancellationToken).ConfigureAwaitFalse();

var database = _databaseFunc();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,15 +14,22 @@ internal class RedisSubscriptionReceiver : RedisMessageReceiver
private readonly Subscription _subscription;
private readonly Func<IDatabase> _databaseFunc;
private readonly IRetry _retry;
private readonly AutoDeleteOnIdleSetting _autoDeleteOnIdle;

// Refreshed well inside the TTL so a receiver never sits an entire heartbeat interval away
// from expiry — a missed poll or two shouldn't be enough to make a live subscriber look dead.
private readonly TimeSpan _heartbeatRefreshInterval;
private DateTime _nextHeartbeatRefreshUtc = DateTime.MinValue;

public RedisSubscriptionReceiver(Subscription subscription,
Func<ConnectionMultiplexer> connectionMultiplexerFunc,
Func<IDatabase> databaseFunc,
ISerializer serializer,
ConcurrentHandlerLimitSetting concurrentHandlerLimit,
IGlobalHandlerThrottle globalHandlerThrottle,
ILogger logger,
IRetry retry)
Func<ConnectionMultiplexer> connectionMultiplexerFunc,
Func<IDatabase> databaseFunc,
ISerializer serializer,
ConcurrentHandlerLimitSetting concurrentHandlerLimit,
IGlobalHandlerThrottle globalHandlerThrottle,
ILogger logger,
IRetry retry,
AutoDeleteOnIdleSetting autoDeleteOnIdle)
: base(
subscription.SubscriptionMessagesRedisKey,
connectionMultiplexerFunc,
Expand All @@ -35,13 +42,32 @@ public RedisSubscriptionReceiver(Subscription subscription,
_subscription = subscription;
_databaseFunc = databaseFunc;
_retry = retry;
_autoDeleteOnIdle = autoDeleteOnIdle;
_heartbeatRefreshInterval = TimeSpan.FromTicks(autoDeleteOnIdle.Value.Ticks / 4);
}

protected override async Task WarmUp()
{
var database = _databaseFunc();
await _retry.DoAsync(() => database.SetAddAsync(_subscription.TopicSubscribersRedisKey, _subscription.SubscriptionMessagesRedisKey)).ConfigureAwaitFalse();
await _retry
.DoAsync(() => database.SetAddAsync(_subscription.TopicSubscribersRedisKey,
_subscription.SubscriptionMessagesRedisKey)).ConfigureAwaitFalse();
RefreshHeartbeat(database);
await base.WarmUp().ConfigureAwaitFalse();
}

protected override void OnPoll()
{
var now = DateTime.UtcNow;
if (now < _nextHeartbeatRefreshUtc) return;

RefreshHeartbeat(_databaseFunc());
_nextHeartbeatRefreshUtc = now + _heartbeatRefreshInterval;
}

private void RefreshHeartbeat(IDatabase database)
{
_retry.Do(() => database.StringSet(_subscription.SubscriberAliveRedisKey, true, _autoDeleteOnIdle.Value));
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -27,18 +27,28 @@ public async Task Send(NimbusMessage message)

var subscribersRedisKey = Subscription.TopicSubscribersRedisKeyFor(_topicPath);
var subscribers = database.SetMembers(subscribersRedisKey)
.Select(s => s.ToString())
.ToArray();
.Select(s => s.ToString())
.ToArray();

await subscribers
.Select(subscriberPath => Task.Run(() =>
{
var clone = (NimbusMessage) _serializer.Deserialize(_serializer.Serialize(message), typeof (NimbusMessage));
clone.DeliverTo = subscriberPath;
var serialized = _serializer.Serialize(clone);
database.ListRightPush(subscriberPath, serialized);
database.Publish(subscriberPath, string.Empty);
}).ConfigureAwaitFalse())
{

var aliveKey = Subscription.SubscriberAliveRedisKeyFor(subscriberPath);
if (!database.KeyExists(aliveKey))
{
database.SetRemove(subscribersRedisKey, subscriberPath);
database.KeyDelete(subscriberPath);
return;
}

var clone = (NimbusMessage)_serializer.Deserialize(_serializer.Serialize(message),
typeof(NimbusMessage));
clone.DeliverTo = subscriberPath;
var serialized = _serializer.Serialize(clone);
database.ListRightPush(subscriberPath, serialized);
database.Publish(subscriberPath, string.Empty);
}).ConfigureAwaitFalse())
.WhenAll();
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,18 +8,25 @@ internal class Subscription
public string SubscriptionName { get; }
public string TopicSubscribersRedisKey { get; }
public string SubscriptionMessagesRedisKey { get; }
public string SubscriberAliveRedisKey { get; }

public Subscription(string topicPath, string subscriptionName)
{
TopicPath = topicPath;
SubscriptionName = subscriptionName;
TopicSubscribersRedisKey = TopicSubscribersRedisKeyFor(topicPath);
SubscriptionMessagesRedisKey = $"{topicPath}.{subscriptionName}";
SubscriberAliveRedisKey = SubscriberAliveRedisKeyFor(SubscriptionMessagesRedisKey);
}

public static string TopicSubscribersRedisKeyFor(string topicPath)
{
return $"{SubscriptionsPrefix}.{topicPath}";
}

public static string SubscriberAliveRedisKeyFor(string subscriptionMessagesRedisKey)
{
return $"{subscriptionMessagesRedisKey}.alive";
}
}
}
2 changes: 1 addition & 1 deletion src/Nimbus.Transports.Redis/Nimbus.Transports.Redis.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
<PackageReference Include="StackExchange.Redis" Version="2.8.0" />
</ItemGroup>
<PropertyGroup>
<Version>4.3.3</Version>
<Version>4.4.0</Version>
<GeneratePackageOnBuild>true</GeneratePackageOnBuild>
<Description>Redis transport for the Nimbus messaging framework using StackExchange.Redis.</Description>
<PackageTags>nimbus;messaging;transport;redis</PackageTags>
Expand Down
3 changes: 3 additions & 0 deletions src/Nimbus.Transports.Redis/Properties/AssemblyVisibility.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
using System.Runtime.CompilerServices;

[assembly: InternalsVisibleTo("Nimbus.Tests.Integration", AllInternalsVisible = true)]
Loading
Loading