diff --git a/src/Dafda.Tests/Builders/ConsumerBuilder.cs b/src/Dafda.Tests/Builders/ConsumerBuilder.cs index d0861fd..bbbc094 100644 --- a/src/Dafda.Tests/Builders/ConsumerBuilder.cs +++ b/src/Dafda.Tests/Builders/ConsumerBuilder.cs @@ -1,5 +1,6 @@ namespace Dafda.Tests.Builders; +using System; using Dafda.Consuming; using Dafda.Consuming.MessageFilters; using TestDoubles; @@ -16,6 +17,7 @@ internal class ConsumerBuilder private MessageFilter _messageFilter = MessageFilter.Default; private IDeadLetterQueue _deadLetterQueue = NullDeadLetterQueue.Instance; private int _maxRetries; + private Func _deadLetterQueueBypass; public ConsumerBuilder WithUnitOfWork(IHandlerUnitOfWork unitOfWork) { @@ -70,6 +72,12 @@ public ConsumerBuilder WithMaxRetries(int maxRetries) return this; } + public ConsumerBuilder WithDeadLetterQueueBypass(Func deadLetterQueueBypass) + { + _deadLetterQueueBypass = deadLetterQueueBypass; + return this; + } + public Consumer Build() => new Consumer( _registry, @@ -80,5 +88,6 @@ public Consumer Build() => _messageHandlerExecutionStrategy, _enableAutoCommit, _deadLetterQueue, - _maxRetries); + _maxRetries, + _deadLetterQueueBypass); } \ No newline at end of file diff --git a/src/Dafda.Tests/Configuration/TestDeadLetterQueueOptions.cs b/src/Dafda.Tests/Configuration/TestDeadLetterQueueOptions.cs new file mode 100644 index 0000000..242957b --- /dev/null +++ b/src/Dafda.Tests/Configuration/TestDeadLetterQueueOptions.cs @@ -0,0 +1,69 @@ +namespace Dafda.Tests.Configuration; + +using System; +using Dafda.Configuration; +using Xunit; + +public class TestDeadLetterQueueOptions +{ + [Fact] + public void has_no_bypass_predicate_by_default() + { + var sut = new DeadLetterQueueOptions("dlq"); + + Assert.Null(sut.BypassPredicate); + } + + [Fact] + public void bypass_predicate_matches_the_registered_exception_type() + { + var sut = new DeadLetterQueueOptions("dlq") + .BypassFor(); + + Assert.True(sut.BypassPredicate(new InvalidOperationException())); + Assert.False(sut.BypassPredicate(new FormatException())); + } + + [Fact] + public void bypass_predicate_matches_types_derived_from_the_registered_exception_type() + { + var sut = new DeadLetterQueueOptions("dlq") + .BypassFor(); + + Assert.True(sut.BypassPredicate(new ArgumentNullException())); + } + + [Fact] + public void bypass_predicate_matches_any_of_the_registered_predicates() + { + var sut = new DeadLetterQueueOptions("dlq") + .BypassFor() + .BypassWhen(exception => exception is FormatException); + + Assert.True(sut.BypassPredicate(new InvalidOperationException())); + Assert.True(sut.BypassPredicate(new FormatException())); + Assert.False(sut.BypassPredicate(new NotSupportedException())); + } + + [Fact] + public void throws_when_the_bypass_predicate_is_null() + { + var sut = new DeadLetterQueueOptions("dlq"); + + Assert.Throws(() => sut.BypassWhen(null)); + } + + [Fact] + public void bypass_predicate_is_not_affected_by_predicates_registered_afterwards() + { + var sut = new DeadLetterQueueOptions("dlq") + .BypassFor(); + + var predicate = sut.BypassPredicate; + + sut.BypassFor(); + + Assert.False(predicate(new FormatException())); + Assert.True(sut.BypassPredicate(new FormatException())); + } +} diff --git a/src/Dafda.Tests/Consuming/TestConsumer.cs b/src/Dafda.Tests/Consuming/TestConsumer.cs index ad9df7c..cdfb21a 100644 --- a/src/Dafda.Tests/Consuming/TestConsumer.cs +++ b/src/Dafda.Tests/Consuming/TestConsumer.cs @@ -469,11 +469,61 @@ public void disposing_consumer_disposes_the_dead_letter_queue() Assert.Equal(1, deadLetterQueueSpy.DisposedCount); } + [Fact] + public async Task propagates_exception_and_bypasses_dead_letter_queue_for_bypassed_exception_type() + { + var handlerInvocations = 0; + var handler = new MessageHandlerSpy(() => + { + handlerInvocations++; + throw new InvalidOperationException("fatal"); + }); + + var deadLetterQueueSpy = new DeadLetterQueueSpy(); + var committed = false; + + var sut = BuildConsumerWithHandler( + handler, + onCommit: _ => + { + committed = true; + return Task.CompletedTask; + }, + deadLetterQueue: deadLetterQueueSpy, + maxRetries: 3, + deadLetterQueueBypass: exception => exception is InvalidOperationException); + + await Assert.ThrowsAsync( + () => sut.ConsumeSingle(CancellationToken.None)); + + Assert.Equal(1, handlerInvocations); + Assert.Equal(0, deadLetterQueueSpy.SendCount); + Assert.False(committed); + } + + [Fact] + public async Task dead_letters_exceptions_that_do_not_match_the_bypass() + { + var handler = new MessageHandlerSpy(() => throw new InvalidOperationException("boom")); + + var deadLetterQueueSpy = new DeadLetterQueueSpy(); + + var sut = BuildConsumerWithHandler( + handler, + deadLetterQueue: deadLetterQueueSpy, + deadLetterQueueBypass: exception => exception is FormatException); + + await sut.ConsumeSingle(CancellationToken.None); + + Assert.Equal(1, deadLetterQueueSpy.SendCount); + } + private static Consumer BuildConsumerWithHandler( IMessageHandler handler, Func onCommit = null, IDeadLetterQueue deadLetterQueue = null, - int maxRetries = 0) + int maxRetries = 0, + Func deadLetterQueueBypass = null) { var registration = new MessageRegistrationBuilder() .WithHandlerInstanceType(handler.GetType()) @@ -495,7 +545,8 @@ private static Consumer BuildConsumerWithHandler( .WithConsumerScopeFactory(new ConsumerScopeFactoryStub(new ConsumerScopeStub(messageResult))) .WithUnitOfWork(new UnitOfWorkStub(handler)) .WithMessageHandlerRegistry(registry) - .WithMaxRetries(maxRetries); + .WithMaxRetries(maxRetries) + .WithDeadLetterQueueBypass(deadLetterQueueBypass); if (deadLetterQueue != null) { diff --git a/src/Dafda/Configuration/ConsumerConfiguration.cs b/src/Dafda/Configuration/ConsumerConfiguration.cs index 0859c41..b7ac7bc 100644 --- a/src/Dafda/Configuration/ConsumerConfiguration.cs +++ b/src/Dafda/Configuration/ConsumerConfiguration.cs @@ -13,7 +13,8 @@ internal class ConsumerConfiguration( MessageFilter messageFilter, IConsumerErrorHandler consumerErrorHandler, Func deadLetterQueueFactory, - int maxRetries) + int maxRetries, + Func deadLetterQueueBypass) : ConsumerConfigurationBase(configuration, factories.UnitOfWorkFactory, consumerErrorHandler) { public ConsumerConfigurationFactories Factories { get; } = factories; @@ -21,4 +22,5 @@ internal class ConsumerConfiguration( public MessageFilter MessageFilter { get; } = messageFilter; public Func DeadLetterQueueFactory { get; } = deadLetterQueueFactory; public int MaxRetries { get; } = maxRetries; + public Func DeadLetterQueueBypass { get; } = deadLetterQueueBypass; } \ No newline at end of file diff --git a/src/Dafda/Configuration/ConsumerConfigurationBuilder.cs b/src/Dafda/Configuration/ConsumerConfigurationBuilder.cs index b6d0bc2..a10ab1b 100644 --- a/src/Dafda/Configuration/ConsumerConfigurationBuilder.cs +++ b/src/Dafda/Configuration/ConsumerConfigurationBuilder.cs @@ -207,6 +207,7 @@ internal ConsumerConfiguration Build() var deadLetterQueueFactory = BuildDeadLetterQueueFactory(configurations); var maxRetries = _deadLetterQueueOptions?.MaxRetries ?? 0; + var deadLetterQueueBypass = _deadLetterQueueOptions?.BypassPredicate; return new ConsumerConfiguration( configuration: configurations, @@ -215,7 +216,8 @@ internal ConsumerConfiguration Build() messageFilter: _messageFilter, consumerErrorHandler: _consumerErrorHandler, deadLetterQueueFactory: deadLetterQueueFactory, - maxRetries: maxRetries); + maxRetries: maxRetries, + deadLetterQueueBypass: deadLetterQueueBypass); } private Func BuildDeadLetterQueueFactory(IDictionary configurations) diff --git a/src/Dafda/Configuration/ConsumerServiceCollectionExtensions.cs b/src/Dafda/Configuration/ConsumerServiceCollectionExtensions.cs index 63e0dfb..c41a8a9 100644 --- a/src/Dafda/Configuration/ConsumerServiceCollectionExtensions.cs +++ b/src/Dafda/Configuration/ConsumerServiceCollectionExtensions.cs @@ -46,7 +46,8 @@ 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 { + private readonly List> _bypassPredicates = new(); + internal DeadLetterQueueOptions(string topicName) { TopicName = topicName; @@ -38,4 +44,87 @@ public DeadLetterQueueOptions WithMaxRetries(int maxRetries) MaxRetries = maxRetries; return this; } + + /// + /// A predicate matching exceptions that should bypass the dead letter queue. + /// When an exception matches, it is rethrown instead of being retried or + /// forwarded to the dead letter queue. Returns null when no bypass has + /// been configured. + /// + internal Func BypassPredicate + { + get + { + if (_bypassPredicates.Count == 0) + { + return null; + } + + var snapshot = _bypassPredicates.ToArray(); + return exception => snapshot.Any(predicate => predicate(exception)); + } + } + + /// + /// Bypass the dead letter queue for the specified exception type (and any + /// derived types). A matching exception is neither retried nor forwarded to the + /// dead letter queue: it propagates out of message dispatch, and Dafda does not + /// commit the offset for the message. + /// + /// + /// The exception is then passed to the configured consumer error handler (see + /// ). With the default + /// handler, stops the application. + /// If the handler returns + /// the consumer is restarted and the redelivered message fails again, so only + /// combine a bypass with a restart strategy that backs off. + /// + /// Whether the bypassed message is actually redelivered depends on the commit + /// strategy. Dafda only commits the offset itself when enable.auto.commit + /// is false, so manual commits are required for redelivery. With automatic + /// commits (the default) the Kafka client stores and commits offsets on its own, + /// including when the consumer is closed, so a bypassed message may still be + /// marked as consumed and will not be redelivered. + /// + /// + /// The exception type to bypass the dead letter queue for. + public DeadLetterQueueOptions BypassFor() where TException : Exception + { + _bypassPredicates.Add(exception => exception is TException); + return this; + } + + /// + /// Bypass the dead letter queue for exceptions matching the supplied + /// . When it returns true, the exception is + /// neither retried nor forwarded to the dead letter queue: it propagates out of + /// message dispatch, and Dafda does not commit the offset for the message. + /// + /// + /// The exception is then passed to the configured consumer error handler (see + /// ). With the default + /// handler, stops the application. + /// If the handler returns + /// the consumer is restarted and the redelivered message fails again, so only + /// combine a bypass with a restart strategy that backs off. + /// + /// Whether the bypassed message is actually redelivered depends on the commit + /// strategy. Dafda only commits the offset itself when enable.auto.commit + /// is false, so manual commits are required for redelivery. With automatic + /// commits (the default) the Kafka client stores and commits offsets on its own, + /// including when the consumer is closed, so a bypassed message may still be + /// marked as consumed and will not be redelivered. + /// + /// + /// Evaluates a thrown exception and returns true to bypass the dead letter queue. + public DeadLetterQueueOptions BypassWhen(Func predicate) + { + if (predicate == null) + { + throw new InvalidConfigurationException("The dead letter queue bypass predicate cannot be null."); + } + + _bypassPredicates.Add(predicate); + return this; + } } \ No newline at end of file diff --git a/src/Dafda/Consuming/Consumer.cs b/src/Dafda/Consuming/Consumer.cs index 0ca3943..166453b 100644 --- a/src/Dafda/Consuming/Consumer.cs +++ b/src/Dafda/Consuming/Consumer.cs @@ -16,7 +16,8 @@ internal class Consumer( IMessageHandlerExecutionStrategy messageHandlerExecutionStrategy, bool isAutoCommitEnabled = false, IDeadLetterQueue deadLetterQueue = null, - int maxRetries = 0) + int maxRetries = 0, + Func deadLetterQueueBypass = null) : IConsumer, IDisposable { private readonly LocalMessageDispatcher _localMessageDispatcher = new( @@ -70,7 +71,7 @@ private async Task Dispatch(MessageResult messageResult, CancellationToken cance await _localMessageDispatcher.Dispatch(messageResult, cancellationToken); return; } - catch (Exception exception) when (deadLetterQueueEnabled && !cancellationToken.IsCancellationRequested) + catch (Exception exception) when (deadLetterQueueEnabled && !cancellationToken.IsCancellationRequested && !ShouldBypassDeadLetterQueue(exception)) { if (attempt++ < maxRetries) { @@ -83,6 +84,11 @@ private async Task Dispatch(MessageResult messageResult, CancellationToken cance } } + private bool ShouldBypassDeadLetterQueue(Exception exception) + { + return deadLetterQueueBypass != null && deadLetterQueueBypass(exception); + } + public void Dispose() { (_deadLetterQueue as IDisposable)?.Dispose();