Skip to content

Commit 30fa4ed

Browse files
saby1101claude
andcommitted
feat: add DLQ support for RabbitMQ messages exceeding max retries
- Add DLQ queue creation in rabbitmq-utils.ts for messages that fail after 10 retries - Add dlqPublisher in consumer to move failed messages to DLQ with headers: - x-original-queue: source queue name - x-final-status-code or x-final-error: last failure reason - x-final-retry-count: number of retry attempts - Fix critical bug in getRetryCount(): was summing x-death counts from ALL queues (both main + retry queue), effectively counting each retry twice. Now only counts entries with reason="rejected" (actual consumer rejections) - Add comprehensive e2e test verifying full 10-retry cycle and DLQ behavior - Add RABBITMQ.md architecture documentation with mermaid flowchart 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude <noreply@anthropic.com>
1 parent 418ddb0 commit 30fa4ed

7 files changed

Lines changed: 1231 additions & 15 deletions

File tree

src/event-bus/RABBITMQ.md

Lines changed: 132 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,132 @@
1+
# RabbitMQ Architecture
2+
3+
This document describes the exchange and queue topology used by the event-bus RabbitMQ implementation.
4+
5+
## Naming Convention
6+
7+
All resources use a dynamic prefix derived from the service name (part before the first hyphen):
8+
- `wms-cincout` → prefix: `wms`
9+
- `xyz-service` → prefix: `xyz`
10+
- `noprefix` → prefix: `default`
11+
12+
## Topology Flowchart
13+
14+
```mermaid
15+
flowchart TD
16+
subgraph Publisher
17+
P[Producer]
18+
end
19+
20+
subgraph "Main Exchange (fanout)"
21+
ME["{prefix}.main-exchange"]
22+
end
23+
24+
subgraph "Service Queues"
25+
SQ1["{prefix}.queue.{service-a}"]
26+
SQ2["{prefix}.queue.{service-b}"]
27+
end
28+
29+
subgraph "Retry System (per service)"
30+
RE["{prefix}.retry-exchange.{service}<br/>(direct)"]
31+
RQ["{prefix}.retry-queue.{service}<br/>(TTL: 5s)"]
32+
end
33+
34+
subgraph "Dead Letter Queue (per service)"
35+
DLQ["{prefix}.dlq.{service}<br/>(manual retry)"]
36+
end
37+
38+
subgraph Consumers
39+
C1[Consumer A]
40+
C2[Consumer B]
41+
end
42+
43+
P -->|publish| ME
44+
ME -->|fanout| SQ1
45+
ME -->|fanout| SQ2
46+
SQ1 -->|consume| C1
47+
SQ2 -->|consume| C2
48+
49+
C1 -->|nack/reject| SQ1
50+
SQ1 -.->|dead-letter| RE
51+
RE -->|route| RQ
52+
RQ -.->|"dead-letter after TTL<br/>(via default exchange)"| SQ1
53+
C1 -->|"retry >= 10"| DLQ
54+
```
55+
56+
## Message Flow
57+
58+
### Happy Path
59+
1. Producer publishes message to `{prefix}.main-exchange`
60+
2. Exchange fans out message to all bound service queues
61+
3. Consumer reads message from `{prefix}.queue.{service}`
62+
4. Consumer acknowledges (ack) → message removed
63+
64+
### Retry Path (dead-letter)
65+
Used for: 5xx errors, 429 (rate-limit), 409 (lock conflict)
66+
67+
1. Consumer returns `DROP` (nack with requeue=false)
68+
2. Message dead-letters to `{prefix}.retry-exchange.{service}`
69+
3. Retry exchange routes to `{prefix}.retry-queue.{service}`
70+
4. Message sits in retry queue for 5 seconds (TTL)
71+
5. After TTL expires, message dead-letters directly to `{prefix}.queue.{service}` (via default exchange)
72+
6. Message is re-delivered only to the failed service (not fanned out to all services)
73+
7. **Max 10 retries** - after 10 failed attempts, message is moved to `{prefix}.dlq.{service}` for manual retry (logged as `RABBITMQ_MESSAGE_MAX_RETRIES_EXCEEDED`)
74+
75+
### Delayed Message Path (local sleep)
76+
Used for: 425 (too early - `processAfterDelayMs` not yet reached)
77+
78+
1. Consumer sleeps locally (randomDelay)
79+
2. Consumer returns `REQUEUE` (nack with requeue=true)
80+
3. Message returns to the same queue immediately for retry
81+
4. This avoids multiple DLX cycles when delay exceeds 5s TTL
82+
83+
## Consumer Status Handling
84+
85+
| HTTP Status | ConsumerStatus | Behavior |
86+
|-------------|----------------|----------|
87+
| 2xx | `ACK` | Success, message removed |
88+
| 429, 409 | `DROP` | Dead-letter retry (rate-limit/lock conflict) |
89+
| 425 | `REQUEUE` | Local sleep + immediate requeue (delayed message) |
90+
| 5xx | `DROP` | Dead-letter retry (transient error) |
91+
| Other 4xx | `DROP` | Dead-letter (bad message, will likely fail again) |
92+
| Exception | `DROP` | Dead-letter (consumer error) |
93+
| Any (retry >= 10) | `ACK` | Max retries exceeded, message moved to DLQ |
94+
95+
## Queue Configuration
96+
97+
| Queue | Type | Dead-Letter Exchange | Dead-Letter Routing Key | TTL |
98+
|-------|------|---------------------|------------------------|-----|
99+
| `{prefix}.queue.{service}` | classic | `{prefix}.retry-exchange.{service}` | `retry` | - |
100+
| `{prefix}.retry-queue.{service}` | classic | `""` (default) | `{prefix}.queue.{service}` | 5000ms |
101+
| `{prefix}.dlq.{service}` | classic | - | - | - |
102+
103+
## Exchange Configuration
104+
105+
| Exchange | Type | Purpose |
106+
|----------|------|---------|
107+
| `{prefix}.main-exchange` | fanout | Distribute messages to all service queues |
108+
| `{prefix}.retry-exchange.{service}` | direct | Route failed messages to retry queue (binding key: `retry`) |
109+
110+
## Dead Letter Queue (DLQ)
111+
112+
Messages that fail after 10 retry attempts are moved to `{prefix}.dlq.{service}` for manual inspection and retry.
113+
114+
### DLQ Message Headers
115+
116+
Messages in the DLQ include additional headers for debugging:
117+
118+
| Header | Description |
119+
|--------|-------------|
120+
| `x-original-queue` | The queue the message was consumed from |
121+
| `x-final-status-code` | HTTP status code from last processing attempt (if applicable) |
122+
| `x-final-error` | Error message from last processing attempt (if exception) |
123+
| `x-final-retry-count` | Number of retry attempts before moving to DLQ |
124+
| `x-death` | Standard RabbitMQ dead-letter history |
125+
126+
### Manual Retry
127+
128+
To retry messages from the DLQ:
129+
1. Inspect messages in `{prefix}.dlq.{service}` queue
130+
2. Fix the underlying issue (e.g., dependent service, data issue)
131+
3. Move message back to `{prefix}.queue.{service}` for reprocessing
132+
4. The `x-death` header will be preserved, so retry count continues from where it left off

src/event-bus/commons.ts

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -174,6 +174,13 @@ export function CreateHandlerRunner(
174174
const handlers = handlersMap.get(eventMsg.event) ?? [];
175175
const actions: Action[] = [];
176176
const specifiedFile = eventMsg.attributes.file;
177+
req.log.debug({
178+
tag: "HANDLER_LOOKUP",
179+
event: eventMsg.event,
180+
specifiedFile,
181+
registeredHandlers: handlers.map((h) => h.file),
182+
handlersCount: handlers.length,
183+
});
177184
for (const { file, handler } of handlers) {
178185
if (specifiedFile && file !== specifiedFile) {
179186
continue;
@@ -190,8 +197,18 @@ export function CreateHandlerRunner(
190197
actions.push(act);
191198
}
192199

200+
req.log.debug({
201+
tag: "HANDLER_ACTIONS_CREATED",
202+
event: eventMsg.event,
203+
actionsCount: actions.length,
204+
});
193205
for (let i = 0; i < actions.length; i += CONCURRENCY) {
194206
await Promise.all(actions.slice(i, i + CONCURRENCY).map((act) => act()));
195207
}
208+
req.log.debug({
209+
tag: "HANDLER_ACTIONS_COMPLETED",
210+
event: eventMsg.event,
211+
actionsCount: actions.length,
212+
});
196213
};
197214
}

src/event-bus/event-consumer/rabbitmq.ts

Lines changed: 99 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { Connection, ConsumerStatus } from "rabbitmq-client";
1+
import { AsyncMessage, Connection, ConsumerStatus } from "rabbitmq-client";
22
import {
33
RABBITMQ_TAG,
44
ensureRabbitMqExchangesAndQueues,
@@ -7,6 +7,27 @@ import {
77
import { EventConsumerBuilder } from "./interface";
88
import { randomDelay } from "./utils";
99

10+
const MAX_RETRY_COUNT = 10;
11+
12+
/**
13+
* Extract retry count from x-death header added by RabbitMQ DLX.
14+
* Only counts rejections (consumer drops), not TTL expirations from retry queue.
15+
* x-death entries look like: { queue, reason, count, ... }
16+
* - reason "rejected" = consumer nacked/dropped the message
17+
* - reason "expired" = TTL expired in retry queue (not a retry attempt)
18+
*/
19+
function getRetryCount(msg: AsyncMessage): number {
20+
const xDeath = msg.headers?.["x-death"];
21+
if (!Array.isArray(xDeath) || xDeath.length === 0) {
22+
return 0;
23+
}
24+
// Only count "rejected" entries (actual consumer rejections)
25+
// Don't count "expired" entries from retry queue TTL
26+
return xDeath
27+
.filter((death) => death.reason === "rejected")
28+
.reduce((total, death) => total + (death.count || 0), 0);
29+
}
30+
1031
/**
1132
* RabbitMq supports
1233
* 1. prefetchCount -> this is to support parallel processing
@@ -16,6 +37,14 @@ import { randomDelay } from "./utils";
1637
export const RabbitMqServiceBusConsumerBuilder: EventConsumerBuilder = async (
1738
instance,
1839
) => {
40+
// Skip consumer creation if no handlers are registered
41+
if ((instance as any)._hasEventHandlers === false) {
42+
instance.log.info("No event handlers registered, skipping RabbitMQ consumer");
43+
return {
44+
close: async () => {},
45+
};
46+
}
47+
1948
if (!process.env.RABBITMQ_URL) {
2049
throw new Error("RabbitMq requires RABBITMQ_URL");
2150
}
@@ -29,6 +58,12 @@ export const RabbitMqServiceBusConsumerBuilder: EventConsumerBuilder = async (
2958
const service = process.env.K_SERVICE;
3059
const prefix = getServicePrefix(service);
3160

61+
// Publisher for DLQ messages
62+
const dlqPublisher = connection.createPublisher({
63+
confirm: true,
64+
maxAttempts: 3,
65+
});
66+
3267
const sub = connection.createConsumer(
3368
{
3469
queue: `${prefix}.queue.${service}`,
@@ -41,7 +76,7 @@ export const RabbitMqServiceBusConsumerBuilder: EventConsumerBuilder = async (
4176
concurrency: 10,
4277
consumerTag: `${prefix}.consumer.${service}.${RABBITMQ_TAG}`,
4378
},
44-
async (msg) => {
79+
async (msg: AsyncMessage) => {
4580
if (ctrl.signal.aborted) {
4681
return ConsumerStatus.REQUEUE;
4782
}
@@ -60,18 +95,46 @@ export const RabbitMqServiceBusConsumerBuilder: EventConsumerBuilder = async (
6095
});
6196
if (resp.statusCode >= 200 && resp.statusCode < 300) {
6297
return ConsumerStatus.ACK;
63-
} else if (resp.statusCode === 429 || resp.statusCode === 409) {
64-
await randomDelay();
65-
// rate-limited or lock-conflict
66-
return ConsumerStatus.REQUEUE;
98+
}
99+
100+
// Check retry count before dead-lettering
101+
const retryCount = getRetryCount(msg);
102+
if (retryCount >= MAX_RETRY_COUNT) {
103+
instance.log.error({
104+
tag: "RABBITMQ_MESSAGE_MAX_RETRIES_EXCEEDED",
105+
messageId: msg.messageId,
106+
retryCount,
107+
statusCode: resp.statusCode,
108+
body: resp.body,
109+
});
110+
// Store in DLQ for manual retry before ACKing
111+
await dlqPublisher.send(
112+
{
113+
routingKey: `${prefix}.dlq.${service}`,
114+
contentType: "application/json",
115+
headers: {
116+
...msg.headers,
117+
"x-original-queue": `${prefix}.queue.${service}`,
118+
"x-final-status-code": resp.statusCode,
119+
"x-final-retry-count": retryCount,
120+
},
121+
},
122+
msg.body,
123+
);
124+
// ACK to remove from queue permanently (don't dead-letter again)
125+
return ConsumerStatus.ACK;
126+
}
127+
128+
if (resp.statusCode === 429 || resp.statusCode === 409) {
129+
// rate-limited or lock-conflict, use dead-letter retry
130+
return ConsumerStatus.DROP;
67131
} else if (resp.statusCode === 425) {
68-
// delayed message. requeue with delay to avoid tight loop
132+
// delayed message, use local sleep since delay may exceed DLX TTL
69133
await randomDelay();
70134
return ConsumerStatus.REQUEUE;
71135
} else if (resp.statusCode >= 500 && resp.statusCode < 600) {
72-
// transient server error, retry
73-
await randomDelay();
74-
return ConsumerStatus.REQUEUE;
136+
// transient server error, use dead-letter retry
137+
return ConsumerStatus.DROP;
75138
} else {
76139
instance.log.warn({
77140
tag: "RABBITMQ_MESSAGE_DROPPED",
@@ -82,6 +145,31 @@ export const RabbitMqServiceBusConsumerBuilder: EventConsumerBuilder = async (
82145
return ConsumerStatus.DROP;
83146
}
84147
} catch (err) {
148+
// Check retry count before dead-lettering on exception
149+
const retryCount = getRetryCount(msg);
150+
if (retryCount >= MAX_RETRY_COUNT) {
151+
instance.log.error({
152+
tag: "RABBITMQ_MESSAGE_MAX_RETRIES_EXCEEDED",
153+
messageId: msg.messageId,
154+
retryCount,
155+
err,
156+
});
157+
// Store in DLQ for manual retry before ACKing
158+
await dlqPublisher.send(
159+
{
160+
routingKey: `${prefix}.dlq.${service}`,
161+
contentType: "application/json",
162+
headers: {
163+
...msg.headers,
164+
"x-original-queue": `${prefix}.queue.${service}`,
165+
"x-final-error": err instanceof Error ? err.message : String(err),
166+
"x-final-retry-count": retryCount,
167+
},
168+
},
169+
msg.body,
170+
);
171+
return ConsumerStatus.ACK;
172+
}
85173
instance.log.error({
86174
tag: "RABBITMQ_CONSUMER_ERROR",
87175
err: err,
@@ -97,6 +185,7 @@ export const RabbitMqServiceBusConsumerBuilder: EventConsumerBuilder = async (
97185
close: async () => {
98186
ctrl.abort();
99187
await sub.close();
188+
await dlqPublisher.close();
100189
await connection.close();
101190
},
102191
};

src/event-bus/index.ts

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ const plugin: FastifyPluginAsync<EventBusOptions> = async function (
99
switch (options.busType) {
1010
case "gcp-pubsub":
1111
// eslint-disable-next-line @typescript-eslint/no-require-imports
12-
f.register(require("./gcp-pubsub"), options);
12+
await f.register(require("./gcp-pubsub"), options);
1313
break;
1414
case "azure-servicebus":
1515
if (!options.namespace) {
@@ -18,15 +18,15 @@ const plugin: FastifyPluginAsync<EventBusOptions> = async function (
1818
);
1919
}
2020
// eslint-disable-next-line @typescript-eslint/no-require-imports
21-
f.register(require("./azure-servicebus"), options);
21+
await f.register(require("./azure-servicebus"), options);
2222
break;
2323
case "rabbitmq":
2424
// eslint-disable-next-line @typescript-eslint/no-require-imports
25-
f.register(require("./rabbitmq"), options);
25+
await f.register(require("./rabbitmq"), options);
2626
break;
2727
default:
2828
// eslint-disable-next-line @typescript-eslint/no-require-imports
29-
f.register(require("./local"), options);
29+
await f.register(require("./local"), options);
3030
}
3131

3232
///

src/event-bus/rabbitmq-utils.ts

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@ export async function ensureRabbitMqExchangesAndQueues(
2929
queue: `${prefix}.queue.${service}`,
3030
arguments: {
3131
"x-dead-letter-exchange": `${prefix}.retry-exchange.${service}`,
32+
"x-dead-letter-routing-key": "retry", // Fixed routing key for DLX
3233
"x-queue-type": "classic",
3334
},
3435
autoDelete: false,
@@ -53,7 +54,8 @@ export async function ensureRabbitMqExchangesAndQueues(
5354
queue: `${prefix}.retry-queue.${service}`,
5455
arguments: {
5556
"x-message-ttl": 5000,
56-
"x-dead-letter-exchange": `${prefix}.main-exchange`,
57+
"x-dead-letter-exchange": "", // default exchange routes by queue name
58+
"x-dead-letter-routing-key": `${prefix}.queue.${service}`,
5759
"x-queue-type": "classic",
5860
},
5961
autoDelete: false,
@@ -64,5 +66,17 @@ export async function ensureRabbitMqExchangesAndQueues(
6466
await connection.queueBind({
6567
exchange: `${prefix}.retry-exchange.${service}`,
6668
queue: `${prefix}.retry-queue.${service}`,
69+
routingKey: "retry", // Must match x-dead-letter-routing-key from main queue
70+
});
71+
// Dead Letter Queue for messages that exceed max retries (manual retry)
72+
await connection.queueDeclare({
73+
queue: `${prefix}.dlq.${service}`,
74+
arguments: {
75+
"x-queue-type": "classic",
76+
},
77+
autoDelete: false,
78+
durable: true,
79+
exclusive: false,
80+
passive: false,
6781
});
6882
}

0 commit comments

Comments
 (0)