using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using RabbitMQ.Client;
using System;
using System.Collections.Generic;
using System.Text;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
namespace FrameworkBLL.SchedulerTaskGeneratorPublisher
{
public interface IRabbitMqPublisher
{
Task PublishAsync(string queue, object payload);
Task PurgeQueueAsync(string queue);
}
///
/// RabbitMQ Publisher — initializes connection asynchronously via IHostedService.StartAsync.
///
/// Queue topology:
/// Scheduler.Ready — main work queue (jobs due for execution)
/// Scheduler.Retry — retry queue (jobs that failed, RetryCount < MaxRetry)
/// Scheduler.DLQ — dead letter queue (jobs that exhausted retries)
///
/// Retry/DLQ routing is done by the CONSUMER — this publisher just sends to named queues.
///
/// Features:
/// Async connection (no sync-over-async)
/// Queues declared once at startup (idempotent)
/// Publisher confirms on every publish (no silent message loss)
/// Durable queues + persistent messages
/// DLX (Dead Letter Exchange) wired on Ready and Retry queues
///
public sealed class RabbitMqPublisher : IHostedService, IRabbitMqPublisher, IAsyncDisposable
{
// ─────────────────────────────────────────────────────────────────
// Queue / Exchange names
// ─────────────────────────────────────────────────────────────────
private const string ReadyQueue = "Scheduler.Ready";
private const string RetryQueue = "Scheduler.Retry";
private const string DlqQueue = "Scheduler.DLQ";
private const string DlxExchange = "Scheduler.DLX"; // Dead Letter Exchange
private static readonly string[] AllQueues =
[ReadyQueue, RetryQueue, DlqQueue];
private readonly IConfiguration _config;
private readonly ILogger _logger;
private IConnection? _connection;
public RabbitMqPublisher(IConfiguration config, ILogger logger)
{
_config = config ?? throw new ArgumentNullException(nameof(config));
_logger = logger ?? throw new ArgumentNullException(nameof(logger));
}
// ─────────────────────────────────────────────────────────────────
// IHostedService — startup / shutdown
// ─────────────────────────────────────────────────────────────────
///
/// Called by the host at startup.
/// Opens the connection and declares all queues + DLX topology once.
///
public Task StartAsync(CancellationToken cancellationToken)
{
// Fire-and-forget retry loop — app startup is not blocked
_ = ConnectWithRetryAsync(cancellationToken);
return Task.CompletedTask;
}
private async Task ConnectWithRetryAsync(CancellationToken ct)
{
var host = _config["RabbitMQ:Host"] ?? "localhost";
var user = _config["RabbitMQ:UserName"] ?? "guest";
var pass = _config["RabbitMQ:Password"] ?? "guest";
var vhost = _config["RabbitMQ:VirtualHost"] ?? "/";
var delay = TimeSpan.FromSeconds(5);
while (!ct.IsCancellationRequested)
{
try
{
await ConnectAsync(host, user, pass, vhost, ct).ConfigureAwait(false);
return; // success
}
catch (OperationCanceledException)
{
return;
}
catch (Exception ex)
{
_connection = null;
_logger.LogWarning(ex,
"RabbitMqPublisher: connect failed ({Host}) — retrying in {Delay}s",
host, delay.TotalSeconds);
try { await Task.Delay(delay, ct).ConfigureAwait(false); }
catch (OperationCanceledException) { return; }
delay = TimeSpan.FromSeconds(Math.Min(delay.TotalSeconds * 2, 60));
}
}
}
private async Task ConnectAsync(
string host, string user, string pass, string vhost,
CancellationToken cancellationToken)
{
var factory = new ConnectionFactory
{
HostName = host,
UserName = user,
Password = pass,
VirtualHost = vhost,
AutomaticRecoveryEnabled = true
};
_connection = await factory.CreateConnectionAsync(cancellationToken)
.ConfigureAwait(false);
await using var channel = await _connection
.CreateChannelAsync(cancellationToken: cancellationToken)
.ConfigureAwait(false);
// Declare DLX (fanout) — routes expired/rejected messages to DLQ
await channel.ExchangeDeclareAsync(
exchange: DlxExchange,
type: ExchangeType.Fanout,
durable: true,
autoDelete: false,
cancellationToken: cancellationToken
).ConfigureAwait(false);
// Declare DLQ first (must exist before other queues reference it via DLX)
await channel.QueueDeclareAsync(
queue: DlqQueue,
durable: true,
exclusive: false,
autoDelete: false,
arguments: null,
cancellationToken: cancellationToken
).ConfigureAwait(false);
await channel.QueueBindAsync(
queue: DlqQueue,
exchange: DlxExchange,
routingKey: "",
cancellationToken: cancellationToken
).ConfigureAwait(false);
await channel.QueueDeclareAsync(
queue: ReadyQueue,
durable: true,
exclusive: false,
autoDelete: false,
arguments: new Dictionary
{
["x-dead-letter-exchange"] = DlxExchange
},
cancellationToken: cancellationToken
).ConfigureAwait(false);
await channel.QueueDeclareAsync(
queue: RetryQueue,
durable: true,
exclusive: false,
autoDelete: false,
arguments: new Dictionary
{
["x-dead-letter-exchange"] = DlxExchange
},
cancellationToken: cancellationToken
).ConfigureAwait(false);
// Declare action-exec partition queues so messages can be enqueued
// before ActionProcessorWorker connects and starts consuming.
const int partitionCount = 5;
for (int i = 0; i < partitionCount; i++)
{
await channel.QueueDeclareAsync(
queue: $"action-exec-p{i}",
durable: true,
exclusive: false,
autoDelete: false,
arguments: null,
cancellationToken: cancellationToken
).ConfigureAwait(false);
}
_logger.LogInformation(
"RabbitMqPublisher connected. Queues declared: {Queues}, action-exec-p0..p{Max}. DLX={DLX}",
string.Join(", ", AllQueues), partitionCount - 1, DlxExchange);
}
public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask;
// ─────────────────────────────────────────────────────────────────
// IRabbitMqPublisher — publish
// ─────────────────────────────────────────────────────────────────
///
/// Publishes a message to the named queue.
/// Publisher confirms are enabled — method awaits broker acknowledgement.
/// A new channel is created per call (acceptable at 30s poll intervals).
/// For high-throughput scenarios, replace with a channel pool.
///
public async Task PublishAsync(string queue, object payload)
{
if (_connection is null)
throw new InvalidOperationException(
$"RabbitMqPublisher: cannot publish to '{queue}' — RabbitMQ is not connected.");
// Publisher confirms enabled at channel level (RabbitMQ.Client v7 API)
await using var channel = await _connection.CreateChannelAsync(
new CreateChannelOptions(
publisherConfirmationsEnabled: true,
publisherConfirmationTrackingEnabled: true)
).ConfigureAwait(false);
var body = Encoding.UTF8.GetBytes(JsonSerializer.Serialize(payload));
var props = new BasicProperties
{
Persistent = true, // survive broker restart
ContentType = "application/json"
};
// Awaiting BasicPublishAsync with confirms enabled blocks until broker acks
await channel.BasicPublishAsync(
exchange: "",
routingKey: queue,
mandatory: false,
basicProperties: props,
body: body
).ConfigureAwait(false);
_logger.LogDebug(
"RabbitMqPublisher published | Queue={Queue} PayloadType={Type}",
queue, payload.GetType().Name);
}
// ─────────────────────────────────────────────────────────────────
// Debug utility
// ─────────────────────────────────────────────────────────────────
///
/// DEBUG ONLY — clears all messages from a queue.
/// Never call in production.
///
public async Task PurgeQueueAsync(string queue)
{
if (_connection is null)
{
_logger.LogWarning(
"RabbitMqPublisher: skipping purge of {Queue} — RabbitMQ is not connected.",
queue);
return;
}
await using var channel = await _connection.CreateChannelAsync()
.ConfigureAwait(false);
await channel.QueuePurgeAsync(queue).ConfigureAwait(false);
_logger.LogWarning("RabbitMqPublisher queue purged | Queue={Queue}", queue);
}
// ─────────────────────────────────────────────────────────────────
// IAsyncDisposable
// ─────────────────────────────────────────────────────────────────
public async ValueTask DisposeAsync()
{
if (_connection is not null)
await _connection.DisposeAsync().ConfigureAwait(false);
}
}
}