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); } } }