MailQueueWorker.cs 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157
  1. using Application.Abstractions.Messaging.Email;
  2. using Domain.Entities.Common;
  3. using Domain.Entities.Common.ValueObject;
  4. using Infrastructure.Messaging.Email;
  5. using Infrastructure.Persistence;
  6. using Microsoft.EntityFrameworkCore;
  7. using Microsoft.Extensions.DependencyInjection;
  8. using Microsoft.Extensions.Hosting;
  9. using Microsoft.Extensions.Logging;
  10. namespace MailWorker;
  11. /// <summary>
  12. /// EmailLog 큐를 폴링하여 SMTP로 송신하는 BackgroundService.
  13. /// - WITH (ROWLOCK, UPDLOCK, READPAST) 로 다중 워커 동시성 안전
  14. /// - Pending + NextRetryAt <= now 만 픽업 (시작 시각 cutoff 없음)
  15. /// - 실패 시 1차 30s / 2차 5m / 3차 30m 백오프 후 Failed (Dead Letter)
  16. /// </summary>
  17. public sealed class MailQueueWorker : BackgroundService
  18. {
  19. private const int WorkerCount = 2;
  20. private const int IdlePollDelayMs = 3000;
  21. private static readonly TimeSpan[] BackoffDelays =
  22. [
  23. TimeSpan.FromSeconds(30),
  24. TimeSpan.FromMinutes(5),
  25. TimeSpan.FromMinutes(30)
  26. ];
  27. private readonly ILogger<MailQueueWorker> _logger;
  28. private readonly IServiceScopeFactory _scopeFactory;
  29. public MailQueueWorker(ILogger<MailQueueWorker> logger, IServiceScopeFactory scopeFactory)
  30. {
  31. _logger = logger;
  32. _scopeFactory = scopeFactory;
  33. }
  34. protected override Task ExecuteAsync(CancellationToken stoppingToken)
  35. {
  36. _logger.LogInformation("MailQueueWorker started with {WorkerCount} parallel workers", WorkerCount);
  37. var tasks = Enumerable.Range(0, WorkerCount).Select(i => RunLoopAsync(i, stoppingToken)).ToArray();
  38. return Task.WhenAll(tasks);
  39. }
  40. private async Task RunLoopAsync(int workerIndex, CancellationToken stoppingToken)
  41. {
  42. while (!stoppingToken.IsCancellationRequested)
  43. {
  44. try
  45. {
  46. var processed = await ProcessOneAsync(workerIndex, stoppingToken);
  47. if (!processed)
  48. {
  49. await Task.Delay(IdlePollDelayMs, stoppingToken);
  50. }
  51. }
  52. catch (OperationCanceledException)
  53. {
  54. break;
  55. }
  56. catch (Exception ex)
  57. {
  58. _logger.LogError(ex, "Worker[{WorkerIndex}] unhandled error in poll loop", workerIndex);
  59. await Task.Delay(IdlePollDelayMs, stoppingToken);
  60. }
  61. }
  62. }
  63. private async Task<bool> ProcessOneAsync(int workerIndex, CancellationToken ct)
  64. {
  65. using var scope = _scopeFactory.CreateScope();
  66. var db = scope.ServiceProvider.GetRequiredService<AppDbContext>();
  67. var mailService = scope.ServiceProvider.GetRequiredService<DirectMailService>();
  68. EmailLog? entry;
  69. // 트랜잭션으로 큐 항목을 잠그고 Status=Processing 으로 변경
  70. await using (var tx = await db.Database.BeginTransactionAsync(ct))
  71. {
  72. entry = await db.EmailLog
  73. .FromSqlInterpolated($@"
  74. SELECT TOP(1) *
  75. FROM EmailLog WITH (ROWLOCK, UPDLOCK, READPAST)
  76. WHERE [Status] = {(byte)MailStatus.Pending}
  77. AND [NextRetryAt] <= GETUTCDATE()
  78. ORDER BY [NextRetryAt] ASC")
  79. .AsTracking()
  80. .SingleOrDefaultAsync(ct);
  81. if (entry is null)
  82. {
  83. await tx.RollbackAsync(ct);
  84. return false;
  85. }
  86. entry.Status = MailStatus.Processing;
  87. await db.SaveChangesAsync(ct);
  88. await tx.CommitAsync(ct);
  89. }
  90. // 트랜잭션 밖에서 SMTP 발송 시도
  91. var sendData = new SendData(entry.ToAddress, entry.Subject, entry.MessageHtml, entry.MessageText);
  92. try
  93. {
  94. await mailService.SendAsync(sendData, ct);
  95. entry.Status = MailStatus.Sent;
  96. entry.ProcessedAt = DateTime.UtcNow;
  97. entry.LastError = null;
  98. _logger.LogInformation("Worker[{WorkerIndex}] sent (ID={ID}) to {ToAddress}", workerIndex, entry.ID, entry.ToAddress);
  99. }
  100. catch (OperationCanceledException)
  101. {
  102. // 셧다운 중. 다음 시작 시 재시도되도록 Pending으로 되돌림 (RetryCount 그대로)
  103. entry.Status = MailStatus.Pending;
  104. await db.SaveChangesAsync(CancellationToken.None);
  105. throw;
  106. }
  107. catch (Exception ex)
  108. {
  109. entry.RetryCount += 1;
  110. entry.LastError = Truncate(ex.Message, 500);
  111. if (entry.RetryCount > entry.MaxRetryCount)
  112. {
  113. entry.Status = MailStatus.Failed;
  114. entry.FailedAt = DateTime.UtcNow;
  115. _logger.LogError(ex, "Worker[{WorkerIndex}] dead-letter (ID={ID}) to {ToAddress} after {RetryCount} retries",
  116. workerIndex, entry.ID, entry.ToAddress, entry.RetryCount);
  117. }
  118. else
  119. {
  120. var delay = BackoffDelays[Math.Min(entry.RetryCount - 1, BackoffDelays.Length - 1)];
  121. entry.Status = MailStatus.Pending;
  122. entry.NextRetryAt = DateTime.UtcNow + delay;
  123. _logger.LogWarning(ex, "Worker[{WorkerIndex}] retry {RetryCount}/{MaxRetryCount} (ID={ID}) to {ToAddress} after {DelaySeconds}s",
  124. workerIndex, entry.RetryCount, entry.MaxRetryCount, entry.ID, entry.ToAddress, delay.TotalSeconds);
  125. }
  126. }
  127. await db.SaveChangesAsync(ct);
  128. return true;
  129. }
  130. private static string Truncate(string value, int maxLength)
  131. {
  132. if (string.IsNullOrEmpty(value)) {
  133. return string.Empty;
  134. }
  135. return value.Length <= maxLength ? value : value[..maxLength];
  136. }
  137. }