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: 10 additions & 1 deletion src/Dafda.Tests/Builders/ConsumerBuilder.cs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
namespace Dafda.Tests.Builders;

using System;
using Dafda.Consuming;
using Dafda.Consuming.MessageFilters;
using TestDoubles;
Expand All @@ -16,6 +17,7 @@ internal class ConsumerBuilder
private MessageFilter _messageFilter = MessageFilter.Default;
private IDeadLetterQueue _deadLetterQueue = NullDeadLetterQueue.Instance;
private int _maxRetries;
private Func<Exception, bool> _deadLetterQueueBypass;

public ConsumerBuilder WithUnitOfWork(IHandlerUnitOfWork unitOfWork)
{
Expand Down Expand Up @@ -70,6 +72,12 @@ public ConsumerBuilder WithMaxRetries(int maxRetries)
return this;
}

public ConsumerBuilder WithDeadLetterQueueBypass(Func<Exception, bool> deadLetterQueueBypass)
{
_deadLetterQueueBypass = deadLetterQueueBypass;
return this;
}

public Consumer Build() =>
new Consumer(
_registry,
Expand All @@ -80,5 +88,6 @@ public Consumer Build() =>
_messageHandlerExecutionStrategy,
_enableAutoCommit,
_deadLetterQueue,
_maxRetries);
_maxRetries,
_deadLetterQueueBypass);
}
69 changes: 69 additions & 0 deletions src/Dafda.Tests/Configuration/TestDeadLetterQueueOptions.cs
Original file line number Diff line number Diff line change
@@ -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<InvalidOperationException>();

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<ArgumentException>();

Assert.True(sut.BypassPredicate(new ArgumentNullException()));
}

[Fact]
public void bypass_predicate_matches_any_of_the_registered_predicates()
{
var sut = new DeadLetterQueueOptions("dlq")
.BypassFor<InvalidOperationException>()
.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<InvalidConfigurationException>(() => sut.BypassWhen(null));
}

[Fact]
public void bypass_predicate_is_not_affected_by_predicates_registered_afterwards()
{
var sut = new DeadLetterQueueOptions("dlq")
.BypassFor<InvalidOperationException>();

var predicate = sut.BypassPredicate;

sut.BypassFor<FormatException>();

Assert.False(predicate(new FormatException()));
Assert.True(sut.BypassPredicate(new FormatException()));
}
}
55 changes: 53 additions & 2 deletions src/Dafda.Tests/Consuming/TestConsumer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<FooMessage>(() =>
{
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<InvalidOperationException>(
() => 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<FooMessage>(() => 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<FooMessage> handler,
Func<CancellationToken, Task> onCommit = null,
IDeadLetterQueue deadLetterQueue = null,
int maxRetries = 0)
int maxRetries = 0,
Func<Exception, bool> deadLetterQueueBypass = null)
{
var registration = new MessageRegistrationBuilder()
.WithHandlerInstanceType(handler.GetType())
Expand All @@ -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)
{
Expand Down
4 changes: 3 additions & 1 deletion src/Dafda/Configuration/ConsumerConfiguration.cs
Original file line number Diff line number Diff line change
Expand Up @@ -13,12 +13,14 @@ internal class ConsumerConfiguration(
MessageFilter messageFilter,
IConsumerErrorHandler consumerErrorHandler,
Func<IServiceProvider, IDeadLetterQueue> deadLetterQueueFactory,
int maxRetries)
int maxRetries,
Func<Exception, bool> deadLetterQueueBypass)
: ConsumerConfigurationBase(configuration, factories.UnitOfWorkFactory, consumerErrorHandler)
{
public ConsumerConfigurationFactories Factories { get; } = factories;
public MessageHandlerRegistry MessageHandlerRegistry { get; } = messageHandlerRegistry;
public MessageFilter MessageFilter { get; } = messageFilter;
public Func<IServiceProvider, IDeadLetterQueue> DeadLetterQueueFactory { get; } = deadLetterQueueFactory;
public int MaxRetries { get; } = maxRetries;
public Func<Exception, bool> DeadLetterQueueBypass { get; } = deadLetterQueueBypass;
}
4 changes: 3 additions & 1 deletion src/Dafda/Configuration/ConsumerConfigurationBuilder.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -215,7 +216,8 @@ internal ConsumerConfiguration Build()
messageFilter: _messageFilter,
consumerErrorHandler: _consumerErrorHandler,
deadLetterQueueFactory: deadLetterQueueFactory,
maxRetries: maxRetries);
maxRetries: maxRetries,
deadLetterQueueBypass: deadLetterQueueBypass);
}

private Func<IServiceProvider, IDeadLetterQueue> BuildDeadLetterQueueFactory(IDictionary<string, string> configurations)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,8 @@ public static void AddConsumer(this IServiceCollection services, Action<Consumer
configuration.Factories.MessageHandlerExecutionStrategyFactory(provider),
configuration.EnableAutoCommit,
configuration.DeadLetterQueueFactory(provider),
configuration.MaxRetries
configuration.MaxRetries,
configuration.DeadLetterQueueBypass
),
configuration.GroupId,
configuration.ConsumerErrorHandler
Expand Down Expand Up @@ -82,7 +83,8 @@ public static void AddConsumer(this IServiceCollection services, Func<IServicePr
configuration.Factories.MessageHandlerExecutionStrategyFactory(provider),
configuration.EnableAutoCommit,
configuration.DeadLetterQueueFactory(provider),
configuration.MaxRetries
configuration.MaxRetries,
configuration.DeadLetterQueueBypass
),
configuration.GroupId,
configuration.ConsumerErrorHandler
Expand Down
89 changes: 89 additions & 0 deletions src/Dafda/Configuration/DeadLetterQueueOptions.cs
Original file line number Diff line number Diff line change
@@ -1,11 +1,17 @@
namespace Dafda.Configuration;

using System;
using System.Collections.Generic;
using System.Linq;

/// <summary>
/// Fluent options for configuring a dead letter queue on a consumer.
/// Returned by <see cref="ConsumerOptions.WithDeadLetterQueue"/>.
/// </summary>
public sealed class DeadLetterQueueOptions
{
private readonly List<Func<Exception, bool>> _bypassPredicates = new();

internal DeadLetterQueueOptions(string topicName)
{
TopicName = topicName;
Expand Down Expand Up @@ -38,4 +44,87 @@ public DeadLetterQueueOptions WithMaxRetries(int maxRetries)
MaxRetries = maxRetries;
return this;
}

/// <summary>
/// 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 <c>null</c> when no bypass has
/// been configured.
/// </summary>
internal Func<Exception, bool> BypassPredicate
{
get
{
if (_bypassPredicates.Count == 0)
{
return null;
}

var snapshot = _bypassPredicates.ToArray();
return exception => snapshot.Any(predicate => predicate(exception));
}
}

/// <summary>
/// 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.
/// </summary>
/// <remarks>
/// The exception is then passed to the configured consumer error handler (see
/// <see cref="ConsumerOptions.WithConsumerErrorHandler"/>). With the default
/// handler, <see cref="ConsumerFailureStrategy.Default"/> stops the application.
/// If the handler returns <see cref="ConsumerFailureStrategy.RestartConsumer"/>
/// the consumer is restarted and the redelivered message fails again, so only
/// combine a bypass with a restart strategy that backs off.
/// <para>
/// Whether the bypassed message is actually redelivered depends on the commit
/// strategy. Dafda only commits the offset itself when <c>enable.auto.commit</c>
/// is <c>false</c>, 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.
/// </para>
/// </remarks>
/// <typeparam name="TException">The exception type to bypass the dead letter queue for.</typeparam>
public DeadLetterQueueOptions BypassFor<TException>() where TException : Exception
{
_bypassPredicates.Add(exception => exception is TException);
return this;
}

/// <summary>
/// Bypass the dead letter queue for exceptions matching the supplied
/// <paramref name="predicate"/>. When it returns <c>true</c>, 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.
/// </summary>
/// <remarks>
/// The exception is then passed to the configured consumer error handler (see
/// <see cref="ConsumerOptions.WithConsumerErrorHandler"/>). With the default
/// handler, <see cref="ConsumerFailureStrategy.Default"/> stops the application.
/// If the handler returns <see cref="ConsumerFailureStrategy.RestartConsumer"/>
/// the consumer is restarted and the redelivered message fails again, so only
/// combine a bypass with a restart strategy that backs off.
/// <para>
/// Whether the bypassed message is actually redelivered depends on the commit
/// strategy. Dafda only commits the offset itself when <c>enable.auto.commit</c>
/// is <c>false</c>, 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.
/// </para>
/// </remarks>
/// <param name="predicate">Evaluates a thrown exception and returns <c>true</c> to bypass the dead letter queue.</param>
public DeadLetterQueueOptions BypassWhen(Func<Exception, bool> predicate)
{
if (predicate == null)
{
throw new InvalidConfigurationException("The dead letter queue bypass predicate cannot be null.");
}

_bypassPredicates.Add(predicate);
return this;
}
}
10 changes: 8 additions & 2 deletions src/Dafda/Consuming/Consumer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,8 @@ internal class Consumer(
IMessageHandlerExecutionStrategy messageHandlerExecutionStrategy,
bool isAutoCommitEnabled = false,
IDeadLetterQueue deadLetterQueue = null,
int maxRetries = 0)
int maxRetries = 0,
Func<Exception, bool> deadLetterQueueBypass = null)
: IConsumer, IDisposable
{
private readonly LocalMessageDispatcher _localMessageDispatcher = new(
Expand Down Expand Up @@ -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))
Comment thread
nabisobhi marked this conversation as resolved.
{
Comment thread
nabisobhi marked this conversation as resolved.
if (attempt++ < maxRetries)
{
Expand All @@ -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();
Expand Down
Loading