using System.Collections.Concurrent; using AccountsDAL.CustomCode.Warehouse; using GB5Shared.DTO.Framework.Login; using GB5Shared.EventLogPublish; using GB5Shared.QueryExecutor; using GB5Shared.Telemetry; using Microsoft.Data.SqlClient; using Microsoft.Extensions.Configuration; using static GB5Shared.GB5Constant.Constant; namespace AccountsBLL.Warehouse { public class WarehouseFactPostingBLL : IWarehouseFactPostingBLL { private const int DefaultMaxConcurrentOus = 8; private const int SqlDeadlockVictimErrorNumber = 1205; private const int MaxDeadlockRetries = 3; private readonly IWarehouseFactPostingDAL _warehouseFactPostingDAL; private readonly IQueryExecutor _queryExecutor; private readonly EventLogPublish _eventLog; private readonly int _maxConcurrentOus; public WarehouseFactPostingBLL( IWarehouseFactPostingDAL warehouseFactPostingDAL, IQueryExecutor queryExecutor, EventLogPublish eventLog, IConfiguration configuration) { _warehouseFactPostingDAL = warehouseFactPostingDAL; _queryExecutor = queryExecutor; _eventLog = eventLog; var configuredValue = configuration.GetValue("Warehouse:MaxConcurrentOus"); _maxConcurrentOus = configuredValue is > 0 ? configuredValue.Value : DefaultMaxConcurrentOus; } public async Task PostFactAsync( int factId, DateTime dateFrom, DateTime dateTo, List? ouIds, LoginDTO loginDTO, CancellationToken ct) { if (dateFrom > dateTo) throw new InvalidOperationException("dateFrom must not be after dateTo."); var targetDates = new List(); for (var d = dateFrom.Date; d <= dateTo.Date; d = d.AddDays(1)) targetDates.Add(d); var resolvedOuIds = ouIds is { Count: > 0 } ? ouIds : await _warehouseFactPostingDAL.GetActiveOuIdsAsync(loginDTO, ct).ConfigureAwait(false); var succeededOuIds = new ConcurrentBag(); var failedOuIds = new ConcurrentDictionary(); await Parallel.ForEachAsync( resolvedOuIds, new ParallelOptions { MaxDegreeOfParallelism = _maxConcurrentOus, CancellationToken = ct }, async (ouId, innerCt) => { for (var attempt = 1; attempt <= MaxDeadlockRetries; attempt++) { GB5Trace.Step("post-fact-ou-start", new { factId, ouId, dateFrom, dateTo, attempt }); var tx = await _queryExecutor.BeginTransactionAsync(loginDTO).ConfigureAwait(false); try { await _warehouseFactPostingDAL .PostFactForOuAsync(factId, ouId, targetDates, loginDTO, tx, innerCt) .ConfigureAwait(false); await _queryExecutor.CommitAsync(tx).ConfigureAwait(false); succeededOuIds.Add(ouId); return; } catch (SqlException ex) when (ex.Number == SqlDeadlockVictimErrorNumber && attempt < MaxDeadlockRetries) { await _queryExecutor.RollbackAsync(tx).ConfigureAwait(false); GB5Trace.Step("post-fact-ou-deadlock-retry", new { factId, ouId, attempt }); await Task.Delay(TimeSpan.FromMilliseconds(200 * attempt), innerCt).ConfigureAwait(false); } catch (Exception ex) { await _queryExecutor.RollbackAsync(tx).ConfigureAwait(false); GB5Trace.MarkFailed("post-fact-ou-failed", ex); failedOuIds[ouId] = ex.Message; return; } } }).ConfigureAwait(false); var result = new WarehouseFactPostingResultDTO(); result.SucceededOuIds.AddRange(succeededOuIds); foreach (var kv in failedOuIds) result.FailedOuIds[kv.Key] = kv.Value; var allSucceeded = result.FailedOuIds.Count == 0; // Audit event — fired for both success and partial failure so dashboards and alerting // can track posting runs without polling MWAREHOUSEFACT.LASTPOSTEDON directly. GB5Trace.Step("event-publish", new { factId, allSucceeded }); await _eventLog.PublishEventLogAsync( allSucceeded ? "Warehouse Fact Posted" : "Warehouse Fact Partially Failed", new { factId, result.SucceededOuIds, result.FailedOuIds }, allSucceeded ? EventTypeConstant.POSTWAREHOUSEFACTSUCCESSEVENTTYPEID : EventTypeConstant.POSTWAREHOUSEFACTFAILEDEVENTTYPEID, factId, loginDTO, ct: ct).ConfigureAwait(false); // Freshness stamp — only on complete success (all OUs posted without error). // LASTPOSTEDDATEID = DimDate.DateKey format (YYYYMMDD) for dateTo. if (allSucceeded) { var lastPostedDateId = int.Parse(dateTo.ToString("yyyyMMdd")); await _warehouseFactPostingDAL .UpdateLastPostedAsync(factId, lastPostedDateId, loginDTO, ct) .ConfigureAwait(false); } return result; } public async Task> PostAllFactsAsync( DateTime dateFrom, DateTime dateTo, LoginDTO loginDTO, CancellationToken ct) { GB5Trace.Step("post-all-facts", new { dateFrom, dateTo }); var factIds = await _warehouseFactPostingDAL .GetAllActiveFactIdsAsync(loginDTO, ct) .ConfigureAwait(false); var results = new Dictionary(); foreach (var factId in factIds) { ct.ThrowIfCancellationRequested(); try { GB5Trace.Step("post-all-facts-factid", new { factId }); var result = await PostFactAsync(factId, dateFrom, dateTo, null, loginDTO, ct) .ConfigureAwait(false); results[factId] = result; } catch (Exception ex) { GB5Trace.MarkFailed("post-all-facts-factid-failed", ex); var failed = new WarehouseFactPostingResultDTO(); failed.FailedOuIds[-1] = ex.Message; results[factId] = failed; } } return results; } } }