using System.Text.Json; using Jiaowu.Api.Domain.System; using RabbitMQ.Client; using RabbitMQ.Client.Events; namespace Jiaowu.Api.Infrastructure.BackgroundJobs; public sealed class RabbitMqBackgroundJobTransport( BackgroundJobOptions jobOptions, RabbitMqOptions rabbitOptions, ILogger logger) : IBackgroundJobTransport, IAsyncDisposable { private readonly SemaphoreSlim _gate = new(1, 1); private IConnection? _connection; private IChannel? _channel; public bool IsDurable => true; public async ValueTask PublishAsync( BackgroundJobEnvelope message, CancellationToken cancellationToken) { var body = JsonSerializer.SerializeToUtf8Bytes(message); await _gate.WaitAsync(cancellationToken); try { var channel = await GetChannelAsync(cancellationToken); var properties = new BasicProperties { Persistent = true, ContentType = "application/json", MessageId = message.OutboxMessageId.ToString("D"), Type = message.JobKind.ToString(), AppId = "jiaowu-api" }; await channel.BasicPublishAsync( jobOptions.Exchange, RabbitMqBackgroundJobTopology.RoutingKey(message.JobKind), mandatory: true, properties, body, cancellationToken); } catch { await ResetConnectionAsync(); throw; } finally { _gate.Release(); } } public async Task CheckHealthAsync(CancellationToken cancellationToken) { await _gate.WaitAsync(cancellationToken); try { var channel = await GetChannelAsync(cancellationToken); return channel.IsOpen; } catch (Exception exception) { logger.LogWarning(exception, "RabbitMQ health check failed."); await ResetConnectionAsync(); return false; } finally { _gate.Release(); } } private async Task GetChannelAsync(CancellationToken cancellationToken) { if (_channel is { IsOpen: true }) return _channel; await ResetConnectionAsync(); _connection = await RabbitMqBackgroundJobTopology.CreateConnectionAsync( rabbitOptions, "jiaowu-background-job-publisher", consumerDispatchConcurrency: 1, cancellationToken); _channel = await _connection.CreateChannelAsync( new CreateChannelOptions( publisherConfirmationsEnabled: true, publisherConfirmationTrackingEnabled: true), cancellationToken); await RabbitMqBackgroundJobTopology.DeclareAsync( _channel, jobOptions, cancellationToken); return _channel; } private async Task ResetConnectionAsync() { if (_channel is not null) { try { await _channel.DisposeAsync(); } catch { // The broker may already have closed the channel. } _channel = null; } if (_connection is not null) { try { await _connection.DisposeAsync(); } catch { // The broker may already have closed the connection. } _connection = null; } } public async ValueTask DisposeAsync() { await _gate.WaitAsync(); try { await ResetConnectionAsync(); } finally { _gate.Release(); _gate.Dispose(); } } } public sealed class RabbitMqBackgroundJobWorker( BackgroundJobOptions jobOptions, RabbitMqOptions rabbitOptions, BackgroundJobRunner runner, ILogger logger) : BackgroundService { protected override async Task ExecuteAsync(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { IConnection? connection = null; var channels = new List(); try { var consumerCount = RabbitMqBackgroundJobTopology.JobKinds.Sum( jobOptions.ConsumerConcurrency); connection = await RabbitMqBackgroundJobTopology.CreateConnectionAsync( rabbitOptions, "jiaowu-background-job-worker", (ushort)consumerCount, stoppingToken); foreach (var kind in RabbitMqBackgroundJobTopology.JobKinds) { for (var workerIndex = 0; workerIndex < jobOptions.ConsumerConcurrency(kind); workerIndex++) { var channel = await connection.CreateChannelAsync( new CreateChannelOptions( publisherConfirmationsEnabled: false, publisherConfirmationTrackingEnabled: false), stoppingToken); channels.Add(channel); await RabbitMqBackgroundJobTopology.DeclareAsync( channel, jobOptions, stoppingToken); await channel.BasicQosAsync( 0, jobOptions.PrefetchCount, global: false, stoppingToken); var consumer = new AsyncEventingBasicConsumer(channel); consumer.ReceivedAsync += async (_, eventArgs) => { using var deliveryCancellation = CancellationTokenSource.CreateLinkedTokenSource( eventArgs.CancellationToken, stoppingToken); var deliveryToken = deliveryCancellation.Token; try { var message = JsonSerializer .Deserialize( eventArgs.Body.Span); if (message is null || message.JobKind != kind) { await channel.BasicNackAsync( eventArgs.DeliveryTag, multiple: false, requeue: false, deliveryToken); return; } var result = await runner.RunAsync( message, deliveryToken); if (result.Outcome == BackgroundJobRunOutcome.Completed) { await channel.BasicAckAsync( eventArgs.DeliveryTag, multiple: false, deliveryToken); return; } await Task.Delay( result.RetryAfter, deliveryToken); await channel.BasicNackAsync( eventArgs.DeliveryTag, multiple: false, requeue: true, deliveryToken); } catch (OperationCanceledException) { // Closing the channel requeues unacknowledged deliveries. } catch (Exception exception) { logger.LogError( exception, "RabbitMQ delivery for {JobKind} failed and will be requeued.", kind); if (!channel.IsOpen) return; await channel.BasicNackAsync( eventArgs.DeliveryTag, multiple: false, requeue: true, CancellationToken.None); } }; await channel.BasicConsumeAsync( RabbitMqBackgroundJobTopology.QueueName( jobOptions, kind), autoAck: false, consumer, stoppingToken); } } logger.LogInformation( "{ConsumerCount} RabbitMQ background job consumers are connected to {HostName}:{Port}.", consumerCount, rabbitOptions.HostName, rabbitOptions.Port); while (connection.IsOpen && !stoppingToken.IsCancellationRequested) await Task.Delay(TimeSpan.FromSeconds(1), stoppingToken); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { break; } catch (Exception exception) { logger.LogError( exception, "RabbitMQ background job consumer connection failed; retrying."); } finally { foreach (var channel in channels) { try { await channel.DisposeAsync(); } catch { // The connection may already have disposed its channels. } } if (connection is not null) { try { await connection.DisposeAsync(); } catch { // The broker may already have closed the connection. } } } if (!stoppingToken.IsCancellationRequested) await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken); } } } internal static class RabbitMqBackgroundJobTopology { public static readonly BackgroundJobKind[] JobKinds = [ BackgroundJobKind.AutomaticSchedule, BackgroundJobKind.SchedulePublish, BackgroundJobKind.MakeupExamAuto, BackgroundJobKind.ExamArrangement, BackgroundJobKind.ExamSignInExport, BackgroundJobKind.ExamPublish ]; public static async Task CreateConnectionAsync( RabbitMqOptions options, string clientName, ushort consumerDispatchConcurrency, CancellationToken cancellationToken) { var factory = new ConnectionFactory { HostName = options.HostName, Port = options.Port, UserName = options.UserName, Password = options.Password, VirtualHost = options.VirtualHost, ClientProvidedName = clientName, AutomaticRecoveryEnabled = false, TopologyRecoveryEnabled = false, RequestedHeartbeat = TimeSpan.FromSeconds(30), ConsumerDispatchConcurrency = consumerDispatchConcurrency }; if (options.UseTls) { factory.Ssl = new SslOption { Enabled = true, ServerName = string.IsNullOrWhiteSpace(options.TlsServerName) ? options.HostName : options.TlsServerName }; } return await factory.CreateConnectionAsync(cancellationToken); } public static async Task DeclareAsync( IChannel channel, BackgroundJobOptions options, CancellationToken cancellationToken) { var deadExchange = options.Exchange + ".dead"; await channel.ExchangeDeclareAsync( options.Exchange, ExchangeType.Direct, durable: true, autoDelete: false, cancellationToken: cancellationToken); await channel.ExchangeDeclareAsync( deadExchange, ExchangeType.Direct, durable: true, autoDelete: false, cancellationToken: cancellationToken); foreach (var kind in JobKinds) { var routingKey = RoutingKey(kind); var queueName = QueueName(options, kind); var deadQueueName = queueName + ".dead"; var queueArguments = QueueArguments(options); queueArguments["x-dead-letter-exchange"] = deadExchange; queueArguments["x-dead-letter-routing-key"] = routingKey; await channel.QueueDeclareAsync( queueName, durable: true, exclusive: false, autoDelete: false, arguments: queueArguments, cancellationToken: cancellationToken); await channel.QueueBindAsync( queueName, options.Exchange, routingKey, cancellationToken: cancellationToken); await channel.QueueDeclareAsync( deadQueueName, durable: true, exclusive: false, autoDelete: false, arguments: QueueArguments(options), cancellationToken: cancellationToken); await channel.QueueBindAsync( deadQueueName, deadExchange, routingKey, cancellationToken: cancellationToken); } } public static string RoutingKey(BackgroundJobKind kind) => kind switch { BackgroundJobKind.AutomaticSchedule => "schedule.automatic", BackgroundJobKind.SchedulePublish => "schedule.publish", BackgroundJobKind.MakeupExamAuto => "makeup-exam.automatic", BackgroundJobKind.ExamArrangement => "exam.arrangement", BackgroundJobKind.ExamSignInExport => "exam.sign-in-export", BackgroundJobKind.ExamPublish => "exam.publish", _ => throw new ArgumentOutOfRangeException(nameof(kind), kind, null) }; public static string QueueName( BackgroundJobOptions options, BackgroundJobKind kind) => $"{options.QueuePrefix}.{RoutingKey(kind)}"; private static Dictionary QueueArguments( BackgroundJobOptions options) { var arguments = new Dictionary(); if (options.UseQuorumQueues) arguments["x-queue-type"] = "quorum"; return arguments; } }