From 795c1d820dc698d76977481a722617c9aa7d7647 Mon Sep 17 00:00:00 2001 From: Akanksha Trehun Date: Tue, 18 Aug 2026 15:53:37 +0530 Subject: [PATCH] fix(connector): clamp exponential backoff shift to avoid uint32 wraparound calculateIncrementalDelay does uint32(1 << redeliveryCount) for exponential backoff. Once redeliveryCount hits 32 that shift wraps a uint32 to 0, so a message that's been redelivered that many times gets an instant redelivery instead of the configured max delay - worse than no backoff at all. Clamp the shift amount before applying it. Signed-off-by: Akanksha Trehun --- pulsar/connector/backoff_test.go | 9 +++++++++ pulsar/connector/dlq.go | 9 ++++++++- 2 files changed, 17 insertions(+), 1 deletion(-) diff --git a/pulsar/connector/backoff_test.go b/pulsar/connector/backoff_test.go index 7bc5a19..77a2a95 100644 --- a/pulsar/connector/backoff_test.go +++ b/pulsar/connector/backoff_test.go @@ -225,6 +225,15 @@ func TestNackBackoffPolicyExponential(t *testing.T) { expectedDelay: 256 * time.Second, // 2^10 = 1024, but capped at 256 expectedValidationErr: false, }, + { + name: "Exponential: redelivery count 32 should still hit the cap, not reset to zero", + minMultiplier: 1, + maxMultiplier: 256, + baseDelay: 1 * time.Second, + redeliveryCount: 32, + expectedDelay: 256 * time.Second, + expectedValidationErr: false, + }, } for _, tt := range tests { diff --git a/pulsar/connector/dlq.go b/pulsar/connector/dlq.go index 913681d..2f23e98 100644 --- a/pulsar/connector/dlq.go +++ b/pulsar/connector/dlq.go @@ -89,7 +89,14 @@ func (nbp *NackBackoffPolicy) calculateIncrementalDelay(redeliveryCount uint32) if useExponential { // Use exponential backoff: base * 2^redeliveryCount - exponentialMultiplier := uint32(1 << redeliveryCount) // 2^redeliveryCount + // 1 << n on a uint32 wraps to 0 once n >= 32, so a message stuck in + // redelivery for that long would otherwise get an instant redelivery + // instead of the (capped) maximum delay. Clamp the shift instead. + shift := redeliveryCount + if shift >= 32 { + shift = 31 + } + exponentialMultiplier := uint32(1) << shift if nbp.maxRedeliveryDelayMultiplier > 0 && exponentialMultiplier > nbp.maxRedeliveryDelayMultiplier { exponentialMultiplier = nbp.maxRedeliveryDelayMultiplier }