218 lines
8.7 KiB
C#
218 lines
8.7 KiB
C#
using CMSMicroservice.Application.CommissionCQ.Commands.CalculateWeeklyBalances;
|
|
using CMSMicroservice.Application.CommissionCQ.Commands.CalculateWeeklyCommissionPool;
|
|
using CMSMicroservice.Application.CommissionCQ.Commands.ProcessUserPayouts;
|
|
using CMSMicroservice.Application.CommissionCQ.Commands.TriggerWeeklyCalculation;
|
|
using CMSMicroservice.Application.Common.Interfaces;
|
|
using CMSMicroservice.Domain.Entities;
|
|
using CMSMicroservice.Domain.Enums;
|
|
using MediatR;
|
|
using Microsoft.EntityFrameworkCore;
|
|
using Microsoft.Extensions.Logging;
|
|
using Polly;
|
|
|
|
namespace CMSMicroservice.Infrastructure.BackgroundJobs;
|
|
|
|
/// <summary>
|
|
/// Hangfire Job for weekly commission calculation
|
|
/// Executes every Sunday at 00:05 (Cron: "5 0 * * 0")
|
|
/// </summary>
|
|
public class WeeklyCommissionJob
|
|
{
|
|
private readonly IMediator _mediator;
|
|
private readonly ILogger<WeeklyCommissionJob> _logger;
|
|
private readonly IApplicationDbContext _context;
|
|
private readonly IWeekDefinitionRepository _weekRepository;
|
|
private readonly ResiliencePipeline _retryPipeline;
|
|
|
|
public WeeklyCommissionJob(
|
|
IMediator mediator,
|
|
ILogger<WeeklyCommissionJob> logger,
|
|
IApplicationDbContext context,
|
|
IWeekDefinitionRepository weekRepository)
|
|
{
|
|
_mediator = mediator;
|
|
_logger = logger;
|
|
_context = context;
|
|
_weekRepository = weekRepository;
|
|
|
|
// Polly Retry: 3 attempts, exponential backoff (5min → 10min → 20min)
|
|
_retryPipeline = new ResiliencePipelineBuilder()
|
|
.AddRetry(new Polly.Retry.RetryStrategyOptions
|
|
{
|
|
MaxRetryAttempts = 3,
|
|
Delay = TimeSpan.FromMinutes(5),
|
|
BackoffType = Polly.DelayBackoffType.Exponential,
|
|
UseJitter = true,
|
|
OnRetry = args =>
|
|
{
|
|
_logger.LogWarning(
|
|
"⚠️ Retry attempt {AttemptNumber} after {Delay}ms delay. Exception: {ExceptionType}",
|
|
args.AttemptNumber,
|
|
args.RetryDelay.TotalMilliseconds,
|
|
args.Outcome.Exception?.GetType().Name ?? "None");
|
|
return ValueTask.CompletedTask;
|
|
}
|
|
})
|
|
.Build();
|
|
}
|
|
|
|
/// <summary>
|
|
/// Execute weekly commission calculation with retry logic
|
|
/// Called by Hangfire scheduler or manually triggered
|
|
/// </summary>
|
|
/// <param name="weekDefinitionId">شناسه هفته. اگر null باشد، هفته قبلی محاسبه میشود</param>
|
|
/// <param name="cancellationToken">Cancellation token</param>
|
|
public async Task ExecuteAsync(long? weekDefinitionId = null, CancellationToken cancellationToken = default)
|
|
{
|
|
var executionId = Guid.NewGuid();
|
|
var startTime = DateTime.Now;
|
|
|
|
// Use provided WeekDefinitionId or calculate for PREVIOUS week (completed week)
|
|
long targetWeekDefinitionId;
|
|
if (weekDefinitionId.HasValue && weekDefinitionId.Value > 0)
|
|
{
|
|
targetWeekDefinitionId = weekDefinitionId.Value;
|
|
_logger.LogInformation("📅 Using manually specified WeekDefinitionId: {WeekDefinitionId}", targetWeekDefinitionId);
|
|
}
|
|
else
|
|
{
|
|
var previousWeek = DateTime.Now.AddDays(-7);
|
|
var weekDef = _weekRepository.GetWeekByDate(previousWeek);
|
|
if (weekDef == null)
|
|
{
|
|
throw new InvalidOperationException($"هفته برای تاریخ {previousWeek:yyyy-MM-dd} تعریف نشده است");
|
|
}
|
|
targetWeekDefinitionId = weekDef.Id;
|
|
_logger.LogInformation("📅 Using previous week (auto-calculated): WeekDefinitionId={WeekDefinitionId}, WeekNumber={WeekNumber}",
|
|
targetWeekDefinitionId, weekDef.GregorianWeekNumber);
|
|
}
|
|
|
|
_logger.LogInformation(
|
|
"🚀 [{ExecutionId}] Starting weekly commission calculation for WeekDefinitionId={WeekDefinitionId}",
|
|
executionId, targetWeekDefinitionId);
|
|
|
|
// Create execution log entry
|
|
var log = new WorkerExecutionLog
|
|
{
|
|
ExecutionId = executionId,
|
|
WeekDefinitionId = targetWeekDefinitionId,
|
|
StartedAt = startTime,
|
|
Status = WorkerExecutionStatus.Running
|
|
};
|
|
_context.WorkerExecutionLogs.Add(log);
|
|
await _context.SaveChangesAsync(cancellationToken);
|
|
|
|
try
|
|
{
|
|
// Execute with retry pipeline
|
|
await _retryPipeline.ExecuteAsync(async ct =>
|
|
{
|
|
await ExecuteWeeklyCalculationAsync(executionId, targetWeekDefinitionId, ct);
|
|
}, cancellationToken);
|
|
|
|
// Update log on success
|
|
var completedAt = DateTime.Now;
|
|
var duration = completedAt - startTime;
|
|
|
|
log.Status = WorkerExecutionStatus.Success;
|
|
log.CompletedAt = completedAt;
|
|
log.DurationMs = (long)duration.TotalMilliseconds;
|
|
|
|
// Get counts from database
|
|
var balancesCount = await _context.NetworkWeeklyBalances
|
|
.CountAsync(x => x.WeekDefinitionId == targetWeekDefinitionId, cancellationToken);
|
|
var payoutsCount = await _context.UserCommissionPayouts
|
|
.CountAsync(x => x.WeekDefinitionId == targetWeekDefinitionId, cancellationToken);
|
|
|
|
log.ProcessedCount = balancesCount + payoutsCount;
|
|
|
|
await _context.SaveChangesAsync(cancellationToken);
|
|
|
|
_logger.LogInformation(
|
|
"✅ [{ExecutionId}] Completed successfully in {Duration}s | Balances: {BalancesCount}, Payouts: {PayoutsCount}",
|
|
executionId, duration.TotalSeconds, balancesCount, payoutsCount);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
// Update log on failure
|
|
var completedAt = DateTime.Now;
|
|
var duration = completedAt - startTime;
|
|
|
|
log.Status = WorkerExecutionStatus.Failed;
|
|
log.CompletedAt = completedAt;
|
|
log.DurationMs = (long)duration.TotalMilliseconds;
|
|
log.ErrorMessage = ex.Message;
|
|
log.ErrorStackTrace = ex.StackTrace;
|
|
|
|
await _context.SaveChangesAsync(cancellationToken);
|
|
|
|
_logger.LogError(ex,
|
|
"❌ [{ExecutionId}] Failed after {Duration}s: {ErrorMessage}",
|
|
executionId, duration.TotalSeconds, ex.Message);
|
|
|
|
throw; // Re-throw for Hangfire to mark job as failed
|
|
}
|
|
}
|
|
|
|
private async Task ExecuteWeeklyCalculationAsync(
|
|
Guid executionId,
|
|
long weekDefinitionId,
|
|
CancellationToken cancellationToken)
|
|
{
|
|
// Check idempotency: Skip if already calculated
|
|
var existingPool = await _context.WeeklyCommissionPools
|
|
.FirstOrDefaultAsync(x => x.WeekDefinitionId == weekDefinitionId, cancellationToken);
|
|
|
|
if (existingPool != null && existingPool.IsCalculated)
|
|
{
|
|
_logger.LogWarning(
|
|
"⚠️ [{ExecutionId}] WeekDefinitionId={WeekDefinitionId} already calculated. Skipping.",
|
|
executionId, weekDefinitionId);
|
|
return;
|
|
}
|
|
|
|
// دریافت WeekNumber برای command ها (فعلاً هنوز از WeekNumber استفاده میکنند)
|
|
var GregorianWeekNumber = _weekRepository.GetGregorianWeekNumber(weekDefinitionId);
|
|
if (string.IsNullOrEmpty(GregorianWeekNumber))
|
|
{
|
|
throw new InvalidOperationException($"WeekDefinitionId={weekDefinitionId} یافت نشد");
|
|
}
|
|
|
|
using var transaction = new System.Transactions.TransactionScope(
|
|
System.Transactions.TransactionScopeOption.Required,
|
|
new System.Transactions.TransactionOptions
|
|
{
|
|
IsolationLevel = System.Transactions.IsolationLevel.ReadCommitted,
|
|
Timeout = TimeSpan.FromMinutes(30)
|
|
},
|
|
System.Transactions.TransactionScopeAsyncFlowOption.Enabled);
|
|
|
|
try
|
|
{
|
|
// Step 1: Calculate user balances (Left/Right leg volumes)
|
|
_logger.LogInformation(
|
|
"📊 [{ExecutionId}] Step 1/3: Calculating weekly balances...",
|
|
executionId);
|
|
|
|
await _mediator.Send(new TriggerWeeklyCalculationCommand
|
|
{
|
|
WeekDefinitionId = weekDefinitionId,
|
|
ForceRecalculate = false
|
|
}, cancellationToken);
|
|
|
|
transaction.Complete();
|
|
|
|
_logger.LogInformation(
|
|
"✅ [{ExecutionId}] All steps completed successfully",
|
|
executionId);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogError(ex,
|
|
"❌ [{ExecutionId}] Transaction rolled back: {ErrorMessage}",
|
|
executionId, ex.Message);
|
|
throw;
|
|
}
|
|
}
|
|
}
|