From 66bd6206c8eacd237e67307ec7fb648672ea5d9a Mon Sep 17 00:00:00 2001 From: Nabi Sobhi Date: Sun, 30 Aug 2026 20:34:23 +0200 Subject: [PATCH 1/5] Allow specific exception types to bypass the dead letter queue Adds a configurable bypass predicate so fatal/systemic exceptions propagate and crash the consumer instead of being retried or dead-lettered. This prevents a systemic outage (e.g. database down) from silently draining an entire topic into the dead letter queue. New fluent API on DeadLetterQueueOptions: .BypassFor() .BypassWhen(ex => ...) Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- src/Dafda.Tests/Builders/ConsumerBuilder.cs | 11 ++++- src/Dafda.Tests/Consuming/TestConsumer.cs | 48 ++++++++++++++++++- .../Configuration/ConsumerConfiguration.cs | 4 +- .../ConsumerConfigurationBuilder.cs | 4 +- .../ConsumerServiceCollectionExtensions.cs | 6 ++- .../Configuration/DeadLetterQueueOptions.cs | 46 ++++++++++++++++++ src/Dafda/Consuming/Consumer.cs | 10 +++- 7 files changed, 120 insertions(+), 9 deletions(-) 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/Consuming/TestConsumer.cs b/src/Dafda.Tests/Consuming/TestConsumer.cs index ad9df7c..cf9f765 100644 --- a/src/Dafda.Tests/Consuming/TestConsumer.cs +++ b/src/Dafda.Tests/Consuming/TestConsumer.cs @@ -469,11 +469,54 @@ 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 sut = BuildConsumerWithHandler( + handler, + 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); + } + + [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 +538,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,44 @@ 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 (crashing the consumer) instead + /// of being retried or forwarded to the dead letter queue. Returns null + /// when no bypass has been configured. + /// + internal Func BypassPredicate => + _bypassPredicates.Count == 0 + ? null + : exception => _bypassPredicates.Any(predicate => predicate(exception)); + + /// + /// Bypass the dead letter queue for the specified exception type (and any + /// derived types). When a message handler throws a matching exception, it is + /// rethrown so the consumer crashes instead of dead-lettering the message. + /// + /// 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 + /// rethrown so the consumer crashes instead of dead-lettering the message. + /// + /// 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(); From bf403d49a5ae841e1b6668faea1fbb914783af1e Mon Sep 17 00:00:00 2001 From: Nabi Sobhi Date: Sun, 30 Aug 2026 21:12:16 +0200 Subject: [PATCH 2/5] Snapshot the dead letter queue bypass predicates when composing them BypassPredicate returned a delegate closing over the mutable backing list, so predicates registered after the configuration was built would change the behavior of an already running consumer. Enumerating the live list from the consumer thread could also throw a collection-modified exception from inside the error handling path. The predicates are now copied when the combined delegate is composed, matching the other options which are read by value at build time. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../TestDeadLetterQueueOptions.cs | 69 +++++++++++++++++++ .../Configuration/DeadLetterQueueOptions.cs | 17 +++-- 2 files changed, 82 insertions(+), 4 deletions(-) create mode 100644 src/Dafda.Tests/Configuration/TestDeadLetterQueueOptions.cs 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/Configuration/DeadLetterQueueOptions.cs b/src/Dafda/Configuration/DeadLetterQueueOptions.cs index 543bdf2..b2fb498 100644 --- a/src/Dafda/Configuration/DeadLetterQueueOptions.cs +++ b/src/Dafda/Configuration/DeadLetterQueueOptions.cs @@ -51,10 +51,19 @@ public DeadLetterQueueOptions WithMaxRetries(int maxRetries) /// of being retried or forwarded to the dead letter queue. Returns null /// when no bypass has been configured. /// - internal Func BypassPredicate => - _bypassPredicates.Count == 0 - ? null - : exception => _bypassPredicates.Any(predicate => predicate(exception)); + 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 From 3814b8e39a6200b46f24ab4e5262f1535cd02f4a Mon Sep 17 00:00:00 2001 From: Nabi Sobhi Date: Sun, 30 Aug 2026 21:12:48 +0200 Subject: [PATCH 3/5] Assert the offset is not committed when an exception bypasses the queue Skipping the commit is what makes a bypassed message redelivered after the consumer restarts, so the bypass test now uses an onCommit spy to assert it, rather than only checking that the message was neither retried nor sent to the dead letter queue. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- src/Dafda.Tests/Consuming/TestConsumer.cs | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/src/Dafda.Tests/Consuming/TestConsumer.cs b/src/Dafda.Tests/Consuming/TestConsumer.cs index cf9f765..cdfb21a 100644 --- a/src/Dafda.Tests/Consuming/TestConsumer.cs +++ b/src/Dafda.Tests/Consuming/TestConsumer.cs @@ -480,9 +480,15 @@ public async Task propagates_exception_and_bypasses_dead_letter_queue_for_bypass }); 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); @@ -492,6 +498,7 @@ await Assert.ThrowsAsync( Assert.Equal(1, handlerInvocations); Assert.Equal(0, deadLetterQueueSpy.SendCount); + Assert.False(committed); } [Fact] From 3faf618ff8e3681ed983b0acb7ebe8c784d6f93a Mon Sep 17 00:00:00 2001 From: Nabi Sobhi Date: Sun, 30 Aug 2026 21:13:45 +0200 Subject: [PATCH 4/5] Describe what actually happens when an exception bypasses the queue The documentation claimed a bypassed exception crashes the consumer. It is in fact rethrown out of message dispatch without committing the offset, and then routed through the configured consumer error handler: the default strategy stops the application, but RestartConsumer restarts the consumer and the redelivered message fails again. Document the propagation, the skipped commit and the restart caveat instead. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../Configuration/DeadLetterQueueOptions.cs | 31 +++++++++++++++---- 1 file changed, 25 insertions(+), 6 deletions(-) diff --git a/src/Dafda/Configuration/DeadLetterQueueOptions.cs b/src/Dafda/Configuration/DeadLetterQueueOptions.cs index b2fb498..5562648 100644 --- a/src/Dafda/Configuration/DeadLetterQueueOptions.cs +++ b/src/Dafda/Configuration/DeadLetterQueueOptions.cs @@ -47,9 +47,9 @@ public DeadLetterQueueOptions WithMaxRetries(int maxRetries) /// /// A predicate matching exceptions that should bypass the dead letter queue. - /// When an exception matches, it is rethrown (crashing the consumer) instead - /// of being retried or forwarded to the dead letter queue. Returns null - /// when no bypass has been configured. + /// 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 { @@ -67,9 +67,18 @@ internal Func BypassPredicate /// /// Bypass the dead letter queue for the specified exception type (and any - /// derived types). When a message handler throws a matching exception, it is - /// rethrown so the consumer crashes instead of dead-lettering the message. + /// derived types). A matching exception is neither retried nor forwarded to the + /// dead letter queue: it propagates out of message dispatch without the offset + /// being committed, so the message is redelivered once consumption resumes. /// + /// + /// 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. + /// /// The exception type to bypass the dead letter queue for. public DeadLetterQueueOptions BypassFor() where TException : Exception { @@ -80,8 +89,18 @@ public DeadLetterQueueOptions BypassFor() where TException : Excepti /// /// Bypass the dead letter queue for exceptions matching the supplied /// . When it returns true, the exception is - /// rethrown so the consumer crashes instead of dead-lettering the message. + /// neither retried nor forwarded to the dead letter queue: it propagates out of + /// message dispatch without the offset being committed, so the message is + /// redelivered once consumption resumes. /// + /// + /// 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. + /// /// Evaluates a thrown exception and returns true to bypass the dead letter queue. public DeadLetterQueueOptions BypassWhen(Func predicate) { From 4331c105c1faee259806c5f3b903c68bd0fc652f Mon Sep 17 00:00:00 2001 From: Nabi Sobhi Date: Sun, 30 Aug 2026 21:22:11 +0200 Subject: [PATCH 5/5] Document that bypass redelivery requires manual commits The bypass documentation promised the message would be redelivered because the offset had not been committed. That only holds when enable.auto.commit is false: Dafda skips its explicit commit, but the Kafka client stores offsets as messages are consumed and commits them on its interval and on close, so a bypassed message can still be marked as consumed. This is a property of the existing auto-commit mode rather than something the bypass introduces, and it applies equally to an exception escaping a consumer with no dead letter queue configured. Only the redelivery claim was wrong, so narrow it to the manual commit case instead of changing the commit behaviour. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../Configuration/DeadLetterQueueOptions.cs | 23 +++++++++++++++---- 1 file changed, 19 insertions(+), 4 deletions(-) diff --git a/src/Dafda/Configuration/DeadLetterQueueOptions.cs b/src/Dafda/Configuration/DeadLetterQueueOptions.cs index 5562648..6700264 100644 --- a/src/Dafda/Configuration/DeadLetterQueueOptions.cs +++ b/src/Dafda/Configuration/DeadLetterQueueOptions.cs @@ -68,8 +68,8 @@ internal Func BypassPredicate /// /// 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 without the offset - /// being committed, so the message is redelivered once consumption resumes. + /// 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 @@ -78,6 +78,14 @@ internal Func BypassPredicate /// 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 @@ -90,8 +98,7 @@ public DeadLetterQueueOptions BypassFor() where TException : Excepti /// 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 without the offset being committed, so the message is - /// redelivered once consumption resumes. + /// message dispatch, and Dafda does not commit the offset for the message. /// /// /// The exception is then passed to the configured consumer error handler (see @@ -100,6 +107,14 @@ public DeadLetterQueueOptions BypassFor() where TException : Excepti /// 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)