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