using System; using System.Collections.Generic; using System.Linq; using System.Security.Cryptography; using System.Text; using System.Text.Json; using AnalyticsBLL.DatasetResolver; using AnalyticsBLL.ExecutionRouting; using AnalyticsDAL.CustomCode.Analysis; using AnalyticsDAL.CustomCode.ApiDataset; using AnalyticsDAL.CustomCode.Warehouse; using AnalyticsDAL.DTO.AnalysisAggregation.V1; using AnalyticsDAL.DTO.BICatalog; using AnalyticsDAL.DTO.BIFieldMapping; using AnalyticsDAL.DTO.Warehouse; using GB5Shared.DALCache; using GB5Shared.Export; using GB5Shared.DTO.Framework.Criteria; using GB5Shared.DTO.Framework.Login; using GB5Shared.Telemetry; using Microsoft.Extensions.Logging; using static GB5Shared.GB5Constant.Constant; namespace AnalyticsBLL.AnalysisAggregation { public class AnalysisAggregationService : IAnalysisAggregationService { private readonly IDatasetResolver _datasetResolver; private readonly IWarehouseDAL _warehouseDAL; private readonly IAnalysisDAL _analysisDAL; private readonly IApiDatasetDAL _apiDatasetDAL; private readonly IDALCache _dalCache; private readonly IAnalysisExecutionRouteResolver _routeResolver; private readonly ILogger _logger; public AnalysisAggregationService( IDatasetResolver datasetResolver, IWarehouseDAL warehouseDAL, IAnalysisDAL analysisDAL, IApiDatasetDAL apiDatasetDAL, IDALCache dalCache, IAnalysisExecutionRouteResolver routeResolver, ILogger logger) { _datasetResolver = datasetResolver; _warehouseDAL = warehouseDAL; _analysisDAL = analysisDAL; _apiDatasetDAL = apiDatasetDAL; _dalCache = dalCache; _routeResolver = routeResolver; _logger = logger; } public async Task ExecuteAsync(AnalysisQueryDefinition definition, LoginDTO login, CancellationToken ct) { try { var resolution = await ResolveAndValidateAsync(definition, login, ct).ConfigureAwait(false); var cacheKey = $"AnalysisResult:{login.ClientId}:{ComputeDefinitionHash(definition, resolution.ComputedFields)}"; var result = await _dalCache.GetOrSetAsync( cacheKey, innerCt => BuildResultAsync(resolution, definition, login, innerCt), DALCache.TtlForLevel(CacheKeyLevel.CLIENT_LEVEL), ct).ConfigureAwait(false); return result ?? throw new InvalidOperationException("AnalysisResult build returned null."); } catch (Exception ex) { GB5Trace.MarkFailed("execute-analysisquery-failed", ex); _logger.LogError(ex, "AnalysisAggregationService.ExecuteAsync failed for DatasetId {Id}", definition.DatasetId); throw; } } /// Resolves the dataset and validates additivity/aggregation only — no /// execution, no caching. Used both by ExecuteAsync (immediately before running the /// query) and by ValidateDefinitionAsync (Phase 1's save-time trust check for /// BIViewBLL.Save), so the two callers can never drift onto different rules. private async Task ResolveAndValidateAsync(AnalysisQueryDefinition definition, LoginDTO login, CancellationToken ct) { GB5Trace.Step("resolve-dataset", new { definition.DatasetId }); var resolution = await _datasetResolver.ResolveAsync(definition.DatasetId, login, ct).ConfigureAwait(false); ValidateAdditivity(definition, resolution); ValidateComputedFieldRequest(definition, resolution); return resolution; } /// Pre-flight checks that don't need row data: computed-field name resolution, /// row-locality, reference/cycle validation (against the REQUESTED Dimensions/Measures /// field strings — a best-effort static check; BuildResultAsync re-validates against the /// actual resolved row keys, which is authoritative since Warehouse-kind dimension /// references can alias, e.g. a requested "DIMOU.REGIONID" resolves to result key /// "REGIONID"), and the Sort conflict below. Running here means BIViewBLL.Save rejects a /// bad computed-field reference at save time via ValidateDefinitionAsync, not just at /// query time. private static void ValidateComputedFieldRequest(AnalysisQueryDefinition definition, DatasetResolutionDTO resolution) { if (!ComputedFieldPlanner.IsRequested(definition)) return; var badSort = definition.Sort.FirstOrDefault(s => definition.ComputedFields.Contains(s.Field, StringComparer.OrdinalIgnoreCase)); if (badSort != null) throw new ArgumentException( $"Sorting by computed field '{badSort.Field}' is not supported — TopN is applied before computed " + "fields are evaluated, so an in-memory post-sort would return the wrong N rows. Sort by a source field instead."); var requestedFieldNames = definition.Dimensions.Select(d => d.Field) .Concat(definition.Measures.Select(m => m.Field)) .ToHashSet(StringComparer.OrdinalIgnoreCase); ComputedFieldPlanner.SelectAndOrder(resolution.ComputedFields ?? new(), definition.ComputedFields, requestedFieldNames); } public async Task ValidateDefinitionAsync(AnalysisQueryDefinition definition, LoginDTO login, CancellationToken ct) { try { await ResolveAndValidateAsync(definition, login, ct).ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("validate-definition-failed", ex); _logger.LogError(ex, "AnalysisAggregationService.ValidateDefinitionAsync failed for DatasetId {Id}", definition.DatasetId); throw; } } // ── Additivity validation — before touching the DB ────────────────────── /// Rejects a requested aggregation that would produce a meaningless number /// for a semi-additive or non-additive measure. Warehouse-kind measures are validated /// against the authoritative MWAREHOUSEMEASURE classification (resolution.Measures); /// ApiService-kind measures are validated against MBIFIELDMAPPING /// (resolution.ApiMeasures) the same way. AnalysisQuery-kind measures have no /// server-side classification yet, so whatever Additivity the caller supplied on the /// definition is trusted as-is (defaults to Additive) — this is a known, documented v1 /// gap for that kind specifically, not an oversight. private static void ValidateAdditivity(AnalysisQueryDefinition definition, DatasetResolutionDTO resolution) { var registryByColumn = resolution.Measures?.ToDictionary(m => m.ColumnName, StringComparer.OrdinalIgnoreCase) ?? new Dictionary(StringComparer.OrdinalIgnoreCase); var apiRegistryByField = resolution.ApiMeasures?.ToDictionary(m => m.SourceFieldName, StringComparer.OrdinalIgnoreCase) ?? new Dictionary(StringComparer.OrdinalIgnoreCase); foreach (var measure in definition.Measures) { AdditivityType additivity; string? semiAdditiveDimension; if (resolution.Kind == DatasetKind.Warehouse) { if (!registryByColumn.TryGetValue(measure.Field, out var registered)) throw new ArgumentException($"Measure '{measure.Field}' is not a registered measure for this warehouse fact.", nameof(definition)); additivity = registered.AdditivityType; semiAdditiveDimension = registered.SemiAdditiveDimension; measure.Additivity = additivity; // echo the authoritative classification back into Meta measure.SemiAdditiveDimension = semiAdditiveDimension; } else if (resolution.Kind == DatasetKind.ApiService) { if (!apiRegistryByField.TryGetValue(measure.Field, out var registeredApi)) throw new ArgumentException($"Measure '{measure.Field}' is not a registered measure mapping for this ApiService dataset.", nameof(definition)); additivity = registeredApi.AdditivityType ?? AdditivityType.Additive; semiAdditiveDimension = registeredApi.SemiAdditiveDimension; measure.Additivity = additivity; // same authoritative-echo as Warehouse — never trust the caller's value measure.SemiAdditiveDimension = semiAdditiveDimension; } else { additivity = measure.Additivity; semiAdditiveDimension = measure.SemiAdditiveDimension; } switch (additivity) { case AdditivityType.Additive: break; // any aggregation allowed case AdditivityType.SemiAdditive: // Only safe to SUM when the query's dimensions still include the exact // semi-additive grain key (e.g. "DATEID") — grouping by a coarser date // rollup (month/quarter/year) collapses multiple point-in-time balances // just as wrongly as omitting the date dimension entirely. if (measure.Aggregation == 1 /* SUM */ && !definition.Dimensions.Any(d => string.Equals(d.Field, semiAdditiveDimension, StringComparison.OrdinalIgnoreCase))) { throw new ArgumentException( $"SUM is not valid across '{semiAdditiveDimension}' for semi-additive measure '{measure.Field}' " + $"— include '{semiAdditiveDimension}' as a dimension, or use MAX/AVG instead."); } break; case AdditivityType.NonAdditive: if (measure.Aggregation is not (0 or 7 or 8)) // None, MAX, MIN throw new ArgumentException( $"Only None/MAX/MIN aggregations are valid for non-additive measure '{measure.Field}' " + "(it is already an aggregated ratio/average) — request its underlying additive components instead."); break; } } } // ── Execution — branch on resolution.Kind ──────────────────────────────── private async Task BuildResultAsync( DatasetResolutionDTO resolution, AnalysisQueryDefinition definition, LoginDTO login, CancellationToken ct) { List> rows; bool truncated = false; switch (resolution.Kind) { case DatasetKind.Warehouse: GB5Trace.Step("execute-warehouse-aggregation", new { resolution.FactId }); rows = await _warehouseDAL.ExecuteAggregationAsync(resolution, definition, login, ct).ConfigureAwait(false); break; case DatasetKind.AnalysisQuery: GB5Trace.Step("execute-analysisquery-dynamicoutput", new { resolution.AnalysisQueryId }); // Route resolution mirrors AnalysisBLL.DynamicOutput's own resolve-then-call // shape — without this, the query silently always executes locally via // IQueryExecutor/LoginDTO, even when the analysis's workspace has an // OverrideDataSourceId pointing at an external SqlWorkbench-managed client DB. var route = resolution.AnalysisId is int analysisId ? await _routeResolver.ResolveAsync(analysisId, login, ct).ConfigureAwait(false) : null; var aqRows = await _analysisDAL.DynamicOutput( resolution.AnalysisQueryId!.Value, resolution.ReportViewId ?? -1, 0, definition.TopN ?? int.MaxValue, BuildCriteriaDTO(definition), login, ct, route).ConfigureAwait(false); rows = aqRows.ConvertAll(r => r.ToDictionary(kv => kv.Key, kv => (object?)kv.Value)); break; case DatasetKind.ApiService: GB5Trace.Step("execute-apiservice-fetch", new { resolution.DataSourceId }); var rawRows = await _apiDatasetDAL.FetchRowsAsync(resolution, definition, login, ct).ConfigureAwait(false); // GROUP BY cannot be pushed into an external HTTP call — group in-memory here. rows = GroupInMemory(rawRows, definition); break; default: throw new InvalidOperationException($"Unhandled DatasetKind '{resolution.Kind}'."); } if (definition.TopN is int topN && rows.Count > topN) { rows = rows.Take(topN).ToList(); truncated = true; } var computedFieldsUsed = ComputedFieldPlanner.IsRequested(definition) ? ApplyComputedFields(resolution, definition, rows) : new List(); return new AnalysisResult { Meta = new AnalysisResultMetaDTO { DimensionsUsed = definition.Dimensions, MeasuresUsed = definition.Measures, GeneratedAt = DateTime.UtcNow, TenantId = login.ClientId, RowCount = rows.Count, Truncated = truncated, Insights = null, ComputedFieldsUsed = computedFieldsUsed }, Data = rows }; } /// Evaluates every requested computed field against the already-resolved result /// rows, mutating them in place, and returns the Meta.ComputedFieldsUsed entries (with /// NullRowCount filled in). Runs AFTER TopN truncation — TopN is applied twice for /// Warehouse-kind (SQL TOP(n), then again above), so computing before truncation would /// waste work on rows about to be discarded; RowCount is unaffected either way. private List ApplyComputedFields( DatasetResolutionDTO resolution, AnalysisQueryDefinition definition, List> rows) { var resultKeys = rows.Count > 0 ? new HashSet(rows[0].Keys, StringComparer.OrdinalIgnoreCase) : new HashSet( definition.Dimensions.Select(d => d.Field).Concat(definition.Measures.Select(m => m.Field)), StringComparer.OrdinalIgnoreCase); var ordered = ComputedFieldPlanner.SelectAndOrder(resolution.ComputedFields ?? new(), definition.ComputedFields, resultKeys); if (ordered.Count == 0 || rows.Count == 0) return new List(); var additivityByKey = BuildAdditivityLookup(resolution, definition); var meta = ComputedFieldPlanner.BuildMeta(ordered, additivityByKey); var evaluatorSpecs = ordered.Select(f => new GB5Shared.DTO.Report.ComputedFieldDTO { FieldName = f.SourceFieldName, FieldType = f.DataType, // MBIFIELDMAPPING.DATATYPE uses the identical 0/1/2/3/4/6 // encoding as MREPORTVSFIELDS.FIELDTYPE (confirmed via // BICatalogBLL.TranslateDbFieldDataType's own output range) Expression = f.ComputeExpression ?? string.Empty, }).ToList(); var evaluator = ComputedFieldEvaluator.Create(evaluatorSpecs, _logger); if (evaluator == null) return new List(); var nullCounts = new Dictionary(StringComparer.OrdinalIgnoreCase); foreach (var row in rows) { evaluator.Apply(row); foreach (var spec in evaluatorSpecs) { if (!row.TryGetValue(spec.FieldName, out var v) || v is null) nullCounts[spec.FieldName] = nullCounts.GetValueOrDefault(spec.FieldName) + 1; } } foreach (var m in meta) m.NullRowCount = nullCounts.GetValueOrDefault(m.Field); return meta; } /// A lookup of every resolved-row key's additivity classification, so /// ComputedFieldPlanner.BuildMeta can report each computed field's InheritedAdditivity. /// Additive is the default for anything not explicitly classified (e.g. a plain /// dimension) — matches ValidateAdditivity's own default for AnalysisQuery-kind /// measures. private static Dictionary BuildAdditivityLookup( DatasetResolutionDTO resolution, AnalysisQueryDefinition definition) { var lookup = new Dictionary(StringComparer.OrdinalIgnoreCase); foreach (var measure in definition.Measures) lookup[measure.Field] = measure.Additivity; return lookup; } /// In-memory GROUP BY for the ApiService path — the target microservice /// returns whatever grain it returns (often already line-level); this reproduces the /// same SUM/AVG/COUNT/MAX/MIN semantics as the SQL paths, just in LINQ. private static List> GroupInMemory( List> rawRows, AnalysisQueryDefinition definition) { if (definition.Dimensions.Count == 0) { // No dimensions requested — a single aggregate row over everything. var singleRow = new Dictionary(); foreach (var measure in definition.Measures) singleRow[measure.Field] = AggregateValues(rawRows.Select(r => r.GetValueOrDefault(measure.Field)), measure.Aggregation); return new List> { singleRow }; } var grouped = rawRows.GroupBy(r => string.Join("", definition.Dimensions.Select(d => r.GetValueOrDefault(d.Field)?.ToString() ?? ""))); var result = new List>(); foreach (var group in grouped) { var row = new Dictionary(); var first = group.First(); foreach (var dim in definition.Dimensions) row[dim.Field] = first.GetValueOrDefault(dim.Field); foreach (var measure in definition.Measures) row[measure.Field] = AggregateValues(group.Select(r => r.GetValueOrDefault(measure.Field)), measure.Aggregation); result.Add(row); } return result; } private static object? AggregateValues(IEnumerable values, byte aggregation) { var numeric = values.Select(v => v is null ? (double?)null : Convert.ToDouble(v)).Where(v => v.HasValue).Select(v => v!.Value).ToList(); return aggregation switch { 0 => values.FirstOrDefault(), // None — raw/first value 1 => numeric.Sum(), // SUM 3 => numeric.Count > 0 ? numeric.Average() : 0d, // AVG 5 => (double)values.Count(), // COUNT 7 => numeric.Count > 0 ? numeric.Max() : (object?)null, // MAX 8 => numeric.Count > 0 ? numeric.Min() : (object?)null, // MIN _ => throw new ArgumentException($"Unknown aggregation code {aggregation}.", nameof(aggregation)) }; } private static CriteriaDTO BuildCriteriaDTO(AnalysisQueryDefinition definition) { var attributeCriteria = definition.Filters.Select(f => new AttributesCriteriaDTO { FieldName = f.Field, OperationType = f.Operator, FieldValue = f.Values.Length > 0 ? f.Values[0] : null, InArray = f.Values, JoinType = CriteriaDTO.AttributeJoinOperationType.And }).ToList(); return new CriteriaDTO { SectionCriteriaList = new List { new SectionCriteriaDTO { SectionId = 0, OperationType = CriteriaDTO.SectionOperationType.And, AttributesCriteriaList = attributeCriteria } } }; } // Folds in the REQUESTED subset of resolution.ComputedFields (name/expression/type) so an // admin editing an existing computed field's COMPUTEEXPRESSION invalidates the cache // immediately rather than serving up to CLIENT_LEVEL's TTL of stale values — the // definition JSON alone only captures the requested NAMES (AnalysisQueryDefinition. // ComputedFields), never the expression text itself, since that lives in the DB catalog, // not on the request. Hashing only the requested subset (not the whole catalog) means // editing an unrelated computed field never invalidates unrelated cached queries. private static string ComputeDefinitionHash(AnalysisQueryDefinition definition, List? computedFields) { var requested = computedFields? .Where(f => definition.ComputedFields.Contains(f.SourceFieldName, StringComparer.OrdinalIgnoreCase)) .OrderBy(f => f.SourceFieldName, StringComparer.OrdinalIgnoreCase) .Select(f => new { f.SourceFieldName, f.ComputeExpression, f.DataType }) .ToList(); var json = JsonSerializer.Serialize(definition) + JsonSerializer.Serialize(requested); var hashBytes = SHA256.HashData(Encoding.UTF8.GetBytes(json)); return Convert.ToBase64String(hashBytes); } } }