| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157 |
- using Application.Abstractions.Messaging.Email;
- using Domain.Entities.Common;
- using Domain.Entities.Common.ValueObject;
- using Infrastructure.Messaging.Email;
- using Infrastructure.Persistence;
- using Microsoft.EntityFrameworkCore;
- using Microsoft.Extensions.DependencyInjection;
- using Microsoft.Extensions.Hosting;
- using Microsoft.Extensions.Logging;
- namespace MailWorker;
- /// <summary>
- /// EmailLog 큐를 폴링하여 SMTP로 송신하는 BackgroundService.
- /// - WITH (ROWLOCK, UPDLOCK, READPAST) 로 다중 워커 동시성 안전
- /// - Pending + NextRetryAt <= now 만 픽업 (시작 시각 cutoff 없음)
- /// - 실패 시 1차 30s / 2차 5m / 3차 30m 백오프 후 Failed (Dead Letter)
- /// </summary>
- public sealed class MailQueueWorker : BackgroundService
- {
- private const int WorkerCount = 2;
- private const int IdlePollDelayMs = 3000;
- private static readonly TimeSpan[] BackoffDelays =
- [
- TimeSpan.FromSeconds(30),
- TimeSpan.FromMinutes(5),
- TimeSpan.FromMinutes(30)
- ];
- private readonly ILogger<MailQueueWorker> _logger;
- private readonly IServiceScopeFactory _scopeFactory;
- public MailQueueWorker(ILogger<MailQueueWorker> logger, IServiceScopeFactory scopeFactory)
- {
- _logger = logger;
- _scopeFactory = scopeFactory;
- }
- protected override Task ExecuteAsync(CancellationToken stoppingToken)
- {
- _logger.LogInformation("MailQueueWorker started with {WorkerCount} parallel workers", WorkerCount);
- var tasks = Enumerable.Range(0, WorkerCount).Select(i => RunLoopAsync(i, stoppingToken)).ToArray();
- return Task.WhenAll(tasks);
- }
- private async Task RunLoopAsync(int workerIndex, CancellationToken stoppingToken)
- {
- while (!stoppingToken.IsCancellationRequested)
- {
- try
- {
- var processed = await ProcessOneAsync(workerIndex, stoppingToken);
- if (!processed)
- {
- await Task.Delay(IdlePollDelayMs, stoppingToken);
- }
- }
- catch (OperationCanceledException)
- {
- break;
- }
- catch (Exception ex)
- {
- _logger.LogError(ex, "Worker[{WorkerIndex}] unhandled error in poll loop", workerIndex);
- await Task.Delay(IdlePollDelayMs, stoppingToken);
- }
- }
- }
- private async Task<bool> ProcessOneAsync(int workerIndex, CancellationToken ct)
- {
- using var scope = _scopeFactory.CreateScope();
- var db = scope.ServiceProvider.GetRequiredService<AppDbContext>();
- var mailService = scope.ServiceProvider.GetRequiredService<DirectMailService>();
- EmailLog? entry;
- // 트랜잭션으로 큐 항목을 잠그고 Status=Processing 으로 변경
- await using (var tx = await db.Database.BeginTransactionAsync(ct))
- {
- entry = await db.EmailLog
- .FromSqlInterpolated($@"
- SELECT TOP(1) *
- FROM EmailLog WITH (ROWLOCK, UPDLOCK, READPAST)
- WHERE [Status] = {(byte)MailStatus.Pending}
- AND [NextRetryAt] <= GETUTCDATE()
- ORDER BY [NextRetryAt] ASC")
- .AsTracking()
- .SingleOrDefaultAsync(ct);
- if (entry is null)
- {
- await tx.RollbackAsync(ct);
- return false;
- }
- entry.Status = MailStatus.Processing;
- await db.SaveChangesAsync(ct);
- await tx.CommitAsync(ct);
- }
- // 트랜잭션 밖에서 SMTP 발송 시도
- var sendData = new SendData(entry.ToAddress, entry.Subject, entry.MessageHtml, entry.MessageText);
- try
- {
- await mailService.SendAsync(sendData, ct);
- entry.Status = MailStatus.Sent;
- entry.ProcessedAt = DateTime.UtcNow;
- entry.LastError = null;
- _logger.LogInformation("Worker[{WorkerIndex}] sent (ID={ID}) to {ToAddress}", workerIndex, entry.ID, entry.ToAddress);
- }
- catch (OperationCanceledException)
- {
- // 셧다운 중. 다음 시작 시 재시도되도록 Pending으로 되돌림 (RetryCount 그대로)
- entry.Status = MailStatus.Pending;
- await db.SaveChangesAsync(CancellationToken.None);
- throw;
- }
- catch (Exception ex)
- {
- entry.RetryCount += 1;
- entry.LastError = Truncate(ex.Message, 500);
- if (entry.RetryCount > entry.MaxRetryCount)
- {
- entry.Status = MailStatus.Failed;
- entry.FailedAt = DateTime.UtcNow;
- _logger.LogError(ex, "Worker[{WorkerIndex}] dead-letter (ID={ID}) to {ToAddress} after {RetryCount} retries",
- workerIndex, entry.ID, entry.ToAddress, entry.RetryCount);
- }
- else
- {
- var delay = BackoffDelays[Math.Min(entry.RetryCount - 1, BackoffDelays.Length - 1)];
- entry.Status = MailStatus.Pending;
- entry.NextRetryAt = DateTime.UtcNow + delay;
- _logger.LogWarning(ex, "Worker[{WorkerIndex}] retry {RetryCount}/{MaxRetryCount} (ID={ID}) to {ToAddress} after {DelaySeconds}s",
- workerIndex, entry.RetryCount, entry.MaxRetryCount, entry.ID, entry.ToAddress, delay.TotalSeconds);
- }
- }
- await db.SaveChangesAsync(ct);
- return true;
- }
- private static string Truncate(string value, int maxLength)
- {
- if (string.IsNullOrEmpty(value)) {
- return string.Empty;
- }
- return value.Length <= maxLength ? value : value[..maxLength];
- }
- }
|