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
Original file line number Diff line number Diff line change
Expand Up @@ -136,12 +136,30 @@ await channel.ExchangeDeclareAsync(
autoDelete: false,
cancellationToken: cancellationToken);

var emailQueueArguments = new Dictionary<string, object?>
{
["x-queue-type"] = "quorum",

["x-delivery-limit"] = _options.DeliveryLimit,

["x-delayed-retry-type"] = "failed",

["x-delayed-retry-min"] = _options.RetryMinDelayMilliseconds,

["x-delayed-retry-max"] = _options.RetryMaxDelayMilliseconds,

["x-dead-letter-exchange"] = _options.DeadLetterExchangeName,

["x-dead-letter-routing-key"] = _options.DeadLetterRoutingKey
};

await channel.QueueDeclareAsync(
queue: _options.EmailQueueName,
durable: true,
exclusive: false,
autoDelete: false,
cancellationToken: cancellationToken);
queue: _options.EmailQueueName,
durable: true,
exclusive: false,
autoDelete: false,
arguments: emailQueueArguments,
cancellationToken: cancellationToken);

await channel.QueueBindAsync(
queue: _options.EmailQueueName,
Expand Down
12 changes: 12 additions & 0 deletions OrderProcessing.Api/Services/Messaging/RabbitMqOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -23,4 +23,16 @@ public sealed class RabbitMqOptions
public string ClientProvidedName { get; set; } = "order-processing-api-publisher";

public int NetworkRecoveryIntervalSeconds { get; set; } = 5;

public string DeadLetterExchangeName { get; set; } = "order-processing.dead-letter";

public string DeadLetterQueueName { get; set; } = "order-processing.email.dead-letter";

public string DeadLetterRoutingKey { get; set; } = "email.failed";

public int DeliveryLimit { get; set; } = 5;

public int RetryMinDelayMilliseconds { get; set; } = 5000;

public int RetryMaxDelayMilliseconds { get; set; } = 30000;
}
8 changes: 7 additions & 1 deletion OrderProcessing.Api/appsettings.json
Original file line number Diff line number Diff line change
Expand Up @@ -44,14 +44,20 @@
"PollingIntervalSeconds": 5
},
"RabbitMq": {
"Enabled": false,
"Enabled": true,
"HostName": "localhost",
"Port": 5672,
"UserName": "guest",
"Password": "guest",
"VirtualHost": "/",
"ExchangeName": "order-processing.events",
"EmailQueueName": "order-processing.email",
"DeadLetterExchangeName": "order-processing.dead-letter",
"DeadLetterQueueName": "order-processing.email.dead-letter",
"DeadLetterRoutingKey": "email.failed",
"DeliveryLimit": 5,
"RetryMinDelayMilliseconds": 5000,
"RetryMaxDelayMilliseconds": 30000,
"ClientProvidedName": "order-processing-api-publisher",
"NetworkRecoveryIntervalSeconds": 5
}
Expand Down
12 changes: 12 additions & 0 deletions OrderProcessing.EmailWorker/Configuration/RabbitMqOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -23,4 +23,16 @@ public sealed class RabbitMqOptions
public int NetworkRecoveryIntervalSeconds { get; set; } = 5;

public ushort PrefetchCount { get; set; } = 1;

public string DeadLetterExchangeName { get; set; } = "order-processing.dead-letter";

public string DeadLetterQueueName { get; set; } = "order-processing.email.dead-letter";

public string DeadLetterRoutingKey { get; set; } = "email.failed";

public int DeliveryLimit { get; set; } = 5;

public int RetryMinDelayMilliseconds { get; set; } = 5000;

public int RetryMaxDelayMilliseconds { get; set; } = 30000;
}
68 changes: 55 additions & 13 deletions OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -189,12 +189,11 @@ await RejectPermanentFailureAsync(
_logger.LogError(
exception,
"Temporary failure processing RabbitMQ " +
"message {MessageId}. The message will be requeued",
"message {MessageId}. The message will be retried",
eventArgs.BasicProperties.MessageId);

await channel.BasicNackAsync(
await channel.BasicRejectAsync(
deliveryTag: eventArgs.DeliveryTag,
multiple: false,
requeue: true);
}
}
Expand All @@ -204,20 +203,62 @@ private async Task RejectPermanentFailureAsync(
BasicDeliverEventArgs eventArgs,
Exception exception)
{
_logger.LogError(
exception,
"Permanent failure processing RabbitMQ " +
"message {MessageId}. The message will not be requeued",
eventArgs.BasicProperties.MessageId);

await channel.BasicNackAsync(
deliveryTag: eventArgs.DeliveryTag,
multiple: false,
requeue: false);
_logger.LogError(exception,
"Permanent failure processing RabbitMQ " +
"message {MessageId}. " +
"The message will be dead-lettered",
eventArgs.BasicProperties.MessageId);

await channel.BasicRejectAsync(
deliveryTag: eventArgs.DeliveryTag,
requeue: false);
}

private async Task DeclareTopologyAsync(IChannel channel, CancellationToken cancellationToken)
{
var emailQueueArguments = new Dictionary<string, object?>
{
["x-queue-type"] = "quorum",

["x-delivery-limit"] = _options.DeliveryLimit,

["x-delayed-retry-type"] = "failed",

["x-delayed-retry-min"] = _options.RetryMinDelayMilliseconds,

["x-delayed-retry-max"] = _options.RetryMaxDelayMilliseconds,

["x-dead-letter-exchange"] = _options.DeadLetterExchangeName,

["x-dead-letter-routing-key"] = _options.DeadLetterRoutingKey
};

await channel.ExchangeDeclareAsync(
exchange: _options.DeadLetterExchangeName,
type: ExchangeType.Direct,
durable: true,
autoDelete: false,
cancellationToken: cancellationToken);

var deadLetterQueueArguments = new Dictionary<string, object?>
{
["x-queue-type"] = "quorum"
};

await channel.QueueDeclareAsync(
queue: _options.DeadLetterQueueName,
durable: true,
exclusive: false,
autoDelete: false,
arguments: deadLetterQueueArguments,
cancellationToken: cancellationToken);

await channel.QueueBindAsync(
queue: _options.DeadLetterQueueName,
exchange: _options.DeadLetterExchangeName,
routingKey: _options.DeadLetterRoutingKey,
cancellationToken: cancellationToken);

await channel.ExchangeDeclareAsync(
exchange: _options.ExchangeName,
type: ExchangeType.Topic,
Expand All @@ -230,6 +271,7 @@ await channel.QueueDeclareAsync(
durable: true,
exclusive: false,
autoDelete: false,
arguments: emailQueueArguments,
cancellationToken: cancellationToken);

await channel.QueueBindAsync(
Expand Down
6 changes: 6 additions & 0 deletions OrderProcessing.EmailWorker/appsettings.json
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,12 @@
"VirtualHost": "/",
"ExchangeName": "order-processing.events",
"EmailQueueName": "order-processing.email",
"DeadLetterExchangeName": "order-processing.dead-letter",
"DeadLetterQueueName": "order-processing.email.dead-letter",
"DeadLetterRoutingKey": "email.failed",
"DeliveryLimit": 5,
"RetryMinDelayMilliseconds": 5000,
"RetryMaxDelayMilliseconds": 30000,
"ClientProvidedName": "order-processing-email-worker",
"NetworkRecoveryIntervalSeconds": 5,
"PrefetchCount": 1
Expand Down