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; /// /// EmailLog 큐를 폴링하여 SMTP로 송신하는 BackgroundService. /// - WITH (ROWLOCK, UPDLOCK, READPAST) 로 다중 워커 동시성 안전 /// - Pending + NextRetryAt <= now 만 픽업 (시작 시각 cutoff 없음) /// - 실패 시 1차 30s / 2차 5m / 3차 30m 백오프 후 Failed (Dead Letter) /// 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 _logger; private readonly IServiceScopeFactory _scopeFactory; public MailQueueWorker(ILogger 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 ProcessOneAsync(int workerIndex, CancellationToken ct) { using var scope = _scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); var mailService = scope.ServiceProvider.GetRequiredService(); 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]; } }