using System.Data.Common; using System.Text.Json; using System.Text.RegularExpressions; using AccountsDAL.DTO.Warehouse; using AccountsDAL.Query.Warehouse; using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; namespace AccountsDAL.CustomCode.Warehouse { // Generic warehouse fact-posting engine — reads declarative rules from // MWAREHOUSEMEASUREPOSTINGRULE and compiles them into parameterized SQL. // // Supported SOURCETYPE values after the 20260827 redesign: // 0 = NatureBucketAggregation (TVOUCHER/MACCOUNT/MACCOUNTGROUP source — existing GL path) // 1 = DerivedFormula (C# arithmetic expression over already-computed measures) // 2 = ModuleAggregation (any module table; columns declared in the rule row) // // 1:N posting rules per measure: two rules can share the same MeasureCode and differ by // CombinationMode (0=Add, 1=Subtract). The engine uses RULE_{PostingRuleId} as the SQL // alias to avoid duplicate column names, then folds contributions into MeasureCode using // MergeContribution() before writing the fact row. public class WarehouseFactPostingDAL : IWarehouseFactPostingDAL { // Defense in depth: all SQL identifier values (table names, column names, measure codes, // nature values) from admin-maintained metadata are still validated before interpolation // into SQL text — a typo'd metadata row should fail loudly, not silently break SQL. private static readonly Regex SafeIdentifier = new(@"^[A-Z0-9_]+$", RegexOptions.Compiled); private readonly IQueryExecutor _queryExecutor; public WarehouseFactPostingDAL(IQueryExecutor queryExecutor) { _queryExecutor = queryExecutor; } public async Task PostFactForOuAsync( int factId, int ouId, List targetDates, LoginDTO loginDTO, DbTransaction tx, CancellationToken ct) { var factTableName = await _queryExecutor .ExecuteScalarAsync(loginDTO, WarehouseFactPostingQB.GET_FACT_TABLE_NAME, new { FactId = factId }, tx, cancellationToken: ct) .ConfigureAwait(false); if (string.IsNullOrEmpty(factTableName)) throw new InvalidOperationException($"No active MWAREHOUSEFACT row found for FactId {factId}."); factTableName = ValidateIdentifier(factTableName); var rules = (await _queryExecutor .QueryAsync(loginDTO, WarehouseFactPostingQB.GET_POSTING_RULES, new { FactId = factId }, tx, cancellationToken: ct) .ConfigureAwait(false)).ToList(); var natureBucketRules = rules.Where(r => r.SourceType == 0).ToList(); var formulaRules = rules.Where(r => r.SourceType == 1).ToList(); var moduleAggRules = rules.Where(r => r.SourceType == 2).ToList(); var dateIdsByDate = await GetDateIdsAsync(targetDates, loginDTO, tx, ct).ConfigureAwait(false); var measureValuesByDate = targetDates.ToDictionary( d => d, _ => new Dictionary(StringComparer.OrdinalIgnoreCase)); // ── SOURCETYPE=0: NatureBucketAggregation ──────────────────────────── var periodFlowGroups = natureBucketRules .Where(r => r.AggregationMode == 0) .GroupBy(r => (r.NatureColumnType, r.AccountTypeFilter)); foreach (var group in periodFlowGroups) { var (natureColumnType, accountTypeFilter) = group.Key; var rowsByDate = await ExecutePeriodFlowMultiDateBatchAsync( group.ToList(), natureColumnType!.Value, accountTypeFilter, ouId, targetDates, loginDTO, tx, ct) .ConfigureAwait(false); foreach (var targetDate in targetDates) { var row = rowsByDate.TryGetValue(targetDate.Date, out var found) ? found : ZeroRow(group); MergeRowIntoAccumulator(measureValuesByDate[targetDate], row); } } var asOfBalanceGroups = natureBucketRules .Where(r => r.AggregationMode == 1) .GroupBy(r => (r.NatureColumnType, r.AccountTypeFilter)); foreach (var targetDate in targetDates) { foreach (var group in asOfBalanceGroups) { var (natureColumnType, accountTypeFilter) = group.Key; var row = await ExecuteFlatBatchAsync( group.ToList(), natureColumnType!.Value, 1, accountTypeFilter, ouId, targetDate, loginDTO, tx, ct).ConfigureAwait(false); MergeRowIntoAccumulator(measureValuesByDate[targetDate], row); } var signSplitGroups = natureBucketRules .Where(r => r.AggregationMode is 2 or 3) .GroupBy(r => (r.NatureColumnType, r.NatureValues, r.AccountTypeFilter)); foreach (var group in signSplitGroups) { var row = await ExecuteSignSplitBatchAsync( group.ToList(), ouId, targetDate, loginDTO, tx, ct).ConfigureAwait(false); MergeRowIntoAccumulator(measureValuesByDate[targetDate], row); } // ── SOURCETYPE=2: ModuleAggregation ────────────────────────────── // Each rule targets a different source table / column combination — NOT batched // (different source tables cannot be collapsed into one query). Processed one rule // at a time; still safe because these are expected to be few and fast SELECTs. foreach (var rule in moduleAggRules) { var contribution = await ExecuteModuleAggregationAsync(rule, ouId, targetDate, loginDTO, tx, ct) .ConfigureAwait(false); MergeContribution(measureValuesByDate[targetDate], rule.MeasureCode, contribution, rule.CombinationMode); } // ── SOURCETYPE=1: DerivedFormula ────────────────────────────────── // Evaluated in MEASUREID, SORTORDER order (the query's ORDER BY) so earlier-defined // formulas are available to later-defined ones. The MVP set has one chain (EBIT → EBITDA); // the ordering dependency is documented here because it is NOT enforced by the schema. var measureValues = measureValuesByDate[targetDate]; foreach (var rule in formulaRules) measureValues[rule.MeasureCode] = EvaluateFormula(rule.FormulaExpression!, measureValues); } await UpsertFactRowsAsync(factTableName, ouId, targetDates, dateIdsByDate, measureValuesByDate, rules, loginDTO, tx, ct) .ConfigureAwait(false); } public async Task UpdateLastPostedAsync(int factId, int lastPostedDateId, LoginDTO loginDTO, CancellationToken ct) { await _queryExecutor .ExecuteAsync(loginDTO, WarehouseFactPostingQB.UPDATE_LAST_POSTED, new { FactId = factId, LastPostedDateId = lastPostedDateId }, cancellationToken: ct) .ConfigureAwait(false); } public async Task> GetActiveOuIdsAsync(LoginDTO loginDTO, CancellationToken ct) { var ouIds = await _queryExecutor .QueryAsync(loginDTO, WarehouseFactPostingQB.GET_ACTIVE_OU_IDS, cancellationToken: ct) .ConfigureAwait(false); return ouIds.ToList(); } public async Task> GetAllActiveFactIdsAsync(LoginDTO loginDTO, CancellationToken ct) { var factIds = await _queryExecutor .QueryAsync(loginDTO, WarehouseFactPostingQB.GET_ALL_ACTIVE_FACT_IDS, cancellationToken: ct) .ConfigureAwait(false); return factIds.ToList(); } // ── SOURCETYPE=2 — ModuleAggregation ───────────────────────────────────── private async Task ExecuteModuleAggregationAsync( WarehouseMeasurePostingRuleDTO rule, int ouId, DateTime targetDate, LoginDTO loginDTO, DbTransaction tx, CancellationToken ct) { var (sql, param) = BuildModuleAggregationSql(rule, ouId, targetDate); var result = await _queryExecutor .ExecuteScalarAsync(loginDTO, sql, param, tx, cancellationToken: ct) .ConfigureAwait(false); return result ?? 0m; } private static (string Sql, Dapper.DynamicParameters Param) BuildModuleAggregationSql( WarehouseMeasurePostingRuleDTO rule, int ouId, DateTime targetDate) { if (string.IsNullOrWhiteSpace(rule.SourceTable)) throw new InvalidOperationException($"PostingRuleId {rule.PostingRuleId}: SOURCETABLE is required for SOURCETYPE=2."); if (string.IsNullOrWhiteSpace(rule.SourceDateColumn)) throw new InvalidOperationException($"PostingRuleId {rule.PostingRuleId}: SOURCEDATECOLUMN is required for SOURCETYPE=2."); if (string.IsNullOrWhiteSpace(rule.SourceOuColumn)) throw new InvalidOperationException($"PostingRuleId {rule.PostingRuleId}: SOURCEOUCOLUMN is required for SOURCETYPE=2."); if (string.IsNullOrWhiteSpace(rule.SourceAmtColumn)) throw new InvalidOperationException($"PostingRuleId {rule.PostingRuleId}: SOURCEAMTCOLUMN is required for SOURCETYPE=2."); var table = ValidateIdentifier(rule.SourceTable.ToUpperInvariant()); var dateCol = ValidateIdentifier(rule.SourceDateColumn.ToUpperInvariant()); var ouCol = ValidateIdentifier(rule.SourceOuColumn.ToUpperInvariant()); var amtCol = ValidateIdentifier(rule.SourceAmtColumn.ToUpperInvariant()); var param = new Dapper.DynamicParameters(); param.Add("TargetDate", targetDate.Date); param.Add("OuId", ouId); var extraClauses = new System.Text.StringBuilder(); if (!string.IsNullOrWhiteSpace(rule.FilterJson)) { var filters = JsonSerializer.Deserialize>(rule.FilterJson) ?? throw new InvalidOperationException($"PostingRuleId {rule.PostingRuleId}: FILTERJSON is not valid JSON."); var i = 0; foreach (var kv in filters) { var col = ValidateIdentifier(kv.Key.ToUpperInvariant()); var pName = $"fp{i++}"; extraClauses.Append($" AND {col} = @{pName}"); if (kv.Value.ValueKind == JsonValueKind.Number) param.Add(pName, kv.Value.GetDecimal()); else if (kv.Value.ValueKind == JsonValueKind.String) param.Add(pName, kv.Value.GetString()); else param.Add(pName, kv.Value.GetRawText()); } } var sql = $"SELECT SUM({amtCol}) FROM {table} WHERE CAST({dateCol} AS DATE) = CAST(@TargetDate AS DATE) AND {ouCol} = @OuId{extraClauses}"; return (sql, param); } // ── NatureBucket helpers (SOURCETYPE=0) ────────────────────────────────── private static Dictionary ZeroRow(IEnumerable rulesInGroup) { var result = new Dictionary(StringComparer.OrdinalIgnoreCase); foreach (var rule in rulesInGroup) result[rule.MeasureCode] = 0m; return result; } private async Task> GetDateIdsAsync( List targetDates, LoginDTO loginDTO, DbTransaction tx, CancellationToken ct) { var parameters = new Dapper.DynamicParameters(); var dateParamNames = new List(); for (var i = 0; i < targetDates.Count; i++) { var paramName = $"d{i}"; parameters.Add(paramName, targetDates[i].Date); dateParamNames.Add("@" + paramName); } var sql = WarehouseFactPostingQB.GET_DATE_KEYS.Replace("{TargetDates}", string.Join(", ", dateParamNames)); var rows = await _queryExecutor .QueryAsync<(DateTime TargetDate, int DateKey)>(loginDTO, sql, parameters, tx, cancellationToken: ct) .ConfigureAwait(false); var result = rows.ToDictionary(r => r.TargetDate.Date, r => r.DateKey); var missing = targetDates.Select(d => d.Date).Where(d => !result.ContainsKey(d)).ToList(); if (missing.Count > 0) throw new InvalidOperationException( $"No DimDate row found for {string.Join(", ", missing.Select(d => d.ToString("yyyy-MM-dd")))} — DimDate needs extending."); return result; } private async Task>> ExecutePeriodFlowMultiDateBatchAsync( List rulesInGroup, byte natureColumnType, string? accountTypeFilter, int ouId, List targetDates, LoginDTO loginDTO, DbTransaction tx, CancellationToken ct) { var natureColumnExpr = natureColumnType == 1 ? "AG.ACCOUNTGROUPNATURE" : "ACC.ACCOUNTSCHEDULENATURE"; var accountTypeClause = string.IsNullOrEmpty(accountTypeFilter) ? string.Empty : $" AND ACC.ACCOUNTTYPE IN ({ValidateIntList(accountTypeFilter)})"; // Use RULE_{PostingRuleId} as alias — unique even when two rules share the same MeasureCode. var selectColumns = string.Join(",\n ", rulesInGroup.Select(r => { var natureValues = ValidateIntList(r.NatureValues!); var signedAmount = r.SignConvention == 1 ? $"CASE WHEN VD.DETAILTYPE = 0 THEN -({WarehouseFactPostingQB.BASE_AMOUNT_EXPR}) ELSE {WarehouseFactPostingQB.BASE_AMOUNT_EXPR} END" : $"CASE WHEN VD.DETAILTYPE = 0 THEN {WarehouseFactPostingQB.BASE_AMOUNT_EXPR} ELSE -({WarehouseFactPostingQB.BASE_AMOUNT_EXPR}) END"; return $"SUM(CASE WHEN {natureColumnExpr} IN ({natureValues}){accountTypeClause} THEN {signedAmount} ELSE 0 END) AS RULE_{r.PostingRuleId}"; })); var template = natureColumnType == 1 ? WarehouseFactPostingQB.BUILD_FLAT_BATCH_GROUP_MULTIDATE : WarehouseFactPostingQB.BUILD_FLAT_BATCH_SCHEDULE_MULTIDATE; var parameters = new Dapper.DynamicParameters(); parameters.Add("OUID", ouId); var dateParamNames = new List(); for (var i = 0; i < targetDates.Count; i++) { var paramName = $"d{i}"; parameters.Add(paramName, targetDates[i].Date); dateParamNames.Add("@" + paramName); } var sql = template .Replace("{SelectColumns}", selectColumns) .Replace("{TargetDates}", string.Join(", ", dateParamNames)); var rows = await _queryExecutor .QueryAsync(loginDTO, sql, parameters, tx, cancellationToken: ct) .ConfigureAwait(false); var result = new Dictionary>(); foreach (var row in rows) { var dict = (IDictionary)row!; var activityDate = Convert.ToDateTime(dict["ActivityDate"]).Date; result[activityDate] = DynamicRowToMeasureValues(row, rulesInGroup); } return result; } private async Task> ExecuteFlatBatchAsync( List rulesInGroup, byte natureColumnType, byte aggregationMode, string? accountTypeFilter, int ouId, DateTime targetDate, LoginDTO loginDTO, DbTransaction tx, CancellationToken ct) { var natureColumnExpr = natureColumnType == 1 ? "AG.ACCOUNTGROUPNATURE" : "ACC.ACCOUNTSCHEDULENATURE"; var accountTypeClause = string.IsNullOrEmpty(accountTypeFilter) ? string.Empty : $" AND ACC.ACCOUNTTYPE IN ({ValidateIntList(accountTypeFilter)})"; var selectColumns = string.Join(",\n ", rulesInGroup.Select(r => { var natureValues = ValidateIntList(r.NatureValues!); var signedAmount = r.SignConvention == 1 ? $"CASE WHEN VD.DETAILTYPE = 0 THEN -({WarehouseFactPostingQB.BASE_AMOUNT_EXPR}) ELSE {WarehouseFactPostingQB.BASE_AMOUNT_EXPR} END" : $"CASE WHEN VD.DETAILTYPE = 0 THEN {WarehouseFactPostingQB.BASE_AMOUNT_EXPR} ELSE -({WarehouseFactPostingQB.BASE_AMOUNT_EXPR}) END"; return $"SUM(CASE WHEN {natureColumnExpr} IN ({natureValues}){accountTypeClause} THEN {signedAmount} ELSE 0 END) AS RULE_{r.PostingRuleId}"; })); var dateFilter = aggregationMode == 0 ? "V.VOUCHERDATE = @TargetDate" : "V.VOUCHERDATE <= @TargetDate"; var template = natureColumnType == 1 ? WarehouseFactPostingQB.BUILD_FLAT_BATCH_GROUP : WarehouseFactPostingQB.BUILD_FLAT_BATCH_SCHEDULE; var sql = template.Replace("{SelectColumns}", selectColumns).Replace("{DateFilter}", dateFilter); var row = await _queryExecutor .QuerySingleAsync(loginDTO, sql, new { OUID = ouId, TargetDate = targetDate }, tx, cancellationToken: ct) .ConfigureAwait(false); return DynamicRowToMeasureValues(row, rulesInGroup); } private async Task> ExecuteSignSplitBatchAsync( List rulesInGroup, int ouId, DateTime targetDate, LoginDTO loginDTO, DbTransaction tx, CancellationToken ct) { var natureValues = ValidateIntList(rulesInGroup[0].NatureValues!); var accountTypeFilter = rulesInGroup[0].AccountTypeFilter; var accountTypeClause = string.IsNullOrEmpty(accountTypeFilter) ? string.Empty : $"AND ACC.ACCOUNTTYPE IN ({ValidateIntList(accountTypeFilter)})"; var outerColumns = string.Join(",\n ", rulesInGroup.Select(r => { var isDebitOnly = r.AggregationMode == 2; var sideExpr = isDebitOnly ? "SUM(CASE WHEN X.NetAmount > 0 THEN X.NetAmount ELSE 0 END)" : "SUM(CASE WHEN X.NetAmount < 0 THEN -X.NetAmount ELSE 0 END)"; return $"{sideExpr} AS RULE_{r.PostingRuleId}"; })); var sql = WarehouseFactPostingQB.BUILD_SIGNSPLIT_BATCH .Replace("{OuterSelectColumns}", outerColumns) .Replace("{NatureValues}", natureValues) .Replace("{AccountTypeFilter}", accountTypeClause); var row = await _queryExecutor .QuerySingleAsync(loginDTO, sql, new { OUID = ouId, TargetDate = targetDate }, tx, cancellationToken: ct) .ConfigureAwait(false); return DynamicRowToMeasureValues(row, rulesInGroup); } // Reads RULE_{PostingRuleId} aliases from the dynamic row and applies CombinationMode // (Add/Subtract) to fold contributions into the final MeasureCode dictionary. private static Dictionary DynamicRowToMeasureValues( dynamic row, List rulesInGroup) { var result = new Dictionary(StringComparer.OrdinalIgnoreCase); var dict = (IDictionary)row!; foreach (var rule in rulesInGroup) { var alias = $"RULE_{rule.PostingRuleId}"; var contribution = dict.TryGetValue(alias, out var value) && value is not null and not DBNull ? Convert.ToDecimal(value) : 0m; MergeContribution(result, rule.MeasureCode, contribution, rule.CombinationMode); } return result; } // Adds or subtracts a rule's contribution to the MeasureCode accumulator. // CombinationMode 0=Add (default), 1=Subtract. private static void MergeContribution( Dictionary accumulator, string measureCode, decimal contribution, byte combinationMode) { accumulator.TryGetValue(measureCode, out var existing); accumulator[measureCode] = combinationMode == 1 ? existing - contribution : existing + contribution; } // Merges a group-level result into the per-date accumulator. // Uses += so multiple groups contributing to the same MeasureCode are summed correctly // (different groups cannot share a MeasureCode by schema design; += is defensive). private static void MergeRowIntoAccumulator( Dictionary accumulator, Dictionary row) { foreach (var kv in row) { accumulator.TryGetValue(kv.Key, out var existing); accumulator[kv.Key] = existing + kv.Value; } } private static decimal EvaluateFormula(string expression, Dictionary measureValues) { var matches = Regex.Matches(expression, @"([+-]?)([A-Z0-9_]+)"); decimal total = 0m; foreach (Match m in matches) { var sign = m.Groups[1].Value == "-" ? -1m : 1m; var code = m.Groups[2].Value; if (!measureValues.TryGetValue(code, out var value)) throw new InvalidOperationException( $"FORMULAEXPRESSION '{expression}' references measure '{code}', which has no computed value yet — check MEASUREID/SORTORDER ordering."); total += sign * value; } return total; } private async Task UpsertFactRowsAsync( string factTableName, int ouId, List targetDates, Dictionary dateIdsByDate, Dictionary> measureValuesByDate, List allRules, LoginDTO loginDTO, DbTransaction tx, CancellationToken ct) { // Use distinct ColumnName list — multiple posting rules may map to the same measure // column (1:N rules), but the fact table column appears only once. var columnNames = allRules.Select(r => r.ColumnName).Distinct(StringComparer.OrdinalIgnoreCase).ToList(); var dateIds = targetDates.Select(d => dateIdsByDate[d.Date]).ToList(); var deleteParameters = new Dapper.DynamicParameters(); deleteParameters.Add("OUID", ouId); var dateIdParamNames = new List(); for (var i = 0; i < dateIds.Count; i++) { var paramName = $"dateId{i}"; deleteParameters.Add(paramName, dateIds[i]); dateIdParamNames.Add("@" + paramName); } var deleteSql = WarehouseFactPostingQB.DELETE_EXISTING_ROWS .Replace("{FactTableName}", factTableName) .Replace("{DateIdList}", string.Join(", ", dateIdParamNames)); await _queryExecutor.ExecuteAsync(loginDTO, deleteSql, deleteParameters, tx, ct).ConfigureAwait(false); var insertParameters = new Dapper.DynamicParameters(); insertParameters.Add("OUID", ouId); // Build column→MeasureCode map for looking up the final value. // If two rules write to the same ColumnName, the last MergeContribution result wins. var columnToMeasureCode = allRules .GroupBy(r => r.ColumnName, StringComparer.OrdinalIgnoreCase) .ToDictionary(g => g.Key, g => g.Last().MeasureCode, StringComparer.OrdinalIgnoreCase); var valueRows = new List(); for (var i = 0; i < targetDates.Count; i++) { var targetDate = targetDates[i]; var measureValues = measureValuesByDate[targetDate]; var dateIdParamName = $"rowDateId{i}"; insertParameters.Add(dateIdParamName, dateIdsByDate[targetDate.Date]); var rowValueParamNames = new List { "@" + dateIdParamName }; foreach (var col in columnNames) { var paramName = $"row{i}_{col}"; var measureCode = columnToMeasureCode.TryGetValue(col, out var mc) ? mc : col; insertParameters.Add(paramName, measureValues.TryGetValue(measureCode, out var v) ? v : 0m); rowValueParamNames.Add("@" + paramName); } valueRows.Add($"({string.Join(", ", rowValueParamNames)})"); } var sql = WarehouseFactPostingQB.INSERT_ROWS .Replace("{FactTableName}", factTableName) .Replace("{ColumnList}", string.Join(", ", columnNames)) .Replace("{QualifiedValueColumnList}", string.Join(", ", columnNames.Select(c => "V." + c))) .Replace("{ValueColumnList}", string.Join(", ", columnNames)) .Replace("{ValueRows}", string.Join(",\n ", valueRows)); await _queryExecutor.ExecuteAsync(loginDTO, sql, insertParameters, tx, ct).ConfigureAwait(false); } private static string ValidateIdentifier(string value) { var upper = value.ToUpperInvariant(); if (!SafeIdentifier.IsMatch(upper)) throw new InvalidOperationException( $"'{value}' is not a safe SQL identifier (expected [A-Z0-9_]+) — check the MWAREHOUSEMEASURE/POSTINGRULE metadata."); return upper; } private static string ValidateIntList(string commaSeparatedInts) { var parts = commaSeparatedInts.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries); foreach (var part in parts) { if (!int.TryParse(part, out _)) throw new InvalidOperationException( $"'{commaSeparatedInts}' is not a valid comma-separated integer list — check MWAREHOUSEMEASUREPOSTINGRULE metadata."); } return string.Join(",", parts); } } }