diff --git a/src/Dafda.Tests/Builders/ConsumerBuilder.cs b/src/Dafda.Tests/Builders/ConsumerBuilder.cs index e72959e1..d0861fd6 100644 --- a/src/Dafda.Tests/Builders/ConsumerBuilder.cs +++ b/src/Dafda.Tests/Builders/ConsumerBuilder.cs @@ -1,78 +1,84 @@ -using Dafda.Consuming; +namespace Dafda.Tests.Builders; + +using Dafda.Consuming; using Dafda.Consuming.MessageFilters; -using Dafda.Tests.TestDoubles; +using TestDoubles; -namespace Dafda.Tests.Builders +internal class ConsumerBuilder { - internal class ConsumerBuilder + private IHandlerUnitOfWorkFactory _unitOfWorkFactory = new HandlerUnitOfWorkFactoryStub(null); + private IConsumerScopeFactory _consumerScopeFactory = new ConsumerScopeFactoryStub(new ConsumerScopeStub(new MessageResultBuilder().Build())); + private MessageHandlerRegistry _registry = new(); + private IUnconfiguredMessageHandlingStrategy _unconfiguredMessageStrategy = new RequireExplicitHandlers(); + private readonly IMessageHandlerExecutionStrategy _messageHandlerExecutionStrategy = new DirectMessageHandlerExecutionStrategy(); + + private bool _enableAutoCommit; + private MessageFilter _messageFilter = MessageFilter.Default; + private IDeadLetterQueue _deadLetterQueue = NullDeadLetterQueue.Instance; + private int _maxRetries; + + public ConsumerBuilder WithUnitOfWork(IHandlerUnitOfWork unitOfWork) { - private IHandlerUnitOfWorkFactory _unitOfWorkFactory; - private IConsumerScopeFactory _consumerScopeFactory; - private MessageHandlerRegistry _registry; - private IUnconfiguredMessageHandlingStrategy _unconfiguredMessageStrategy; - private IMessageHandlerExecutionStrategy _messageHandlerExecutionStrategy; - - private bool _enableAutoCommit; - private MessageFilter _messageFilter = MessageFilter.Default; - - public ConsumerBuilder() - { - _unitOfWorkFactory = new HandlerUnitOfWorkFactoryStub(null); - _consumerScopeFactory = new ConsumerScopeFactoryStub(new ConsumerScopeStub(new MessageResultBuilder().Build())); - _registry = new MessageHandlerRegistry(); - _unconfiguredMessageStrategy = new RequireExplicitHandlers(); - _messageHandlerExecutionStrategy = new DirectMessageHandlerExecutionStrategy(); - } + return WithUnitOfWorkFactory(new HandlerUnitOfWorkFactoryStub(unitOfWork)); + } - public ConsumerBuilder WithUnitOfWork(IHandlerUnitOfWork unitOfWork) - { - return WithUnitOfWorkFactory(new HandlerUnitOfWorkFactoryStub(unitOfWork)); - } + public ConsumerBuilder WithUnitOfWorkFactory(IHandlerUnitOfWorkFactory unitofWorkFactory) + { + _unitOfWorkFactory = unitofWorkFactory; + return this; + } - public ConsumerBuilder WithUnitOfWorkFactory(IHandlerUnitOfWorkFactory unitofWorkFactory) - { - _unitOfWorkFactory = unitofWorkFactory; - return this; - } + public ConsumerBuilder WithConsumerScopeFactory(IConsumerScopeFactory consumerScopeFactory) + { + _consumerScopeFactory = consumerScopeFactory; + return this; + } - public ConsumerBuilder WithConsumerScopeFactory(IConsumerScopeFactory consumerScopeFactory) - { - _consumerScopeFactory = consumerScopeFactory; - return this; - } + public ConsumerBuilder WithMessageHandlerRegistry(MessageHandlerRegistry registry) + { + _registry = registry; + return this; + } - public ConsumerBuilder WithMessageHandlerRegistry(MessageHandlerRegistry registry) - { - _registry = registry; - return this; - } + public ConsumerBuilder WithEnableAutoCommit(bool enableAutoCommit) + { + _enableAutoCommit = enableAutoCommit; + return this; + } - public ConsumerBuilder WithEnableAutoCommit(bool enableAutoCommit) - { - _enableAutoCommit = enableAutoCommit; - return this; - } + public void WithMessageFilter(MessageFilter messageFilter) + { + _messageFilter = messageFilter; + } - public void WithMessageFilter(MessageFilter messageFilter) - { - _messageFilter = messageFilter; - } + public ConsumerBuilder WithUnconfiguredMessageStrategy( + IUnconfiguredMessageHandlingStrategy strategy) + { + _unconfiguredMessageStrategy = strategy; + return this; + } - public ConsumerBuilder WithUnconfiguredMessageStrategy( - IUnconfiguredMessageHandlingStrategy strategy) - { - _unconfiguredMessageStrategy = strategy; - return this; - } + public ConsumerBuilder WithDeadLetterQueue(IDeadLetterQueue deadLetterQueue) + { + _deadLetterQueue = deadLetterQueue; + return this; + } - public Consumer Build() => - new Consumer( - _registry, - _unitOfWorkFactory, - _consumerScopeFactory, - _unconfiguredMessageStrategy, - _messageFilter, - _messageHandlerExecutionStrategy, - _enableAutoCommit); + public ConsumerBuilder WithMaxRetries(int maxRetries) + { + _maxRetries = maxRetries; + return this; } -} + + public Consumer Build() => + new Consumer( + _registry, + _unitOfWorkFactory, + _consumerScopeFactory, + _unconfiguredMessageStrategy, + _messageFilter, + _messageHandlerExecutionStrategy, + _enableAutoCommit, + _deadLetterQueue, + _maxRetries); +} \ No newline at end of file diff --git a/src/Dafda.Tests/Consuming/TestConsumer.cs b/src/Dafda.Tests/Consuming/TestConsumer.cs index 63b2fa82..ad9df7c5 100644 --- a/src/Dafda.Tests/Consuming/TestConsumer.cs +++ b/src/Dafda.Tests/Consuming/TestConsumer.cs @@ -1,194 +1,279 @@ -using System; +namespace Dafda.Tests.Consuming; + +using System; using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; +using Builders; using Dafda.Consuming; -using Dafda.Tests.Builders; -using Dafda.Tests.TestDoubles; using Microsoft.Extensions.Logging; using Moq; +using TestDoubles; using Xunit; -namespace Dafda.Tests.Consuming +public class TestConsumer { - public class TestConsumer + [Fact] + public async Task invokes_expected_handler_when_consuming() { - [Fact] - public async Task invokes_expected_handler_when_consuming() - { - var handlerMock = new Mock>(); - var handlerStub = handlerMock.Object; - - var messageRegistrationStub = new MessageRegistrationBuilder() - .WithHandlerInstanceType(handlerStub.GetType()) - .WithMessageInstanceType(typeof(FooMessage)) - .WithMessageType("foo") - .WithTopic("") - .Build(); + var handlerMock = new Mock>(); + var handlerStub = handlerMock.Object; - var registry = new MessageHandlerRegistry(); - registry.Register(messageRegistrationStub); + var messageRegistrationStub = new MessageRegistrationBuilder() + .WithHandlerInstanceType(handlerStub.GetType()) + .WithMessageInstanceType(typeof(FooMessage)) + .WithMessageType("foo") + .WithTopic("") + .Build(); - var sut = new ConsumerBuilder() - .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) - .WithMessageHandlerRegistry(registry) + var registry = new MessageHandlerRegistry(); + registry.Register(messageRegistrationStub); + + var sut = new ConsumerBuilder() + .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) + .WithMessageHandlerRegistry(registry) + .Build(); + + await sut.ConsumeSingle(CancellationToken.None); + + handlerMock.Verify(x => x.Handle(It.IsAny(), It.IsAny(), It.IsAny()), Times.Once); + } + + [Fact] + public async Task throws_when_consuming_an_unknown_message_when_explicit_handlers_are_required() + { + var sut = new ConsumerBuilder().Build(); + + await Assert.ThrowsAsync( + () => sut.ConsumeSingle(CancellationToken.None)); + } + + [Fact] + public async Task does_not_throw_when_consuming_an_unknown_message_with_no_op_strategy() + { + var sut = + new ConsumerBuilder() + .WithUnitOfWork( + new UnitOfWorkStub( + new NoOpHandler(new Mock>().Object))) + .WithUnconfiguredMessageStrategy(new UseNoOpHandler()) .Build(); - await sut.ConsumeSingle(CancellationToken.None); + await sut.ConsumeSingle(CancellationToken.None); + } - handlerMock.Verify(x => x.Handle(It.IsAny(), It.IsAny(), It.IsAny()), Times.Once); - } + [Fact] + public async Task expected_order_of_handler_invocation_in_unit_of_work() + { + var orderOfInvocation = new LinkedList(); + + var dummyMessageResult = new MessageResultBuilder() + .WithTransportLevelMessage(new TransportLevelMessageBuilder().WithType("foo").Build()) + .WithTopic("topic") + .Build(); + var dummyMessageRegistration = new MessageRegistrationBuilder() + .WithMessageType("foo") + .WithTopic("topic") + .Build(); + + var registry = new MessageHandlerRegistry(); + registry.Register(dummyMessageRegistration); + + var sut = new ConsumerBuilder() + .WithUnitOfWork(new UnitOfWorkSpy( + handlerInstance: new MessageHandlerSpy(() => orderOfInvocation.AddLast("during")), + pre: () => orderOfInvocation.AddLast("before"), + post: () => orderOfInvocation.AddLast("after") + )) + .WithConsumerScopeFactory(new ConsumerScopeFactoryStub(new ConsumerScopeStub(dummyMessageResult))) + .WithMessageHandlerRegistry(registry) + .Build(); + + await sut.ConsumeSingle(CancellationToken.None); + + Assert.Equal(new[] { "before", "during", "after" }, orderOfInvocation); + } - [Fact] - public async Task throws_when_consuming_an_unknown_message_when_explicit_handlers_are_required() - { - var sut = new ConsumerBuilder().Build(); + [Fact] + public async Task will_not_call_commit_when_auto_commit_is_enabled() + { + var handlerStub = Dummy.Of>(); - await Assert.ThrowsAsync( - () => sut.ConsumeSingle(CancellationToken.None)); - } + var messageRegistrationStub = new MessageRegistrationBuilder() + .WithHandlerInstanceType(handlerStub.GetType()) + .WithMessageInstanceType(typeof(FooMessage)) + .WithMessageType("foo") + .WithTopic("topic") + .Build(); - [Fact] - public async Task does_not_throw_when_consuming_an_unknown_message_with_no_op_strategy() - { - var sut = - new ConsumerBuilder() - .WithUnitOfWork( - new UnitOfWorkStub( - new NoOpHandler(new Mock>().Object))) - .WithUnconfiguredMessageStrategy(new UseNoOpHandler()) - .Build(); - - await sut.ConsumeSingle(CancellationToken.None); - } + var wasCalled = false; - [Fact] - public async Task expected_order_of_handler_invocation_in_unit_of_work() - { - var orderOfInvocation = new LinkedList(); + var resultSpy = new MessageResultBuilder() + .WithOnCommit((_) => + { + wasCalled = true; + return Task.CompletedTask; + }) + .WithTopic("topic") + .Build(); + + var consumerScopeFactoryStub = new ConsumerScopeFactoryStub(new ConsumerScopeStub(resultSpy)); + var registry = new MessageHandlerRegistry(); + registry.Register(messageRegistrationStub); + + var consumer = new ConsumerBuilder() + .WithConsumerScopeFactory(consumerScopeFactoryStub) + .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) + .WithMessageHandlerRegistry(registry) + .WithEnableAutoCommit(true) + .Build(); + + await consumer.ConsumeSingle(CancellationToken.None); + + Assert.False(wasCalled); + } - var dummyMessageResult = new MessageResultBuilder() - .WithTransportLevelMessage(new TransportLevelMessageBuilder().WithType("foo").Build()) - .WithTopic("topic") - .Build(); - var dummyMessageRegistration = new MessageRegistrationBuilder() - .WithMessageType("foo") - .WithTopic("topic") - .Build(); + [Fact] + public async Task will_call_commit_when_auto_commit_is_disabled() + { + var handlerStub = Dummy.Of>(); - var registry = new MessageHandlerRegistry(); - registry.Register(dummyMessageRegistration); - - var sut = new ConsumerBuilder() - .WithUnitOfWork(new UnitOfWorkSpy( - handlerInstance: new MessageHandlerSpy(() => orderOfInvocation.AddLast("during")), - pre: () => orderOfInvocation.AddLast("before"), - post: () => orderOfInvocation.AddLast("after") - )) - .WithConsumerScopeFactory(new ConsumerScopeFactoryStub(new ConsumerScopeStub(dummyMessageResult))) - .WithMessageHandlerRegistry(registry) - .Build(); + var messageRegistrationStub = new MessageRegistrationBuilder() + .WithHandlerInstanceType(handlerStub.GetType()) + .WithMessageInstanceType(typeof(FooMessage)) + .WithMessageType("foo") + .WithTopic("topic") + .Build(); - await sut.ConsumeSingle(CancellationToken.None); + var wasCalled = false; - Assert.Equal(new[] { "before", "during", "after" }, orderOfInvocation); - } + var resultSpy = new MessageResultBuilder() + .WithTopic("topic") + .WithOnCommit((_) => + { + wasCalled = true; + return Task.CompletedTask; + }) + .Build(); - [Fact] - public async Task will_not_call_commit_when_auto_commit_is_enabled() - { - var handlerStub = Dummy.Of>(); + var consumerScopeFactoryStub = new ConsumerScopeFactoryStub(new ConsumerScopeStub(resultSpy)); + var registry = new MessageHandlerRegistry(); + registry.Register(messageRegistrationStub); - var messageRegistrationStub = new MessageRegistrationBuilder() - .WithHandlerInstanceType(handlerStub.GetType()) - .WithMessageInstanceType(typeof(FooMessage)) - .WithMessageType("foo") - .WithTopic("topic") - .Build(); + var consumer = new ConsumerBuilder() + .WithConsumerScopeFactory(consumerScopeFactoryStub) + .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) + .WithMessageHandlerRegistry(registry) + .WithEnableAutoCommit(false) + .Build(); - var wasCalled = false; + await consumer.ConsumeSingle(CancellationToken.None); - var resultSpy = new MessageResultBuilder() - .WithOnCommit((_) => - { - wasCalled = true; - return Task.CompletedTask; - }) - .WithTopic("topic") - .Build(); + Assert.True(wasCalled); + } - var consumerScopeFactoryStub = new ConsumerScopeFactoryStub(new ConsumerScopeStub(resultSpy)); - var registry = new MessageHandlerRegistry(); - registry.Register(messageRegistrationStub); + [Fact] + public async Task creates_consumer_scope_when_consuming_single_message() + { + var messageResultStub = new MessageResultBuilder() + .WithTransportLevelMessage(new TransportLevelMessageBuilder().WithType("foo").Build()) + .WithTopic("topic") + .Build(); + var handlerStub = Dummy.Of>(); - var consumer = new ConsumerBuilder() - .WithConsumerScopeFactory(consumerScopeFactoryStub) - .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) - .WithMessageHandlerRegistry(registry) - .WithEnableAutoCommit(true) - .Build(); + var messageRegistrationStub = new MessageRegistrationBuilder() + .WithHandlerInstanceType(handlerStub.GetType()) + .WithMessageInstanceType(typeof(FooMessage)) + .WithMessageType("foo") + .WithTopic("topic") + .Build(); - await consumer.ConsumeSingle(CancellationToken.None); + var spy = new ConsumerScopeFactorySpy(new ConsumerScopeStub(messageResultStub)); - Assert.False(wasCalled); - } + var registry = new MessageHandlerRegistry(); + registry.Register(messageRegistrationStub); - [Fact] - public async Task will_call_commit_when_auto_commit_is_disabled() - { - var handlerStub = Dummy.Of>(); + var consumer = new ConsumerBuilder() + .WithConsumerScopeFactory(spy) + .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) + .WithMessageHandlerRegistry(registry) + .Build(); - var messageRegistrationStub = new MessageRegistrationBuilder() - .WithHandlerInstanceType(handlerStub.GetType()) - .WithMessageInstanceType(typeof(FooMessage)) - .WithMessageType("foo") - .WithTopic("topic") - .Build(); + await consumer.ConsumeSingle(CancellationToken.None); - var wasCalled = false; - var resultSpy = new MessageResultBuilder() - .WithTopic("topic") - .WithOnCommit((_) => - { - wasCalled = true; - return Task.CompletedTask; - }) - .Build(); + Assert.Equal(1, spy.CreateConsumerScopeCalled); + } - var consumerScopeFactoryStub = new ConsumerScopeFactoryStub(new ConsumerScopeStub(resultSpy)); - var registry = new MessageHandlerRegistry(); - registry.Register(messageRegistrationStub); + [Fact] + public async Task disposes_consumer_scope_when_consuming_single_message() + { + var messageResultStub = new MessageResultBuilder() + .WithTransportLevelMessage(new TransportLevelMessageBuilder().WithType("foo").Build()) + .WithTopic("topic") + .Build(); + var handlerStub = Dummy.Of>(); - var consumer = new ConsumerBuilder() - .WithConsumerScopeFactory(consumerScopeFactoryStub) - .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) - .WithMessageHandlerRegistry(registry) - .WithEnableAutoCommit(false) - .Build(); + var messageRegistrationStub = new MessageRegistrationBuilder() + .WithHandlerInstanceType(handlerStub.GetType()) + .WithMessageInstanceType(typeof(FooMessage)) + .WithMessageType("foo") + .WithTopic("topic") + .Build(); - await consumer.ConsumeSingle(CancellationToken.None); + var spy = new ConsumerScopeSpy(messageResultStub); - Assert.True(wasCalled); - } + var registry = new MessageHandlerRegistry(); + registry.Register(messageRegistrationStub); - [Fact] - public async Task creates_consumer_scope_when_consuming_single_message() + var consumer = new ConsumerBuilder() + .WithConsumerScopeFactory(new ConsumerScopeFactoryStub(spy)) + .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) + .WithMessageHandlerRegistry(registry) + .Build(); + + await consumer.ConsumeSingle(CancellationToken.None); + + + Assert.Equal(1, spy.Disposed); + } + + [Fact] + public async Task creates_consumer_scope_when_consuming_multiple_messages() + { + var messageResultStub = new MessageResultBuilder() + .WithTransportLevelMessage(new TransportLevelMessageBuilder() + .WithType("foo") + .Build()) + .WithTopic("topic") + .Build(); + var handlerStub = Dummy.Of>(); + + var messageRegistrationStub = new MessageRegistrationBuilder() + .WithHandlerInstanceType(handlerStub.GetType()) + .WithMessageInstanceType(typeof(FooMessage)) + .WithMessageType("foo") + .WithTopic("topic") + .Build(); + + using (var cancellationTokenSource = new CancellationTokenSource(TimeSpan.FromSeconds(5))) { - var messageResultStub = new MessageResultBuilder() - .WithTransportLevelMessage(new TransportLevelMessageBuilder().WithType("foo").Build()) - .WithTopic("topic") - .Build(); - var handlerStub = Dummy.Of>(); + var loops = 0; - var messageRegistrationStub = new MessageRegistrationBuilder() - .WithHandlerInstanceType(handlerStub.GetType()) - .WithMessageInstanceType(typeof(FooMessage)) - .WithMessageType("foo") - .WithTopic("topic") - .Build(); + var subscriberScopeStub = new ConsumerScopeDecoratorWithHooks( + inner: new ConsumerScopeStub(messageResultStub), + postHook: () => + { + loops++; - var spy = new ConsumerScopeFactorySpy(new ConsumerScopeStub(messageResultStub)); + if (loops == 2) + { + cancellationTokenSource.Cancel(); + } + } + ); + + var spy = new ConsumerScopeFactorySpy(subscriberScopeStub); var registry = new MessageHandlerRegistry(); registry.Register(messageRegistrationStub); @@ -199,29 +284,42 @@ public async Task creates_consumer_scope_when_consuming_single_message() .WithMessageHandlerRegistry(registry) .Build(); - await consumer.ConsumeSingle(CancellationToken.None); - + await consumer.ConsumeAll(cancellationTokenSource.Token); + Assert.Equal(2, loops); Assert.Equal(1, spy.CreateConsumerScopeCalled); } + } - [Fact] - public async Task disposes_consumer_scope_when_consuming_single_message() + [Fact] + public async Task disposes_consumer_scope_when_consuming_multiple_messages() + { + var messageResultStub = new MessageResultBuilder() + .WithTransportLevelMessage(new TransportLevelMessageBuilder().WithType("foo").Build()) + .WithTopic("topic") + .Build(); + var handlerStub = Dummy.Of>(); + + var messageRegistrationStub = new MessageRegistrationBuilder() + .WithHandlerInstanceType(handlerStub.GetType()) + .WithMessageInstanceType(typeof(FooMessage)) + .WithMessageType("foo") + .WithTopic("topic") + .Build(); + + using (var cancellationTokenSource = new CancellationTokenSource(TimeSpan.FromSeconds(5))) { - var messageResultStub = new MessageResultBuilder() - .WithTransportLevelMessage(new TransportLevelMessageBuilder().WithType("foo").Build()) - .WithTopic("topic") - .Build(); - var handlerStub = Dummy.Of>(); + var loops = 0; - var messageRegistrationStub = new MessageRegistrationBuilder() - .WithHandlerInstanceType(handlerStub.GetType()) - .WithMessageInstanceType(typeof(FooMessage)) - .WithMessageType("foo") - .WithTopic("topic") - .Build(); + var spy = new ConsumerScopeSpy(messageResultStub, () => + { + loops++; - var spy = new ConsumerScopeSpy(messageResultStub); + if (loops == 2) + { + cancellationTokenSource.Cancel(); + } + }); var registry = new MessageHandlerRegistry(); registry.Register(messageRegistrationStub); @@ -232,199 +330,262 @@ public async Task disposes_consumer_scope_when_consuming_single_message() .WithMessageHandlerRegistry(registry) .Build(); - await consumer.ConsumeSingle(CancellationToken.None); - + await consumer.ConsumeAll(cancellationTokenSource.Token); + Assert.Equal(2, loops); Assert.Equal(1, spy.Disposed); } + } - [Fact] - public async Task creates_consumer_scope_when_consuming_multiple_messages() - { - var messageResultStub = new MessageResultBuilder() - .WithTransportLevelMessage(new TransportLevelMessageBuilder() - .WithType("foo") - .Build()) - .WithTopic("topic") - .Build(); - var handlerStub = Dummy.Of>(); - - var messageRegistrationStub = new MessageRegistrationBuilder() - .WithHandlerInstanceType(handlerStub.GetType()) - .WithMessageInstanceType(typeof(FooMessage)) - .WithMessageType("foo") - .WithTopic("topic") - .Build(); + [Fact] + public async Task throws_when_task_is_canceled() + { + using var cts = new CancellationTokenSource(); + cts.Cancel(); - using (var cancellationTokenSource = new CancellationTokenSource(TimeSpan.FromSeconds(5))) - { - var loops = 0; + var handlerMock = new Mock>(); + var handlerStub = handlerMock.Object; - var subscriberScopeStub = new ConsumerScopeDecoratorWithHooks( - inner: new ConsumerScopeStub(messageResultStub), - postHook: () => - { - loops++; + var messageRegistrationStub = new MessageRegistrationBuilder() + .WithHandlerInstanceType(handlerStub.GetType()) + .WithMessageInstanceType(typeof(FooMessage)) + .WithMessageType("foo") + .WithTopic("") + .Build(); - if (loops == 2) - { - cancellationTokenSource.Cancel(); - } - } - ); + var registry = new MessageHandlerRegistry(); + registry.Register(messageRegistrationStub); - var spy = new ConsumerScopeFactorySpy(subscriberScopeStub); + var sut = new ConsumerBuilder() + .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) + .WithMessageHandlerRegistry(registry) + .Build(); - var registry = new MessageHandlerRegistry(); - registry.Register(messageRegistrationStub); + await Assert.ThrowsAsync(() => sut.ConsumeSingle(cts.Token)); + } - var consumer = new ConsumerBuilder() - .WithConsumerScopeFactory(spy) - .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) - .WithMessageHandlerRegistry(registry) - .Build(); + [Fact] + public async Task dead_letters_and_commits_when_handler_keeps_failing() + { + var handlerInvocations = 0; + var handler = new MessageHandlerSpy(() => + { + handlerInvocations++; + throw new InvalidOperationException("boom"); + }); - await consumer.ConsumeAll(cancellationTokenSource.Token); + var deadLetterQueueSpy = new DeadLetterQueueSpy(); + var committed = false; - Assert.Equal(2, loops); - Assert.Equal(1, spy.CreateConsumerScopeCalled); - } - } + var sut = BuildConsumerWithHandler( + handler, + onCommit: _ => + { + committed = true; + return Task.CompletedTask; + }, + deadLetterQueue: deadLetterQueueSpy, + maxRetries: 2); - [Fact] - public async Task disposes_consumer_scope_when_consuming_multiple_messages() - { - var messageResultStub = new MessageResultBuilder() - .WithTransportLevelMessage(new TransportLevelMessageBuilder().WithType("foo").Build()) - .WithTopic("topic") - .Build(); - var handlerStub = Dummy.Of>(); + await sut.ConsumeSingle(CancellationToken.None); - var messageRegistrationStub = new MessageRegistrationBuilder() - .WithHandlerInstanceType(handlerStub.GetType()) - .WithMessageInstanceType(typeof(FooMessage)) - .WithMessageType("foo") - .WithTopic("topic") - .Build(); + Assert.Equal(3, handlerInvocations); + Assert.Equal(1, deadLetterQueueSpy.SendCount); + Assert.True(committed); + } - using (var cancellationTokenSource = new CancellationTokenSource(TimeSpan.FromSeconds(5))) + [Fact] + public async Task does_not_dead_letter_when_handler_eventually_succeeds() + { + var handlerInvocations = 0; + var handler = new MessageHandlerSpy(() => + { + handlerInvocations++; + if (handlerInvocations < 2) { - var loops = 0; + throw new InvalidOperationException("boom"); + } + }); - var spy = new ConsumerScopeSpy(messageResultStub, () => - { - loops++; + var deadLetterQueueSpy = new DeadLetterQueueSpy(); - if (loops == 2) - { - cancellationTokenSource.Cancel(); - } - }); + var sut = BuildConsumerWithHandler( + handler, + deadLetterQueue: deadLetterQueueSpy, + maxRetries: 2); - var registry = new MessageHandlerRegistry(); - registry.Register(messageRegistrationStub); + await sut.ConsumeSingle(CancellationToken.None); - var consumer = new ConsumerBuilder() - .WithConsumerScopeFactory(new ConsumerScopeFactoryStub(spy)) - .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) - .WithMessageHandlerRegistry(registry) - .Build(); + Assert.Equal(2, handlerInvocations); + Assert.Equal(0, deadLetterQueueSpy.SendCount); + } - await consumer.ConsumeAll(cancellationTokenSource.Token); + [Fact] + public async Task propagates_exception_when_no_dead_letter_queue_is_configured() + { + var handler = new MessageHandlerSpy(() => throw new InvalidOperationException("boom")); - Assert.Equal(2, loops); - Assert.Equal(1, spy.Disposed); - } - } + var sut = BuildConsumerWithHandler(handler); + + await Assert.ThrowsAsync( + () => sut.ConsumeSingle(CancellationToken.None)); + } + + [Fact] + public async Task does_not_dead_letter_when_cancelled_during_dispatch() + { + using var cts = new CancellationTokenSource(); - [Fact] - public async Task throws_when_task_is_canceled() + var handler = new MessageHandlerSpy(() => { - using var cts = new CancellationTokenSource(); cts.Cancel(); + cts.Token.ThrowIfCancellationRequested(); + }); - var handlerMock = new Mock>(); - var handlerStub = handlerMock.Object; + var deadLetterQueueSpy = new DeadLetterQueueSpy(); - var messageRegistrationStub = new MessageRegistrationBuilder() - .WithHandlerInstanceType(handlerStub.GetType()) - .WithMessageInstanceType(typeof(FooMessage)) - .WithMessageType("foo") - .WithTopic("") - .Build(); + var sut = BuildConsumerWithHandler( + handler, + deadLetterQueue: deadLetterQueueSpy, + maxRetries: 3); - var registry = new MessageHandlerRegistry(); - registry.Register(messageRegistrationStub); + await Assert.ThrowsAsync( + () => sut.ConsumeSingle(cts.Token)); - var sut = new ConsumerBuilder() - .WithUnitOfWork(new UnitOfWorkStub(handlerStub)) - .WithMessageHandlerRegistry(registry) - .Build(); - - await Assert.ThrowsAsync(() => sut.ConsumeSingle(cts.Token)); - } - #region helper classes - - private class ConsumerScopeDecoratorWithHooks : ConsumerScope - { - private readonly ConsumerScope _inner; - private readonly Action _preHook; - private readonly Action _postHook; + Assert.Equal(0, deadLetterQueueSpy.SendCount); + } - public ConsumerScopeDecoratorWithHooks(ConsumerScope inner, Action preHook = null, Action postHook = null) - { - _inner = inner; - _preHook = preHook; - _postHook = postHook; - } + [Fact] + public void disposing_consumer_disposes_the_dead_letter_queue() + { + var deadLetterQueueSpy = new DeadLetterQueueSpy(); - public override async Task GetNext(CancellationToken cancellationToken) - { - _preHook?.Invoke(); - var result = await _inner.GetNext(cancellationToken); - _postHook?.Invoke(); + var sut = BuildConsumerWithHandler( + new MessageHandlerSpy(() => { }), + deadLetterQueue: deadLetterQueueSpy); - return result; - } + ((IDisposable)sut).Dispose(); - public override void Dispose() - { - _inner.Dispose(); - } - } + Assert.Equal(1, deadLetterQueueSpy.DisposedCount); + } - public class FooMessage + private static Consumer BuildConsumerWithHandler( + IMessageHandler handler, + Func onCommit = null, + IDeadLetterQueue deadLetterQueue = null, + int maxRetries = 0) + { + var registration = new MessageRegistrationBuilder() + .WithHandlerInstanceType(handler.GetType()) + .WithMessageInstanceType(typeof(FooMessage)) + .WithMessageType("foo") + .WithTopic("topic") + .Build(); + + var registry = new MessageHandlerRegistry(); + registry.Register(registration); + + var messageResult = new MessageResultBuilder() + .WithTransportLevelMessage(new TransportLevelMessageBuilder().WithType("foo").Build()) + .WithTopic("topic") + .WithOnCommit(onCommit ?? (_ => Task.CompletedTask)) + .Build(); + + var builder = new ConsumerBuilder() + .WithConsumerScopeFactory(new ConsumerScopeFactoryStub(new ConsumerScopeStub(messageResult))) + .WithUnitOfWork(new UnitOfWorkStub(handler)) + .WithMessageHandlerRegistry(registry) + .WithMaxRetries(maxRetries); + + if (deadLetterQueue != null) { - public string Value { get; set; } + builder.WithDeadLetterQueue(deadLetterQueue); } - #endregion + return builder.Build(); } + #region helper classes - internal class ConsumerScopeSpy : ConsumerScope + private class ConsumerScopeDecoratorWithHooks : ConsumerScope { - private readonly MessageResult _messageResult; - private readonly Action _onGetNext; + private readonly ConsumerScope _inner; + private readonly Action _preHook; + private readonly Action _postHook; - public ConsumerScopeSpy(MessageResult messageResult, Action onGetNext = null) + public ConsumerScopeDecoratorWithHooks(ConsumerScope inner, Action preHook = null, Action postHook = null) { - _messageResult = messageResult; - _onGetNext = onGetNext; + _inner = inner; + _preHook = preHook; + _postHook = postHook; } - public override Task GetNext(CancellationToken cancellationToken) + public override async Task GetNext(CancellationToken cancellationToken) { - cancellationToken.ThrowIfCancellationRequested(); - _onGetNext?.Invoke(); + _preHook?.Invoke(); + var result = await _inner.GetNext(cancellationToken); + _postHook?.Invoke(); - return Task.FromResult(_messageResult); + return result; } public override void Dispose() { - Disposed++; + _inner.Dispose(); } + } - public int Disposed { get; private set; } + public class FooMessage + { + public string Value { get; set; } } + + private class DeadLetterQueueSpy : IDeadLetterQueue, IDisposable + { + public int SendCount { get; private set; } + public MessageResult LastMessage { get; private set; } + public Exception LastException { get; private set; } + public int DisposedCount { get; private set; } + + public Task Send(MessageResult message, Exception exception, CancellationToken cancellationToken) + { + SendCount++; + LastMessage = message; + LastException = exception; + return Task.CompletedTask; + } + + public void Dispose() + { + DisposedCount++; + } + } + + #endregion +} + +internal class ConsumerScopeSpy : ConsumerScope +{ + private readonly MessageResult _messageResult; + private readonly Action _onGetNext; + + public ConsumerScopeSpy(MessageResult messageResult, Action onGetNext = null) + { + _messageResult = messageResult; + _onGetNext = onGetNext; + } + + public override Task GetNext(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + _onGetNext?.Invoke(); + + return Task.FromResult(_messageResult); + } + + public override void Dispose() + { + Disposed++; + } + + public int Disposed { get; private set; } } \ No newline at end of file diff --git a/src/Dafda.Tests/Consuming/TestKafkaDeadLetterQueue.cs b/src/Dafda.Tests/Consuming/TestKafkaDeadLetterQueue.cs new file mode 100644 index 00000000..bbfa364c --- /dev/null +++ b/src/Dafda.Tests/Consuming/TestKafkaDeadLetterQueue.cs @@ -0,0 +1,31 @@ +namespace Dafda.Tests.Consuming; + +using Dafda.Consuming; +using Xunit; + +public class TestKafkaDeadLetterQueue +{ + [Fact] + public void uses_configured_topic_name_when_provided() + { + var topic = KafkaDeadLetterQueue.ResolveTopicName("orders.dead-letter", "orders", "order-processor"); + + Assert.Equal("orders.dead-letter", topic); + } + + [Fact] + public void derives_topic_name_from_source_topic_and_group_id_when_not_provided() + { + var topic = KafkaDeadLetterQueue.ResolveTopicName(null, "orders", "order-processor"); + + Assert.Equal("orders.order-processor.dead-letter", topic); + } + + [Fact] + public void derives_topic_name_from_source_topic_only_when_group_id_is_missing() + { + var topic = KafkaDeadLetterQueue.ResolveTopicName(null, "orders", null); + + Assert.Equal("orders.dead-letter", topic); + } +} \ No newline at end of file diff --git a/src/Dafda/Configuration/ConsumerConfiguration.cs b/src/Dafda/Configuration/ConsumerConfiguration.cs index 395471c5..0859c410 100644 --- a/src/Dafda/Configuration/ConsumerConfiguration.cs +++ b/src/Dafda/Configuration/ConsumerConfiguration.cs @@ -1,20 +1,24 @@ +namespace Dafda.Configuration; + using System; using System.Collections.Generic; -using Dafda.Consuming; -using Dafda.Consuming.Interfaces; -using Dafda.Consuming.MessageFilters; - -namespace Dafda.Configuration; +using Consuming; +using Consuming.Interfaces; +using Consuming.MessageFilters; internal class ConsumerConfiguration( IDictionary configuration, MessageHandlerRegistry messageHandlerRegistry, ConsumerConfigurationFactories factories, MessageFilter messageFilter, - IConsumerErrorHandler consumerErrorHandler) + IConsumerErrorHandler consumerErrorHandler, + Func deadLetterQueueFactory, + int maxRetries) : ConsumerConfigurationBase(configuration, factories.UnitOfWorkFactory, consumerErrorHandler) { public ConsumerConfigurationFactories Factories { get; } = factories; public MessageHandlerRegistry MessageHandlerRegistry { get; } = messageHandlerRegistry; public MessageFilter MessageFilter { get; } = messageFilter; + public Func DeadLetterQueueFactory { get; } = deadLetterQueueFactory; + public int MaxRetries { get; } = maxRetries; } \ No newline at end of file diff --git a/src/Dafda/Configuration/ConsumerConfigurationBuilder.cs b/src/Dafda/Configuration/ConsumerConfigurationBuilder.cs index 3b31745f..b6d0bc2b 100644 --- a/src/Dafda/Configuration/ConsumerConfigurationBuilder.cs +++ b/src/Dafda/Configuration/ConsumerConfigurationBuilder.cs @@ -1,209 +1,239 @@ +namespace Dafda.Configuration; + using System; using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; -using Dafda.Consuming; -using Dafda.Consuming.MessageFilters; +using Consuming; +using Consuming.MessageFilters; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; -namespace Dafda.Configuration +internal sealed class ConsumerConfigurationBuilder { - internal sealed class ConsumerConfigurationBuilder + private static readonly string[] DefaultConfigurationKeys = { - private static readonly string[] DefaultConfigurationKeys = - { - ConfigurationKey.GroupId, - ConfigurationKey.EnableAutoCommit, - ConfigurationKey.AllowAutoCreateTopics, - ConfigurationKey.BootstrapServers, - ConfigurationKey.BrokerVersionFallback, - ConfigurationKey.ApiVersionFallbackMs, - ConfigurationKey.SslCaLocation, - ConfigurationKey.SaslUsername, - ConfigurationKey.SaslPassword, - ConfigurationKey.SaslMechanisms, - ConfigurationKey.SecurityProtocol, - }; - - private static readonly string[] RequiredConfigurationKeys = - { - ConfigurationKey.GroupId, - ConfigurationKey.BootstrapServers - }; - - private readonly IDictionary _configurations = new Dictionary(); - private readonly IList _namingConventions = new List(); - private readonly MessageHandlerRegistry _messageHandlerRegistry = new MessageHandlerRegistry(); + ConfigurationKey.GroupId, + ConfigurationKey.EnableAutoCommit, + ConfigurationKey.AllowAutoCreateTopics, + ConfigurationKey.BootstrapServers, + ConfigurationKey.BrokerVersionFallback, + ConfigurationKey.ApiVersionFallbackMs, + ConfigurationKey.SslCaLocation, + ConfigurationKey.SaslUsername, + ConfigurationKey.SaslPassword, + ConfigurationKey.SaslMechanisms, + ConfigurationKey.SecurityProtocol, + }; + + private static readonly string[] RequiredConfigurationKeys = + { + ConfigurationKey.GroupId, + ConfigurationKey.BootstrapServers + }; - private ConfigurationSource _configurationSource = ConfigurationSource.Null; + private readonly IDictionary _configurations = new Dictionary(); + private readonly IList _namingConventions = new List(); + private readonly MessageHandlerRegistry _messageHandlerRegistry = new MessageHandlerRegistry(); - private Func _handlerUnitOfWorkFactory; + private ConfigurationSource _configurationSource = ConfigurationSource.Null; - private Func _unconfiguredMessageHandlingStrategy; - private Func _consumerScopeFactory; - private Func _incomingMessageFactory = _ => new JsonIncomingMessageFactory(); + private Func _handlerUnitOfWorkFactory; - private Func _messageHandlerExecutionStrategyFactory; - private bool _readFromBeginning; + private Func _unconfiguredMessageHandlingStrategy; + private Func _consumerScopeFactory; + private Func _incomingMessageFactory = _ => new JsonIncomingMessageFactory(); - private MessageFilter _messageFilter = MessageFilter.Default; - private ConsumerErrorHandler _consumerErrorHandler = ConsumerErrorHandler.Default; + private Func _messageHandlerExecutionStrategyFactory; + private bool _readFromBeginning; - private void ApplyDefaults(IDictionary configurations) - { - _handlerUnitOfWorkFactory ??= sp => ActivatorUtilities.CreateInstance(sp); - _unconfiguredMessageHandlingStrategy ??= sp => ActivatorUtilities.CreateInstance(sp); - _messageHandlerExecutionStrategyFactory ??= sp => ActivatorUtilities.CreateInstance(sp); - _consumerScopeFactory ??= provider => - { - var loggerFactory = provider.GetRequiredService(); - - return new KafkaBasedConsumerScopeFactory( - loggerFactory: loggerFactory, - configuration: configurations, - topics: _messageHandlerRegistry.GetAllSubscribedTopics(), - incomingMessageFactory: _incomingMessageFactory(provider), - readFromBeginning: _readFromBeginning - ); - }; - } + private MessageFilter _messageFilter = MessageFilter.Default; + private ConsumerErrorHandler _consumerErrorHandler = ConsumerErrorHandler.Default; + private DeadLetterQueueOptions _deadLetterQueueOptions; - public ConsumerConfigurationBuilder WithConfigurationSource(ConfigurationSource configurationSource) + private void ApplyDefaults(IDictionary configurations) + { + _handlerUnitOfWorkFactory ??= sp => ActivatorUtilities.CreateInstance(sp); + _unconfiguredMessageHandlingStrategy ??= sp => ActivatorUtilities.CreateInstance(sp); + _messageHandlerExecutionStrategyFactory ??= sp => ActivatorUtilities.CreateInstance(sp); + _consumerScopeFactory ??= provider => { - _configurationSource = configurationSource; - return this; - } + var loggerFactory = provider.GetRequiredService(); - public ConsumerConfigurationBuilder WithNamingConvention(Func converter) - { - _namingConventions.Add(NamingConvention.UseCustom(converter)); - return this; - } + return new KafkaBasedConsumerScopeFactory( + loggerFactory: loggerFactory, + configuration: configurations, + topics: _messageHandlerRegistry.GetAllSubscribedTopics(), + incomingMessageFactory: _incomingMessageFactory(provider), + readFromBeginning: _readFromBeginning + ); + }; + } - internal ConsumerConfigurationBuilder WithNamingConvention(NamingConvention namingConvention) - { - _namingConventions.Add(namingConvention); - return this; - } + public ConsumerConfigurationBuilder WithConfigurationSource(ConfigurationSource configurationSource) + { + _configurationSource = configurationSource; + return this; + } - public ConsumerConfigurationBuilder WithEnvironmentStyle(string prefix = null, params string[] additionalPrefixes) - { - WithNamingConvention(NamingConvention.UseEnvironmentStyle(prefix)); + public ConsumerConfigurationBuilder WithNamingConvention(Func converter) + { + _namingConventions.Add(NamingConvention.UseCustom(converter)); + return this; + } - foreach (var additionalPrefix in additionalPrefixes) - { - WithNamingConvention(NamingConvention.UseEnvironmentStyle(additionalPrefix)); - } + internal ConsumerConfigurationBuilder WithNamingConvention(NamingConvention namingConvention) + { + _namingConventions.Add(namingConvention); + return this; + } - return this; - } + public ConsumerConfigurationBuilder WithEnvironmentStyle(string prefix = null, params string[] additionalPrefixes) + { + WithNamingConvention(NamingConvention.UseEnvironmentStyle(prefix)); - public ConsumerConfigurationBuilder WithConfiguration(string key, string value) + foreach (var additionalPrefix in additionalPrefixes) { - _configurations[key] = value; - return this; + WithNamingConvention(NamingConvention.UseEnvironmentStyle(additionalPrefix)); } - public ConsumerConfigurationBuilder WithGroupId(string groupId) - { - return WithConfiguration(ConfigurationKey.GroupId, groupId); - } + return this; + } - public ConsumerConfigurationBuilder WithBootstrapServers(string bootstrapServers) - { - return WithConfiguration(ConfigurationKey.BootstrapServers, bootstrapServers); - } + public ConsumerConfigurationBuilder WithConfiguration(string key, string value) + { + _configurations[key] = value; + return this; + } - public ConsumerConfigurationBuilder WithUnitOfWorkFactory(Func handlerUnitOfWorkFactory) - { - _handlerUnitOfWorkFactory = handlerUnitOfWorkFactory; - return this; - } + public ConsumerConfigurationBuilder WithGroupId(string groupId) + { + return WithConfiguration(ConfigurationKey.GroupId, groupId); + } - internal ConsumerConfigurationBuilder WithConsumerScopeFactory(Func consumerScopeFactory) - { - _consumerScopeFactory = consumerScopeFactory; - return this; - } + public ConsumerConfigurationBuilder WithBootstrapServers(string bootstrapServers) + { + return WithConfiguration(ConfigurationKey.BootstrapServers, bootstrapServers); + } - public ConsumerConfigurationBuilder WithUnconfiguredMessageHandlingStrategy(Func unconfiguredMessageHandlingStrategy) - { - _unconfiguredMessageHandlingStrategy = unconfiguredMessageHandlingStrategy; - return this; - } + public ConsumerConfigurationBuilder WithUnitOfWorkFactory(Func handlerUnitOfWorkFactory) + { + _handlerUnitOfWorkFactory = handlerUnitOfWorkFactory; + return this; + } + + internal ConsumerConfigurationBuilder WithConsumerScopeFactory(Func consumerScopeFactory) + { + _consumerScopeFactory = consumerScopeFactory; + return this; + } + + public ConsumerConfigurationBuilder WithUnconfiguredMessageHandlingStrategy(Func unconfiguredMessageHandlingStrategy) + { + _unconfiguredMessageHandlingStrategy = unconfiguredMessageHandlingStrategy; + return this; + } - public ConsumerConfigurationBuilder ReadFromBeginning() - { - _readFromBeginning = true; - return this; - } + public ConsumerConfigurationBuilder ReadFromBeginning() + { + _readFromBeginning = true; + return this; + } - public void WithMessageFilter(MessageFilter messageFilter) - { - _messageFilter = messageFilter; - } + public void WithMessageFilter(MessageFilter messageFilter) + { + _messageFilter = messageFilter; + } - public ConsumerConfigurationBuilder RegisterMessageHandler(string topic, string messageType) - where TMessageHandler : IMessageHandler - { - _messageHandlerRegistry.Register(topic, messageType); - return this; - } + public ConsumerConfigurationBuilder RegisterMessageHandler(string topic, string messageType) + where TMessageHandler : IMessageHandler + { + _messageHandlerRegistry.Register(topic, messageType); + return this; + } - public ConsumerConfigurationBuilder WithIncomingMessageFactory(Func incomingMessageFactory) - { - _incomingMessageFactory = incomingMessageFactory; - return this; - } + public ConsumerConfigurationBuilder WithIncomingMessageFactory(Func incomingMessageFactory) + { + _incomingMessageFactory = incomingMessageFactory; + return this; + } - public ConsumerConfigurationBuilder WithPoisonMessageHandling() - { - var inner = _incomingMessageFactory; - _incomingMessageFactory = provider => new PoisonAwareIncomingMessageFactory( - provider.GetRequiredService>(), - inner(provider) - ); - return this; - } + public ConsumerConfigurationBuilder WithPoisonMessageHandling() + { + var inner = _incomingMessageFactory; + _incomingMessageFactory = provider => new PoisonAwareIncomingMessageFactory( + provider.GetRequiredService>(), + inner(provider) + ); + return this; + } - public ConsumerConfigurationBuilder WithConsumerErrorHandler(Func> failureEvaluation) - { - _consumerErrorHandler = new ConsumerErrorHandler(failureEvaluation); - return this; - } + public ConsumerConfigurationBuilder WithConsumerErrorHandler(Func> failureEvaluation) + { + _consumerErrorHandler = new ConsumerErrorHandler(failureEvaluation); + return this; + } - public ConsumerConfigurationBuilder WithMessageHandlerExecutionStrategyFactory(Func factory) - { - _messageHandlerExecutionStrategyFactory = factory; - return this; - } + public ConsumerConfigurationBuilder WithMessageHandlerExecutionStrategyFactory(Func factory) + { + _messageHandlerExecutionStrategyFactory = factory; + return this; + } - internal ConsumerConfiguration Build() - { - var configurations = new ConfigurationBuilder() - .WithConfigurationKeys(DefaultConfigurationKeys) - .WithRequiredConfigurationKeys(RequiredConfigurationKeys) - .WithNamingConventions(_namingConventions.ToArray()) - .WithConfigurationSource(_configurationSource) - .WithConfigurations(_configurations) - .Build(); - - ApplyDefaults(configurations); - - var consumerConfigurationFactories = new ConsumerConfigurationFactories( - UnitOfWorkFactory: _handlerUnitOfWorkFactory, - UnconfiguredMessageHandlingStrategy: _unconfiguredMessageHandlingStrategy, - ConsumerScopeFactory: _consumerScopeFactory, - IncomingMessageFactory: _incomingMessageFactory, - MessageHandlerExecutionStrategyFactory: _messageHandlerExecutionStrategyFactory); + public DeadLetterQueueOptions WithDeadLetterQueue(string topicName = null) + { + _deadLetterQueueOptions = new DeadLetterQueueOptions(topicName); + return _deadLetterQueueOptions; + } + + internal ConsumerConfiguration Build() + { + var configurations = new ConfigurationBuilder() + .WithConfigurationKeys(DefaultConfigurationKeys) + .WithRequiredConfigurationKeys(RequiredConfigurationKeys) + .WithNamingConventions(_namingConventions.ToArray()) + .WithConfigurationSource(_configurationSource) + .WithConfigurations(_configurations) + .Build(); - return new ConsumerConfiguration( - configuration: configurations, - messageHandlerRegistry: _messageHandlerRegistry, - factories: consumerConfigurationFactories, - messageFilter: _messageFilter, - consumerErrorHandler: _consumerErrorHandler); + ApplyDefaults(configurations); + + var consumerConfigurationFactories = new ConsumerConfigurationFactories( + UnitOfWorkFactory: _handlerUnitOfWorkFactory, + UnconfiguredMessageHandlingStrategy: _unconfiguredMessageHandlingStrategy, + ConsumerScopeFactory: _consumerScopeFactory, + IncomingMessageFactory: _incomingMessageFactory, + MessageHandlerExecutionStrategyFactory: _messageHandlerExecutionStrategyFactory); + + var deadLetterQueueFactory = BuildDeadLetterQueueFactory(configurations); + var maxRetries = _deadLetterQueueOptions?.MaxRetries ?? 0; + + return new ConsumerConfiguration( + configuration: configurations, + messageHandlerRegistry: _messageHandlerRegistry, + factories: consumerConfigurationFactories, + messageFilter: _messageFilter, + consumerErrorHandler: _consumerErrorHandler, + deadLetterQueueFactory: deadLetterQueueFactory, + maxRetries: maxRetries); + } + + private Func BuildDeadLetterQueueFactory(IDictionary configurations) + { + if (_deadLetterQueueOptions == null) + { + return _ => NullDeadLetterQueue.Instance; } + + var topicName = _deadLetterQueueOptions.TopicName; + + var producerConfiguration = configurations + .Where(pair => ConfigurationKey.GetAllProducerKeys().Any(key => key.ToString() == pair.Key)) + .ToDictionary(pair => pair.Key, pair => pair.Value); + + return provider => new KafkaDeadLetterQueue( + provider.GetRequiredService(), + producerConfiguration, + topicName); } -} +} \ No newline at end of file diff --git a/src/Dafda/Configuration/ConsumerOptions.cs b/src/Dafda/Configuration/ConsumerOptions.cs index 498ab188..8d47897f 100644 --- a/src/Dafda/Configuration/ConsumerOptions.cs +++ b/src/Dafda/Configuration/ConsumerOptions.cs @@ -1,10 +1,11 @@ +namespace Dafda.Configuration; + using System; using System.Threading.Tasks; -using Dafda.Consuming; -using Dafda.Consuming.MessageFilters; +using Consuming; +using Consuming.MessageFilters; using Microsoft.Extensions.DependencyInjection; -namespace Dafda.Configuration; /// /// Facilitates Dafda configuration in .NET applications using the . @@ -229,6 +230,26 @@ public void WithMessageHandlerExecutionStrategyFactory( Builder.WithMessageHandlerExecutionStrategyFactory(factory); } + /// + /// Enable a dead letter queue for this consumer. When a message handler fails + /// (after any configured retries are exhausted), the original message is + /// published to the dead letter topic and the offset is committed so the + /// consumer can continue past the poison message. + /// + /// + /// The dead letter topic name. When omitted, the topic is derived from the + /// source topic of the failed message (e.g. orders becomes + /// orders.dead-letter). + /// + /// + /// A to further configure the dead letter + /// queue, e.g. WithMaxRetries(...). + /// + public DeadLetterQueueOptions WithDeadLetterQueue(string topicName = null) + { + return Builder.WithDeadLetterQueue(topicName); + } + private class DefaultConfigurationSource(Microsoft.Extensions.Configuration.IConfiguration configuration) : ConfigurationSource { diff --git a/src/Dafda/Configuration/ConsumerServiceCollectionExtensions.cs b/src/Dafda/Configuration/ConsumerServiceCollectionExtensions.cs index 8b3b198f..63e0dfbc 100644 --- a/src/Dafda/Configuration/ConsumerServiceCollectionExtensions.cs +++ b/src/Dafda/Configuration/ConsumerServiceCollectionExtensions.cs @@ -44,7 +44,9 @@ public static void AddConsumer(this IServiceCollection services, Action +/// Fluent options for configuring a dead letter queue on a consumer. +/// Returned by . +/// +public sealed class DeadLetterQueueOptions +{ + internal DeadLetterQueueOptions(string topicName) + { + TopicName = topicName; + } + + /// + /// The explicit dead letter topic name. When null, the topic is + /// derived from the source topic of the failed message. + /// + internal string TopicName { get; } + + /// + /// The number of additional delivery attempts after the first failure, + /// before the message is forwarded to the dead letter queue. + /// + internal int MaxRetries { get; private set; } + + /// + /// Set the maximum number of retries (additional attempts after the first + /// delivery) before a failing message is sent to the dead letter queue. + /// + /// The number of retries. Must be zero or greater. + public DeadLetterQueueOptions WithMaxRetries(int maxRetries) + { + if (maxRetries < 0) + { + throw new InvalidConfigurationException("The number of retries for a dead letter queue cannot be negative."); + } + + MaxRetries = maxRetries; + return this; + } +} \ No newline at end of file diff --git a/src/Dafda/Consuming/Consumer.cs b/src/Dafda/Consuming/Consumer.cs index c1d1d61b..0ca3943a 100644 --- a/src/Dafda/Consuming/Consumer.cs +++ b/src/Dafda/Consuming/Consumer.cs @@ -1,74 +1,90 @@ +namespace Dafda.Consuming; + using System; using System.Threading; using System.Threading.Tasks; -using Dafda.Consuming.Interfaces; -using Dafda.Consuming.MessageFilters; -using Dafda.Diagnostics; +using Diagnostics; +using Interfaces; +using MessageFilters; -namespace Dafda.Consuming +internal class Consumer( + MessageHandlerRegistry messageHandlerRegistry, + IHandlerUnitOfWorkFactory unitOfWorkFactory, + IConsumerScopeFactory consumerScopeFactory, + IUnconfiguredMessageHandlingStrategy fallbackHandler, + MessageFilter messageFilter, + IMessageHandlerExecutionStrategy messageHandlerExecutionStrategy, + bool isAutoCommitEnabled = false, + IDeadLetterQueue deadLetterQueue = null, + int maxRetries = 0) + : IConsumer, IDisposable { - internal class Consumer : IConsumer - { - private readonly LocalMessageDispatcher _localMessageDispatcher; - private readonly IConsumerScopeFactory _consumerScopeFactory; - private readonly MessageFilter _messageFilter; - private readonly bool _isAutoCommitEnabled; + private readonly LocalMessageDispatcher _localMessageDispatcher = new( + messageHandlerRegistry, + unitOfWorkFactory, + fallbackHandler, + messageHandlerExecutionStrategy); - public Consumer( - MessageHandlerRegistry messageHandlerRegistry, - IHandlerUnitOfWorkFactory unitOfWorkFactory, - IConsumerScopeFactory consumerScopeFactory, - IUnconfiguredMessageHandlingStrategy fallbackHandler, - MessageFilter messageFilter, - IMessageHandlerExecutionStrategy messageHandlerExecutionStrategy, - bool isAutoCommitEnabled = false) + private readonly IDeadLetterQueue _deadLetterQueue = deadLetterQueue ?? NullDeadLetterQueue.Instance; + + public async Task ConsumeAll(CancellationToken cancellationToken) + { + using var consumerScope = consumerScopeFactory.CreateConsumerScope(); + while (!cancellationToken.IsCancellationRequested) { - _localMessageDispatcher = - new LocalMessageDispatcher( - messageHandlerRegistry, - unitOfWorkFactory, - fallbackHandler, - messageHandlerExecutionStrategy); - _consumerScopeFactory = - consumerScopeFactory - ?? throw new ArgumentNullException(nameof(consumerScopeFactory)); - _messageFilter = messageFilter; - _isAutoCommitEnabled = isAutoCommitEnabled; + await ProcessNextMessage(consumerScope, cancellationToken); } + } + + public async Task ConsumeSingle(CancellationToken cancellationToken) + { + using var consumerScope = consumerScopeFactory.CreateConsumerScope(); + await ProcessNextMessage(consumerScope, cancellationToken); + } + + private async Task ProcessNextMessage(ConsumerScope consumerScope, CancellationToken cancellationToken) + { + var messageResult = await consumerScope.GetNext(cancellationToken); + using var activity = DafdaActivitySource.StartReceivingActivity(messageResult); - public async Task ConsumeAll(CancellationToken cancellationToken) + if (messageFilter.CanAcceptMessage(messageResult)) { - using (var consumerScope = _consumerScopeFactory.CreateConsumerScope()) - { - while (!cancellationToken.IsCancellationRequested) - { - await ProcessNextMessage(consumerScope, cancellationToken); - } - } + await Dispatch(messageResult, cancellationToken); } - public async Task ConsumeSingle(CancellationToken cancellationToken) + if (!isAutoCommitEnabled) { - using (var consumerScope = _consumerScopeFactory.CreateConsumerScope()) - { - await ProcessNextMessage(consumerScope, cancellationToken); - } + await messageResult.Commit(cancellationToken); } + } - private async Task ProcessNextMessage(ConsumerScope consumerScope, CancellationToken cancellationToken) - { - var messageResult = await consumerScope.GetNext(cancellationToken); - using var activity = DafdaActivitySource.StartReceivingActivity(messageResult); + private async Task Dispatch(MessageResult messageResult, CancellationToken cancellationToken) + { + var deadLetterQueueEnabled = _deadLetterQueue is not NullDeadLetterQueue; + var attempt = 0; - if (_messageFilter.CanAcceptMessage(messageResult)) + while (true) + { + try { await _localMessageDispatcher.Dispatch(messageResult, cancellationToken); + return; } - - if (!_isAutoCommitEnabled) + catch (Exception exception) when (deadLetterQueueEnabled && !cancellationToken.IsCancellationRequested) { - await messageResult.Commit(cancellationToken); + if (attempt++ < maxRetries) + { + continue; + } + + await _deadLetterQueue.Send(messageResult, exception, cancellationToken); + return; } } } + + public void Dispose() + { + (_deadLetterQueue as IDisposable)?.Dispose(); + } } \ No newline at end of file diff --git a/src/Dafda/Consuming/ConsumerHostedService.cs b/src/Dafda/Consuming/ConsumerHostedService.cs index ae020b88..a6266b24 100644 --- a/src/Dafda/Consuming/ConsumerHostedService.cs +++ b/src/Dafda/Consuming/ConsumerHostedService.cs @@ -61,5 +61,12 @@ protected override Task ExecuteAsync(CancellationToken stoppingToken) { return Task.Run(async () => { await ConsumeAll(stoppingToken); }, stoppingToken); } + + /// Disposes the underlying consumer (and any resources it owns, such as a dead letter queue producer). + public override void Dispose() + { + (_consumer as IDisposable)?.Dispose(); + base.Dispose(); + } } } \ No newline at end of file diff --git a/src/Dafda/Consuming/IDeadLetterQueue.cs b/src/Dafda/Consuming/IDeadLetterQueue.cs new file mode 100644 index 00000000..482a0b46 --- /dev/null +++ b/src/Dafda/Consuming/IDeadLetterQueue.cs @@ -0,0 +1,17 @@ +namespace Dafda.Consuming; + +using System; +using System.Threading; +using System.Threading.Tasks; + +/// +/// Forwards a message that could not be handled (after retries are exhausted) +/// to a dead letter queue, so the consumer can move past the poison message. +/// +internal interface IDeadLetterQueue +{ + /// + /// Publish the failed to the dead letter queue. + /// + Task Send(MessageResult message, Exception exception, CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/src/Dafda/Consuming/KafkaConsumerScope.cs b/src/Dafda/Consuming/KafkaConsumerScope.cs index 340b0edf..af8de408 100644 --- a/src/Dafda/Consuming/KafkaConsumerScope.cs +++ b/src/Dafda/Consuming/KafkaConsumerScope.cs @@ -1,53 +1,53 @@ +namespace Dafda.Consuming; + using System.Threading; using System.Threading.Tasks; using Confluent.Kafka; using Microsoft.Extensions.Logging; -namespace Dafda.Consuming +internal class KafkaConsumerScope : ConsumerScope { - internal class KafkaConsumerScope : ConsumerScope + private readonly ILogger _logger; + private readonly IConsumer _innerKafkaConsumer; + private readonly IIncomingMessageFactory _incomingMessageFactory; + private readonly string _groupId; + + internal KafkaConsumerScope(ILoggerFactory loggerFactory, IConsumer innerKafkaConsumer, IIncomingMessageFactory incomingMessageFactory, string groupId) { - private readonly ILogger _logger; - private readonly IConsumer _innerKafkaConsumer; - private readonly IIncomingMessageFactory _incomingMessageFactory; - private readonly string _groupId; + _logger = loggerFactory.CreateLogger(); + _innerKafkaConsumer = innerKafkaConsumer; + _incomingMessageFactory = incomingMessageFactory; + _groupId = groupId; + } - internal KafkaConsumerScope(ILoggerFactory loggerFactory, IConsumer innerKafkaConsumer, IIncomingMessageFactory incomingMessageFactory, string groupId) - { - _logger = loggerFactory.CreateLogger(); - _innerKafkaConsumer = innerKafkaConsumer; - _incomingMessageFactory = incomingMessageFactory; - _groupId = groupId; - } + public override Task GetNext(CancellationToken cancellationToken) + { + var innerResult = _innerKafkaConsumer.Consume(cancellationToken); - public override Task GetNext(CancellationToken cancellationToken) - { - var innerResult = _innerKafkaConsumer.Consume(cancellationToken); - - _logger.LogDebug("Received message {Key}: {RawMessage}", innerResult.Message?.Key, innerResult.Message?.Value); - var result = new MessageResult( - message: _incomingMessageFactory.Create(innerResult.Message.Value), - onCommit: (CancellationToken cancellationToken) => - { - _innerKafkaConsumer.Commit(innerResult); - return Task.CompletedTask; - }) + _logger.LogDebug("Received message {Key}: {RawMessage}", innerResult.Message?.Key, innerResult.Message?.Value); + var result = new MessageResult( + message: _incomingMessageFactory.Create(innerResult.Message.Value), + onCommit: _ => { + _innerKafkaConsumer.Commit(innerResult); + return Task.CompletedTask; + }) + { - Topic = innerResult.Topic, - Partition = innerResult.Partition.Value, - PartitionKey = innerResult.Message.Key, - ClientId = _innerKafkaConsumer.Name, - GroupId = _groupId - }; + Topic = innerResult.Topic, + RawMessage = innerResult.Message.Value, + Partition = innerResult.Partition.Value, + PartitionKey = innerResult.Message.Key, + ClientId = _innerKafkaConsumer.Name, + GroupId = _groupId + }; - return Task.FromResult(result); - } + return Task.FromResult(result); + } - public override void Dispose() - { - _innerKafkaConsumer.Close(); - _innerKafkaConsumer.Dispose(); - } + public override void Dispose() + { + _innerKafkaConsumer.Close(); + _innerKafkaConsumer.Dispose(); } } \ No newline at end of file diff --git a/src/Dafda/Consuming/KafkaDeadLetterQueue.cs b/src/Dafda/Consuming/KafkaDeadLetterQueue.cs new file mode 100644 index 00000000..7ca492da --- /dev/null +++ b/src/Dafda/Consuming/KafkaDeadLetterQueue.cs @@ -0,0 +1,83 @@ +namespace Dafda.Consuming; + +using System; +using System.Collections.Generic; +using System.Text; +using System.Threading; +using System.Threading.Tasks; +using Confluent.Kafka; +using Microsoft.Extensions.Logging; + +/// +/// Publishes failed messages to a Kafka dead letter topic. When no explicit +/// topic name is configured, the target topic is derived from the source +/// topic of the message using the . +/// +internal sealed class KafkaDeadLetterQueue( + ILoggerFactory loggerFactory, + IEnumerable> configuration, + string topicName) + : IDeadLetterQueue, IDisposable +{ + private const string DefaultTopicSuffix = ".dead-letter"; + + private readonly ILogger _logger = loggerFactory.CreateLogger(); + private readonly IProducer _producer = new ProducerBuilder(configuration).Build(); + + private string ResolveTopic(MessageResult message) + { + return ResolveTopicName(topicName, message.Topic, message.GroupId); + } + + internal static string ResolveTopicName(string configuredTopicName, string sourceTopic, string groupId) + { + if (configuredTopicName != null) + { + return configuredTopicName; + } + + return string.IsNullOrEmpty(groupId) + ? $"{sourceTopic}{DefaultTopicSuffix}" + : $"{sourceTopic}.{groupId}{DefaultTopicSuffix}"; + } + + public async Task Send(MessageResult message, Exception exception, CancellationToken cancellationToken) + { + var topic = ResolveTopic(message); + + _logger.LogWarning( + exception, + "Dead-lettering message with key {Key} from topic {SourceTopic} to {DeadLetterTopic}", + message.PartitionKey, + message.Topic, + topic); + + var headers = new Headers + { + { "dafda-dead-letter-source-topic", Encode(message.Topic) }, + { "dafda-dead-letter-exception-type", Encode(exception.GetType().FullName) }, + { "dafda-dead-letter-exception-message", Encode(exception.Message) }, + { "dafda-dead-letter-timestamp", Encode(DateTimeOffset.UtcNow.ToString("O")) }, + }; + + await _producer.ProduceAsync( + topic: topic, + message: new Message + { + Key = message.PartitionKey, + Value = message.RawMessage, + Headers = headers + }, + cancellationToken); + } + + private static byte[] Encode(string value) + { + return Encoding.UTF8.GetBytes(value ?? string.Empty); + } + + public void Dispose() + { + _producer?.Dispose(); + } +} \ No newline at end of file diff --git a/src/Dafda/Consuming/MessageResult.cs b/src/Dafda/Consuming/MessageResult.cs index 80da639e..f0ef93c6 100644 --- a/src/Dafda/Consuming/MessageResult.cs +++ b/src/Dafda/Consuming/MessageResult.cs @@ -1,48 +1,53 @@ +namespace Dafda.Consuming; + using System; using System.Threading; using System.Threading.Tasks; -namespace Dafda.Consuming +/// +/// Object that contains message when consumed from Kafka. +/// To be used for message handling prior to dispatching to the handlers. +/// +public class MessageResult { + private static readonly Func EmptyCommitAction = (_) => Task.CompletedTask; + private readonly Func _onCommit; + /// - /// Object that contains message when consumed from Kafka. - /// To be used for message handling prior to dispatching to the handlers. + /// Resulting Message containing Transport Level Message /// - public class MessageResult + public MessageResult(TransportLevelMessage message, Func onCommit = null) { - private static readonly Func EmptyCommitAction = (_) => Task.CompletedTask; - private readonly Func _onCommit; + Message = message; + _onCommit = onCommit ?? EmptyCommitAction; + } - /// - /// Resulting Message contaning Transport Level Message - /// - public MessageResult(TransportLevelMessage message, Func onCommit = null) - { - Message = message; - _onCommit = onCommit ?? EmptyCommitAction; - } + internal string Topic { get; set; } - internal string Topic { get; set; } + /// + /// The raw, unparsed message value as received from Kafka. Used when + /// forwarding a failed message to a dead letter queue. + /// + public string RawMessage { get; internal set; } - internal int Partition { get; set; } + internal int Partition { get; set; } - internal string PartitionKey { get; set; } + internal string PartitionKey { get; set; } - internal string ClientId { get; set; } + internal string ClientId { get; set; } - internal string GroupId { get; set; } + internal string GroupId { get; set; } - /// - /// Transmitted message consumed from Kafka - /// - public TransportLevelMessage Message { get; } + /// + /// Transmitted message consumed from Kafka + /// + public TransportLevelMessage Message { get; } - /// - /// Commit message to handlers - /// - public async Task Commit(CancellationToken cancellationToken) - { - await _onCommit(cancellationToken); - } + /// + /// Commit message to handlers + /// + public async Task Commit(CancellationToken cancellationToken) + { + await _onCommit(cancellationToken); } } \ No newline at end of file diff --git a/src/Dafda/Consuming/NullDeadLetterQueue.cs b/src/Dafda/Consuming/NullDeadLetterQueue.cs new file mode 100644 index 00000000..8ac3c042 --- /dev/null +++ b/src/Dafda/Consuming/NullDeadLetterQueue.cs @@ -0,0 +1,24 @@ +namespace Dafda.Consuming; + +using System; +using System.Threading; +using System.Threading.Tasks; + +/// +/// A no-op used when no dead letter queue is +/// configured. Its presence signals that failed messages should be rethrown +/// (preserving the pre-existing behavior). +/// +internal sealed class NullDeadLetterQueue : IDeadLetterQueue +{ + public static readonly NullDeadLetterQueue Instance = new(); + + private NullDeadLetterQueue() + { + } + + public Task Send(MessageResult message, Exception exception, CancellationToken cancellationToken) + { + return Task.CompletedTask; + } +} \ No newline at end of file